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
- Let the built-in retry handle it; if flush ultimately fails, the SeaTunnel framework will retry/restore from checkpoint.
- Increase the retry count/backoff options if available in MqttSinkOptions.
- Check broker logs for the publish rejection cause.
- 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
- Keep QoS and inflight settings within broker limits
- Increase retry count/backoff if the broker is intermittently loaded
- Alert on flush failures — they mean retries were exhausted
- Check broker logs for publish rejections matching retry storms
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
- Interrupted during MQTT publish retry
- Airtable API rate limit reached, retry
- Airtable API rate limit reached, retry
- Ambiguous timeout on Couchbase write
- Batch label changed from
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)