{"record":{"id":"9e354c037b0225da","repo":"apache/seatunnel","slug":"clean-session-false-may-cause-broker-side-state-ac","errorCode":null,"errorMessage":"clean_session=false may cause broker-side state accumulation. Ensure proper clientId management.","messagePattern":"clean_session=false may cause broker-side state accumulation\\. Ensure proper clientId management\\.","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":225,"sourceCode":"                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\n        String username = config.get(MqttSinkOptions.USERNAME);\n        if (username != null && !username.isEmpty()) {\n            options.setUserName(username);\n        }\n        String password = config.get(MqttSinkOptions.PASSWORD);\n        if (password != null && !password.isEmpty()) {\n            options.setPassword(password.toCharArray());\n        }\n        return options;\n    }\n\n    private static SerializationSchema createSerializationSchema(\n            SeaTunnelRowType rowType, ReadonlyConfig config) {\n        String format = config.get(MqttSinkOptions.FORMAT);","sourceCodeStart":207,"sourceCodeEnd":243,"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#L207-L243","documentation":"When clean_session=false, the MQTT broker retains subscription and queued-message state per clientId across connections. The sink warns that without careful, unique clientId management, broker-side state (queued messages, session state) can accumulate indefinitely and consume broker resources.","triggerScenarios":"buildConnectOptions reads MqttSinkOptions.CLEAN_SESSION and it is false; the warning is logged during option construction before connecting.","commonSituations":"Users set clean_session=false expecting at-least-once semantics without realizing each writer instance needs a distinct persistent clientId; multiple parallel sink writers sharing a clientId cause session conflicts and state buildup.","solutions":["Set clean_session=true unless you specifically need broker-persisted state for a sink (rare — sinks publish, they don't consume).","If clean_session=false is required, guarantee a unique clientId per writer instance (include subtask index).","Clean up stale sessions on the broker (e.g. MQTT 5 session expiry interval, or broker admin cleanup).","Monitor broker session/storage metrics for growth."],"exampleFix":"// before\nsink {\n  Mqtt {\n    clean_session = false\n  }\n}\n// after\nsink {\n  Mqtt {\n    clean_session = true\n  }\n}","handlingStrategy":"validation","validationCode":"// validate sink config before submit\nif (!config.getBoolean(\"clean_session\")) {\n    log.warn(\"Persistent session requested: ensure unique clientId per writer\");\n}","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Default to clean_session=true for sinks","If persistent sessions are needed, template clientIds with writer index","Periodically purge stale broker sessions","Watch broker retained-state/storage metrics"],"tags":["mqtt","clean-session","client-id","sink"],"backgroundTag":"invalid-config-value","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"}