{"record":{"id":"8353cdff390dada1","repo":"apache/seatunnel","slug":"connection-failed-8353cd","errorCode":"CONNECTION_FAILED","errorMessage":"Failed to connect MQTT source client [${clientId}]","messagePattern":"Failed to connect MQTT source client \\[(.+?)\\]","errorType":"error_code","errorClass":"MqttConnectorException","httpStatus":null,"severity":"critical","filePath":"seatunnel-connectors-v2/connector-mqtt/src/main/java/org/apache/seatunnel/connectors/seatunnel/mqtt/source/MqttSourceReader.java","lineNumber":98,"sourceCode":"\n    @Override\n    public void open() {\n        try {\n            this.mqttClient =\n                    new MqttClient(\n                            sourceConfig.getUrl(),\n                            sourceConfig.getClientId(),\n                            new MemoryPersistence());\n            this.mqttClient.setCallback(this);\n            this.mqttClient.connect(buildConnectOptions(sourceConfig));\n            subscribeTopic();\n            LOG.info(\n                    \"MQTT source reader [{}] subscribed to topic [{}]\",\n                    sourceConfig.getClientId(),\n                    sourceConfig.getTopic());\n        } catch (MqttException e) {\n            closeClientQuietly();\n            throw new MqttConnectorException(\n                    MqttConnectorErrorCode.CONNECTION_FAILED,\n                    \"Failed to connect MQTT source client [\" + sourceConfig.getClientId() + \"]\",\n                    e);\n        }\n    }\n\n    @Override\n    public void pollNext(Collector<SeaTunnelRow> output) throws Exception {\n        checkReceiveException();\n        checkReconnectTimeout();\n\n        byte[] payload = messageQueue.poll(POLL_TIMEOUT_MS, TimeUnit.MILLISECONDS);\n        if (payload == null) {\n            return;\n        }\n\n        SeaTunnelRow row = deserializationSchema.deserialize(payload);\n        if (row == null) {","sourceCodeStart":80,"sourceCodeEnd":116,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-mqtt/src/main/java/org/apache/seatunnel/connectors/seatunnel/mqtt/source/MqttSourceReader.java#L80-L116","documentation":"MqttSourceReader.open throws MqttConnectorException with code CONNECTION_FAILED when the underlying Paho client's connect or subscribe call throws MqttException. Before throwing, it quietly closes the client so no half-connected client leaks. Callers get the clientId in the message and the original MqttException as cause.","triggerScenarios":"MqttClient.connect() or subscribe() fails during reader open: broker unreachable, authentication rejected, TLS handshake failure, client_id already connected (server kick-out), or resubscribe after reconnect failing as in the referenced tests.","commonSituations":"Broker down or wrong broker URL/port; wrong username/password or ACLs denying subscription; duplicate client_id from another process or another SeaTunnel instance; TLS certificate mismatch.","solutions":["Read the chained MqttException reason code for the exact broker rejection cause","Verify broker URL, port, credentials and ACL permissions in the source config","Ensure client_id is unique across all running clients to avoid server kick-outs","If using TLS, validate CA certificates and hostname settings; check network reachability (firewall, DNS)","Let the engine retry the task if the outage is transient, or increase reconnect_timeout for later recovery"],"exampleFix":"// before\nsource {\n  Mqtt {\n    broker = \"tcp://broker.local:1883\"\n    client_id = \"seatunnel-source-1\"\n  }\n}\n// after\n// verify broker reachable and credentials valid, e.g.:\n// mosquitto_sub -h broker.local -p 1883 -t 'topic' -u user -P pass -i seatunnel-source-1\nsource {\n  Mqtt {\n    broker = \"tcp://broker.local:1883\"\n    client_id = \"seatunnel-source-1-unique\"\n    username = \"user\"\n    password = \"pass\"\n  }\n}","handlingStrategy":"try-catch","validationCode":"// preflight: verify broker connectivity before job submission\nMqttClient probe = new MqttClient(brokerUrl, probeId, new MemoryPersistence());\nprobe.connect(new MqttConnectOptions());\nprobe.disconnect();","typeGuard":"null","tryCatchPattern":"try {\n  reader.open(...);\n} catch (MqttConnectorException e) {\n  if (e.getCode() == MqttConnectorErrorCode.CONNECTION_FAILED\n      && e.getCause() instanceof MqttException mEx) {\n    LOG.error(\"MQTT connect failed, reason code {}\", mEx.getReasonCode(), mEx);\n  }\n  throw e;\n}","preventionTips":["Verify broker URL, port, credentials and ACLs before deployment","Guarantee client_id uniqueness across processes and clusters","Validate TLS/CA configuration when using ssl:// brokers","Monitor broker availability and set reconnect_timeout for transient outage recovery"],"tags":["mqtt","connection-failed","mqtt-exception","source"],"backgroundTag":"connection-refused","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"}