apache/seatunnel · warning · IOException
Interrupted during MQTT publish retry
Error message
Interrupted during MQTT publish retry
What it means
Thrown by publishWithRetry in MqttSinkWriter when the retry loop's backoff Thread.sleep is interrupted while waiting to re-publish a transiently failed MQTT message. The interrupt flag is restored before throwing so upstream shutdown logic still sees the interrupt. This distinguishes a job cancellation/shutdown during retry from a genuine publish failure.
Solutions
- Inspect the chained InterruptedException cause: this is expected during job shutdown, so treat it as a cancellation, not a data error
- If it fires unexpectedly, check for code that interrupts the writer thread prematurely (e.g. custom thread pools around the sink)
- Fix the underlying transient publish failure (broker availability, QoS settings) so the retry backoff path is rarely entered
- Increase tolerance by lowering publish rate or QoS to reduce retry pressure during shutdown windows
Example fix
// before
try {
publish(message);
} catch (MqttException e) {
log.warn("Transient MQTT publish failure, retrying...", e);
Thread.sleep(RETRY_BACKOFF_MS);
}
// after
try {
publish(message);
} catch (MqttException 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);
}
} Defensive patterns
Strategy: try-catch
Validate before calling
null
Type guard
null
Try / catch
try {
writer.flush();
} catch (IOException e) {
if (Thread.currentThread().isInterrupted() || e.getCause() instanceof InterruptedException) {
// job shutdown; stop gracefully, do not alert on data failure
return;
}
throw e;
} Prevention
- Avoid long blocking backoff sleeps by checking thread interrupt status before retrying
- Monitor broker health so transient failures (and thus retry sleeps) are rare
- Treat interrupt-based errors during shutdown as normal cancellation, not data loss
When it happens
Trigger: A transient MQTT publish fails, the writer enters its backoff sleep (RETRY_BACKOFF_MS), and the thread is interrupted by task cancellation, checkpoint/flush shutdown, or job stop before the sleep completes.
Common situations: Stopping a SeaTunnel job while the sink is retrying against a slow or flapping MQTT broker; cluster node shutdown during backpressure; misjudging this as a broker failure when it is really a cancellation.
Related errors
- Connection retry sleep interrupted by exception:
- Failed to get connection, interrupted while doing another…
- MqttConnectorErrorCode.PUBLISH_FAILED
- RECEIVE_FAILED
- Transient MQTT publish failure, retrying...
AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10).
Data as JSON: /api/errors/ce5f1e5146b8b046.
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:208
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();
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.");View on GitHub (pinned to cf67b549a7)