{"record":{"id":"ce5f1e5146b8b046","repo":"apache/seatunnel","slug":"interrupted-during-mqtt-publish-retry","errorCode":null,"errorMessage":"Interrupted during MQTT publish retry","messagePattern":"Interrupted during MQTT publish retry","errorType":"exception","errorClass":"IOException","httpStatus":null,"severity":"warning","filePath":"seatunnel-connectors-v2/connector-mqtt/src/main/java/org/apache/seatunnel/connectors/seatunnel/mqtt/sink/MqttSinkWriter.java","lineNumber":208,"sourceCode":"\n    private void publishWithRetry(MqttMessage message) throws IOException {\n        long deadline = System.currentTimeMillis() + retryTimeoutMs;\n        MqttException lastException = null;\n        while (System.currentTimeMillis() < deadline) {\n            try {\n                if (mqttClient.isConnected()) {\n                    mqttClient.publish(topic, message);\n                    return;\n                }\n            } catch (MqttException e) {\n                lastException = e;\n                log.warn(\"Transient MQTT publish failure, retrying...\", e);\n            }\n            try {\n                Thread.sleep(RETRY_BACKOFF_MS);\n            } catch (InterruptedException ie) {\n                Thread.currentThread().interrupt();\n                throw new IOException(\"Interrupted during MQTT publish retry\", ie);\n            }\n        }\n        throw new IOException(\n                new MqttConnectorException(\n                                MqttConnectorErrorCode.PUBLISH_FAILED,\n                                \"Failed to publish MQTT message after \" + retryTimeoutMs + \"ms\")\n                        .getMessage(),\n                lastException);\n    }\n\n    private static MqttConnectOptions buildConnectOptions(ReadonlyConfig config) {\n        MqttConnectOptions options = new MqttConnectOptions();\n        options.setAutomaticReconnect(true);\n        boolean cleanSession = config.get(MqttSinkOptions.CLEAN_SESSION);\n        options.setCleanSession(cleanSession);\n        if (!cleanSession) {\n            log.warn(\n                    \"clean_session=false may cause broker-side state accumulation. Ensure proper clientId management.\");","sourceCodeStart":190,"sourceCodeEnd":226,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-mqtt/src/main/java/org/apache/seatunnel/connectors/seatunnel/mqtt/sink/MqttSinkWriter.java#L190-L226","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","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"],"exampleFix":"// before\ntry {\n    publish(message);\n} catch (MqttException e) {\n    log.warn(\"Transient MQTT publish failure, retrying...\", e);\n    Thread.sleep(RETRY_BACKOFF_MS);\n}\n// after\ntry {\n    publish(message);\n} catch (MqttException e) {\n    log.warn(\"Transient MQTT publish failure, retrying...\", e);\n    try {\n        Thread.sleep(RETRY_BACKOFF_MS);\n    } catch (InterruptedException ie) {\n        Thread.currentThread().interrupt();\n        throw new IOException(\"Interrupted during MQTT publish retry\", ie);\n    }\n}","handlingStrategy":"try-catch","validationCode":"null","typeGuard":"null","tryCatchPattern":"try {\n  writer.flush();\n} catch (IOException e) {\n  if (Thread.currentThread().isInterrupted() || e.getCause() instanceof InterruptedException) {\n    // job shutdown; stop gracefully, do not alert on data failure\n    return;\n  }\n  throw e;\n}","preventionTips":["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"],"tags":["mqtt","interrupted","retry","io-exception"],"backgroundTag":"thread-interrupted","analyzedSha":"cf67b549a7a6c35fa0beb12d83c62892427ea919","analyzedAt":"2026-09-10T21:44:55.265Z","contentChangedAt":"2026-09-10T21:44:55.265Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}