{"record":{"id":"46515941878b7b15","repo":"apache/seatunnel","slug":"writer-operation-failed-465159","errorCode":"WRITER_OPERATION_FAILED","errorMessage":"Close RocketMq sink writer error","messagePattern":"Close RocketMq sink writer error","errorType":"error_code","errorClass":"RocketMqConnectorException","httpStatus":null,"severity":"error","filePath":"seatunnel-connectors-v2/connector-rocketmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rocketmq/sink/RocketMqSinkWriter.java","lineNumber":69,"sourceCode":"        }\n        // Set `rocketmq.client.logUseSlf4j` to `true` to avoid create many\n        // `AsyncAppender-Dispatcher-Thread`\n        System.setProperty(\"rocketmq.client.logUseSlf4j\", \"true\");\n    }\n\n    @Override\n    public void write(SeaTunnelRow element) throws IOException {\n        Message message = seaTunnelRowSerializer.serializeRow(element);\n        rocketMqProducerSender.send(message);\n    }\n\n    @Override\n    public void close() throws IOException {\n        if (this.rocketMqProducerSender != null) {\n            try {\n                this.rocketMqProducerSender.close();\n            } catch (Exception e) {\n                throw new RocketMqConnectorException(\n                        CommonErrorCodeDeprecated.WRITER_OPERATION_FAILED,\n                        \"Close RocketMq sink writer error\",\n                        e);\n            }\n        }\n    }\n\n    private SeaTunnelRowSerializer<byte[], byte[]> getSerializer(\n            SeaTunnelRowType seaTunnelRowType) {\n        return new DefaultSeaTunnelRowSerializer(\n                producerMetadata.getTopic(),\n                producerMetadata.getTag(),\n                getPartitionKeyFields(seaTunnelRowType),\n                seaTunnelRowType,\n                producerMetadata.getFormat(),\n                producerMetadata.getFieldDelimiter());\n    }\n","sourceCodeStart":51,"sourceCodeEnd":87,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-rocketmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rocketmq/sink/RocketMqSinkWriter.java#L51-L87","documentation":"RocketMqSinkWriter.close closes the underlying RocketMQ producer sender. If sender.close() throws any Exception, it is rewrapped as RocketMqConnectorException with WRITER_OPERATION_FAILED and the message 'Close RocketMq sink writer error', preserving the original as the cause.","triggerScenarios":"Writer cleanup (job shutdown, checkpoint teardown, or test code closing writers) when rocketMqProducerSender.close() fails — e.g. in-flight send shutdown errors or producer client cleanup exceptions.","commonSituations":"Producer already half-closed due to an earlier failure; network problems during producer shutdown; broker unreachable while flushing on close; shared client closed twice across multi-table writers.","solutions":["Inspect the wrapped cause (getCause()) for the real RocketMQ client error","Ensure producer shutdown completes/flushes before close and that clients aren't closed twice in multi-writer setups","Verify broker connectivity and producer shutdownTimeout settings; log-and-suppress in close if failing close shouldn't mask the original job error"],"exampleFix":"// before\npublic void close() throws IOException {\n    this.rocketMqProducerSender.close();\n}\n// after\npublic void close() throws IOException {\n    try {\n        this.rocketMqProducerSender.close();\n    } catch (Exception e) {\n        log.warn(\"Error closing RocketMQ producer\", e);\n    }\n}","handlingStrategy":"try-catch","validationCode":"if (rocketMqProducerSender != null && !rocketMqProducerSender.isClosed()) rocketMqProducerSender.close();","typeGuard":null,"tryCatchPattern":"try { writer.close(); } catch (RocketMqConnectorException e) { log.warn(\"RocketMQ producer close failed\", e.getCause()); }","preventionTips":["Make close idempotent and guard with a closed flag","Never let close() mask the original job exception; log-and-warn instead","Test shutdown paths (multi-writer shared clients) to avoid double-close"],"tags":["rocketmq","sink","close","cleanup"],"backgroundTag":"writer-close-failed","analyzedSha":"cf67b549a7a6c35fa0beb12d83c62892427ea919","analyzedAt":"2026-09-10T21:44:55.265Z","contentChangedAt":"2026-09-10T21:44:55.265Z","schemaVersion":2},"datasetVersion":"2026-09-14T11:17:12.474Z"}