apache/seatunnel · warning

Transient MQTT publish failure, retrying...

Error message

Transient MQTT publish failure, retrying...

What it means

publishWithRetry caught an MqttException during publish while the client reports connected; the error is treated as transient and the method retries after RETRY_BACKOFF_MS. After exhausting retries it throws an IOException carrying lastException, which fails the flush.

Solutions

  1. Let the built-in retry handle it; if flush ultimately fails, the SeaTunnel framework will retry/restore from checkpoint.
  2. Increase the retry count/backoff options if available in MqttSinkOptions.
  3. Check broker logs for the publish rejection cause.
  4. Reduce flush batch size/frequency to lower inflight pressure.
Defensive patterns

Strategy: retry

Try / catch

try {
    publishWithRetry(topic, message);
} catch (IOException e) {
    // retries exhausted; rely on checkpoint restart or surface the lastException
    throw new IOException("MQTT publish failed after retries", e.getCause());
}

Prevention

When it happens

Trigger: mqttClient.publish(topic, message) throws MqttException (transient broker error, QoS timeout) inside publishWithRetry called from flushBuffer.

Common situations: Broker under load rejecting publishes, QoS1 timeouts, brief network blips mid-batch, broker hitting max inflight window.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/fd2131f30bcf0678. 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:202

        }
        for (MqttMessage message : messageBuffer) {
            publishWithRetry(message);
        }
        messageBuffer.clear();
    }

    private void publishWithRetry(MqttMessage message) throws IOException {
        long deadline = System.currentTimeMillis() + retryTimeoutMs;
        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();

View on GitHub (pinned to cf67b549a7)