apache/seatunnel · error · KafkaConnectorException

COMMON_WRITER_OPERATION_FAILED (CommonErrorCodeDeprecated.WRITER_OPERATION_FAILED)

COMMON_WRITER_OPERATION_FAILED (CommonErrorCodeDeprecated.WRITER_OPERATION_FAILED)

Error message

Close kafka sink writer error

What it means

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.

Solutions

  1. Check broker connectivity and that the Kafka cluster is healthy at shutdown time
  2. Inspect the wrapped cause ('Caused by') for producer-level errors such as TimeoutException or authentication failure
  3. Verify transactional.id uniqueness across concurrent jobs to avoid fencing that breaks close
  4. Increase close timeout settings or ensure the task thread is not interrupted mid-close
Defensive patterns

Strategy: try-catch

Try / catch

try {
  sinkWriter.close();
} catch (KafkaConnectorException e) {
  log.error("Kafka sink close failed", e.getCause());
  // ensure producer resources are reclaimed; rely on engine failure handling
}

Prevention

When it happens

Trigger: 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.

Common situations: 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.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/f1b871395b6f9b18. Report an issue: GitHub.

Appendix: source

Thrown at seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaSinkWriter.java:150

        return states;
    }

    @Override
    public Optional<KafkaCommitInfo> prepareCommit() {
        return kafkaProducerSender.prepareCommit();
    }

    @Override
    public void abortPrepare() {
        kafkaProducerSender.abortTransaction();
    }

    @Override
    public void close() {
        try {
            kafkaProducerSender.close();
        } catch (Exception e) {
            throw new KafkaConnectorException(
                    CommonErrorCodeDeprecated.WRITER_OPERATION_FAILED,
                    "Close kafka sink writer error",
                    e);
        }
    }

    private Properties getKafkaProperties(ReadonlyConfig pluginConfig) {
        Properties kafkaProperties = new Properties();
        if (pluginConfig.get(KAFKA_CONFIG) != null) {
            pluginConfig.get(KAFKA_CONFIG).forEach((key, value) -> kafkaProperties.put(key, value));
        }

        if (pluginConfig.get(ASSIGN_PARTITIONS) != null) {
            kafkaProperties.put(
                    ProducerConfig.PARTITIONER_CLASS_CONFIG,
                    "org.apache.seatunnel.connectors.seatunnel.kafka.sink.MessageContentPartitioner");
            kafkaProperties.put(
                    MessageContentPartitioner.ASSIGN_PARTITIONS_CONFIG,

View on GitHub (pinned to cf67b549a7)