{"record":{"id":"04a195b478be083a","repo":"apache/seatunnel","slug":"mqtt-connection-lost-auto-reconnect-will-attempt","errorCode":null,"errorMessage":"MQTT connection lost, auto-reconnect will attempt recovery","messagePattern":"MQTT connection lost, auto-reconnect will attempt recovery","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"seatunnel-connectors-v2/connector-mqtt/src/main/java/org/apache/seatunnel/connectors/seatunnel/mqtt/sink/MqttSinkWriter.java","lineNumber":166,"sourceCode":"                try {\n                    if (mqttClient.isConnected()) {\n                        mqttClient.disconnect();\n                    }\n                    mqttClient.close();\n                    log.info(\"MQTT sink writer closed\");\n                } catch (MqttException e) {\n                    throw new IOException(\"Error closing MQTT client\", e);\n                }\n            }\n        }\n    }\n\n    // ---- MqttCallback implementation ----\n\n    @Override\n    public void connectionLost(Throwable cause) {\n        // Auto-reconnect is enabled; log for observability but do not throw.\n        log.warn(\"MQTT connection lost, auto-reconnect will attempt recovery\", cause);\n    }\n\n    @Override\n    public void messageArrived(String topic, MqttMessage message) {\n        // Sink-only client — inbound messages are not expected.\n    }\n\n    @Override\n    public void deliveryComplete(IMqttDeliveryToken token) {\n        // QoS acknowledgement received from broker.\n    }\n\n    // ---- private helpers ----\n\n    private void flushBuffer() throws IOException {\n        if (messageBuffer.isEmpty()) {\n            return;\n        }","sourceCodeStart":148,"sourceCodeEnd":184,"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#L148-L184","documentation":"The MQTT sink's MqttCallback.connectionLost fired because the broker connection dropped. Automatic reconnect is enabled, so the writer deliberately logs a warning and does not throw; recovery is left to Paho's auto-reconnect. Buffered messages are flushed on the next successful connection.","triggerScenarios":"The broker closes the socket or the network breaks while MqttSinkWriter is connected; Paho invokes connectionLost(Throwable).","commonSituations":"Broker restarts, keep-alive timeouts, LB idle connection eviction, network partitions, broker-side client kick due to duplicate clientId.","solutions":["No immediate action needed if automatic reconnect is enabled (default in this sink); watch for repeated occurrences.","Ensure the clientId is unique per sink writer to avoid brokers kicking the connection.","Tune keep-alive/connection timeout options to fit network conditions.","Investigate the attached 'cause' for the root disconnect reason (e.g. 32109 socket loss)."],"exampleFix":"// before\nurl = \"tcp://broker:1883\"\n// after (TLS for unstable networks)\nurl = \"ssl://broker:8883\"","handlingStrategy":"try-catch","validationCode":null,"typeGuard":null,"tryCatchPattern":"// connectionLost is informational; detect stuck reconnection instead\n@Override\npublic void connectionLost(Throwable cause) {\n    log.warn(\"MQTT connection lost\", cause);\n    scheduleReconnectDeadlineCheck();\n}","preventionTips":["Use unique clientIds per writer to prevent broker kicks","Enable automaticReconnect and verify it stays connected after broker restarts","Monitor connectionLost frequency; clusters of events mean broker/network trouble","Use TLS for unstable networks"],"tags":["mqtt","connection-lost","reconnect","sink"],"backgroundTag":"connection-lost","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"}