apache/seatunnel · critical · IOException

MqttConnectorErrorCode.PUBLISH_FAILED

MqttConnectorErrorCode.PUBLISH_FAILED

Error message

Failed to publish MQTT message after 

What it means

Thrown by publishWithRetry in MqttSinkWriter when an MQTT message still cannot be published after exhausting the retry window (retryTimeoutMs). It wraps MqttConnectorException with code MqttConnectorErrorCode.PUBLISH_FAILED and preserves the last broker-side exception as cause. This is the sink's terminal failure for an unflushable message.

Solutions

  1. Check the lastException cause chain for the root MqttException (reason code, socket error)
  2. Verify the broker is reachable and the connection options (broker URL, keep-alive) are correct in sink config
  3. Increase retry-timeout to tolerate longer broker outages, or add broker-side redundancy
  4. Check for client_id collisions if multiple writers share the same client_id

Example fix

// before
throw new IOException(
    new MqttConnectorException(MqttConnectorErrorCode.PUBLISH_FAILED,
        "Failed to publish MQTT message after " + retryTimeoutMs + "ms").getMessage(),
    lastException);
// after
// keep the error, but prevent it: ensure broker availability and a unique client_id,
// or raise retryTimeoutMs so transient broker restarts are absorbed by the retry loop
Defensive patterns

Strategy: retry

Validate before calling

// preflight broker reachability
MqttClient c = new MqttClient(brokerUrl, clientId, new MemoryPersistence());
c.connect(opts); // fail fast before job submission

Type guard

null

Try / catch

try {
  writer.flush();
} catch (IOException e) {
  Throwable cause = e.getCause();
  if (cause instanceof MqttConnectorException
      && ((MqttConnectorException) cause).getCode() == MqttConnectorErrorCode.PUBLISH_FAILED) {
    // inspect cause.getCause() for root MqttException reason code
    LOG.error("MQTT publish exhausted retries", cause);
  }
  throw e;
}

Prevention

When it happens

Trigger: publish() keeps throwing (connection loss, broker unavailable, client disconnected) for longer than retryTimeoutMs, so the retry loop exits and the final IOException is thrown from flushBuffer.

Common situations: MQTT broker down or restarting; network partition between SeaTunnel worker and broker; client_id conflict causing repeated disconnects; QoS 1/2 broker overload.

Understand the failure class

Background: "API request failed": what wrapped HTTP errors from external APIs mean and how to find the real cause — this error's family across 29 libraries.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/cc090d97b3600156. Report an issue: GitHub.

Appendix: source

Thrown at seatunnel-connectors-v2/connector-mqtt/src/main/java/org/apache/seatunnel/connectors/seatunnel/mqtt/sink/MqttSinkWriter.java:211

        MqttException lastException = null;
        while (System.currentTimeMillis() < deadline) {
            try {
                if (mqttClient.isConnected()) {
                    mqttClient.publish(topic, message);
                    return;
                }
            } catch (MqttException e) {
                lastException = e;
                log.warn("Transient MQTT publish failure, retrying...", e);
            }
            try {
                Thread.sleep(RETRY_BACKOFF_MS);
            } catch (InterruptedException ie) {
                Thread.currentThread().interrupt();
                throw new IOException("Interrupted during MQTT publish retry", ie);
            }
        }
        throw new IOException(
                new MqttConnectorException(
                                MqttConnectorErrorCode.PUBLISH_FAILED,
                                "Failed to publish MQTT message after " + retryTimeoutMs + "ms")
                        .getMessage(),
                lastException);
    }

    private static MqttConnectOptions buildConnectOptions(ReadonlyConfig config) {
        MqttConnectOptions options = new MqttConnectOptions();
        options.setAutomaticReconnect(true);
        boolean cleanSession = config.get(MqttSinkOptions.CLEAN_SESSION);
        options.setCleanSession(cleanSession);
        if (!cleanSession) {
            log.warn(
                    "clean_session=false may cause broker-side state accumulation. Ensure proper clientId management.");
        }
        options.setConnectionTimeout(config.get(MqttSinkOptions.CONNECTION_TIMEOUT));

View on GitHub (pinned to cf67b549a7)