{"record":{"id":"f1b871395b6f9b18","repo":"apache/seatunnel","slug":"common-writer-operation-failed-commonerrorcodedep","errorCode":"COMMON_WRITER_OPERATION_FAILED (CommonErrorCodeDeprecated.WRITER_OPERATION_FAILED)","errorMessage":"Close kafka sink writer error","messagePattern":"Close kafka sink writer error","errorType":"error_code","errorClass":"KafkaConnectorException","httpStatus":null,"severity":"error","filePath":"seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaSinkWriter.java","lineNumber":150,"sourceCode":"        return states;\n    }\n\n    @Override\n    public Optional<KafkaCommitInfo> prepareCommit() {\n        return kafkaProducerSender.prepareCommit();\n    }\n\n    @Override\n    public void abortPrepare() {\n        kafkaProducerSender.abortTransaction();\n    }\n\n    @Override\n    public void close() {\n        try {\n            kafkaProducerSender.close();\n        } catch (Exception e) {\n            throw new KafkaConnectorException(\n                    CommonErrorCodeDeprecated.WRITER_OPERATION_FAILED,\n                    \"Close kafka sink writer error\",\n                    e);\n        }\n    }\n\n    private Properties getKafkaProperties(ReadonlyConfig pluginConfig) {\n        Properties kafkaProperties = new Properties();\n        if (pluginConfig.get(KAFKA_CONFIG) != null) {\n            pluginConfig.get(KAFKA_CONFIG).forEach((key, value) -> kafkaProperties.put(key, value));\n        }\n\n        if (pluginConfig.get(ASSIGN_PARTITIONS) != null) {\n            kafkaProperties.put(\n                    ProducerConfig.PARTITIONER_CLASS_CONFIG,\n                    \"org.apache.seatunnel.connectors.seatunnel.kafka.sink.MessageContentPartitioner\");\n            kafkaProperties.put(\n                    MessageContentPartitioner.ASSIGN_PARTITIONS_CONFIG,","sourceCodeStart":132,"sourceCodeEnd":168,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaSinkWriter.java#L132-L168","documentation":"KafkaSinkWriter.close() wraps any exception thrown while closing the underlying KafkaProducerSender (producer.close()) into a KafkaConnectorException with WRITER_OPERATION_FAILED. It signals that the sink writer could not shut down its Kafka producer cleanly, usually because the producer was in a bad state or the close timed out / was interrupted.","triggerScenarios":"Thrown from KafkaSinkWriter.close() when kafkaProducerSender.close() raises any Exception — e.g. producer.close() timing out, an in-flight transaction stuck in abortable-error state, or the producer already closed fatally.","commonSituations":"Job cancellation or checkpoint-complete shutdown with records still buffered; Kafka broker unreachable during final flush; transactional producer left in an abortable error state after a broker failure; calling close() twice.","solutions":["Check broker connectivity and that the Kafka cluster is healthy at shutdown time","Inspect the wrapped cause ('Caused by') for producer-level errors such as TimeoutException or authentication failure","Verify transactional.id uniqueness across concurrent jobs to avoid fencing that breaks close","Increase close timeout settings or ensure the task thread is not interrupted mid-close"],"exampleFix":null,"handlingStrategy":"try-catch","validationCode":null,"typeGuard":null,"tryCatchPattern":"try {\n  sinkWriter.close();\n} catch (KafkaConnectorException e) {\n  log.error(\"Kafka sink close failed\", e.getCause());\n  // ensure producer resources are reclaimed; rely on engine failure handling\n}","preventionTips":["Keep broker connectivity healthy through job shutdown","Use unique transactional.id per writer to avoid fencing","Avoid interrupting the task thread during close"],"tags":["kafka","sink-writer","resource-close","producer"],"backgroundTag":"resource-cleanup-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"}