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

  1. Inspect the chained InterruptedException cause: this is expected during job shutdown, so treat it as a cancellation, not a data error
  2. If it fires unexpectedly, check for code that interrupts the writer thread prematurely (e.g. custom thread pools around the sink)
  3. Fix the underlying transient publish failure (broker availability, QoS settings) so the retry backoff path is rarely entered
  4. 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

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


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)