{"record":{"id":"cc090d97b3600156","repo":"apache/seatunnel","slug":"mqttconnectorerrorcode-publish-failed","errorCode":"MqttConnectorErrorCode.PUBLISH_FAILED","errorMessage":"Failed to publish MQTT message after ","messagePattern":"Failed to publish MQTT message after ","errorType":"error_code","errorClass":"IOException","httpStatus":null,"severity":"critical","filePath":"seatunnel-connectors-v2/connector-mqtt/src/main/java/org/apache/seatunnel/connectors/seatunnel/mqtt/sink/MqttSinkWriter.java","lineNumber":211,"sourceCode":"        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.\");\n        }\n        options.setConnectionTimeout(config.get(MqttSinkOptions.CONNECTION_TIMEOUT));\n","sourceCodeStart":193,"sourceCodeEnd":229,"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#L193-L229","documentation":"Thrown by publishWithRetry in MqttSinkWriter when an MQTT message still cannot be published after exhausting the retry window (retryTimeoutMs). It wraps MqttConnectorException with code MqttConnectorErrorCode.PUBLISH_FAILED and preserves the last broker-side exception as cause. This is the sink's terminal failure for an unflushable message.","triggerScenarios":"publish() keeps throwing (connection loss, broker unavailable, client disconnected) for longer than retryTimeoutMs, so the retry loop exits and the final IOException is thrown from flushBuffer.","commonSituations":"MQTT broker down or restarting; network partition between SeaTunnel worker and broker; client_id conflict causing repeated disconnects; QoS 1/2 broker overload.","solutions":["Check the lastException cause chain for the root MqttException (reason code, socket error)","Verify the broker is reachable and the connection options (broker URL, keep-alive) are correct in sink config","Increase retry-timeout to tolerate longer broker outages, or add broker-side redundancy","Check for client_id collisions if multiple writers share the same client_id"],"exampleFix":"// before\nthrow new IOException(\n    new MqttConnectorException(MqttConnectorErrorCode.PUBLISH_FAILED,\n        \"Failed to publish MQTT message after \" + retryTimeoutMs + \"ms\").getMessage(),\n    lastException);\n// after\n// keep the error, but prevent it: ensure broker availability and a unique client_id,\n// or raise retryTimeoutMs so transient broker restarts are absorbed by the retry loop","handlingStrategy":"retry","validationCode":"// preflight broker reachability\nMqttClient c = new MqttClient(brokerUrl, clientId, new MemoryPersistence());\nc.connect(opts); // fail fast before job submission","typeGuard":"null","tryCatchPattern":"try {\n  writer.flush();\n} catch (IOException e) {\n  Throwable cause = e.getCause();\n  if (cause instanceof MqttConnectorException\n      && ((MqttConnectorException) cause).getCode() == MqttConnectorErrorCode.PUBLISH_FAILED) {\n    // inspect cause.getCause() for root MqttException reason code\n    LOG.error(\"MQTT publish exhausted retries\", cause);\n  }\n  throw e;\n}","preventionTips":["Configure a retry-timeout that covers expected broker restart windows","Ensure unique client_id per writer to avoid disconnect loops","Add broker monitoring/alerting for availability before running long jobs","Keep messages small to reduce broker timeout risk"],"tags":["mqtt","publish-failed","retry-exhausted","io-exception"],"backgroundTag":"api-request-failed","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"}