apache/seatunnel · critical · KafkaConnectorException

KAFKA_PRODUCE_DATA_FAILED

KAFKA_PRODUCE_DATA_FAILED

Error message

Kafka transaction [%s] failed to send one or more of its %d record(s) asynchronously.

What it means

KafkaTransactionSender.checkAsyncSendException() rethrows a stored producer-callback failure as KAFKA_PRODUCE_DATA_FAILED when at least one of the transaction's records failed to send asynchronously. It is invoked from send() and prepareCommit() so a background send failure surfaces at the next checkpoint instead of being swallowed.

Solutions

  1. Read the 'Caused by' exception to identify the producer-level failure (timeout, size, auth, serialization)
  2. Check broker availability, network, and topic ACLs for the producing user
  3. Reduce batch/record size or increase max.request.size / delivery timeouts if record-too-large or timeout
  4. After fixing the cause, restart the pipeline from the last successful checkpoint (the transaction will be aborted)
Defensive patterns

Strategy: try-catch

Try / catch

try {
  writer.write(row);
} catch (KafkaConnectorException e) {
  if (e.getCode() == KafkaConnectorErrorCode.PRODUCE_DATA_FAILED) {
    Throwable cause = e.getCause(); // inspect producer failure: timeout, size, auth
    log.error("Async Kafka send failed", cause);
  }
}

Prevention

When it happens

Trigger: A KafkaProducer send callback recorded an exception in asyncSendException; the next send() call or prepareCommit() calls checkAsyncSendException() and throws with the transaction id, record count, and the original exception as cause.

Common situations: Broker unavailable or record-too-large / message-size-limit exceeded; serialization failures; TimeoutException under load; topic authorization failures discovered asynchronously.

Related errors


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

Appendix: source

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

                                    + "started after flushing pending sends. The transaction "
                                    + "registration did not complete. Refusing to commit to avoid "
                                    + "data loss.",
                            transactionId, recordNumInTransaction));
        }
        KafkaCommitInfo kafkaCommitInfo =
                new KafkaCommitInfo(
                        transactionId,
                        kafkaProperties,
                        this.kafkaProducer.getProducerId(),
                        this.kafkaProducer.getEpoch(),
                        txnStarted);
        return Optional.of(kafkaCommitInfo);
    }

    private void checkAsyncSendException() {
        Exception exception = asyncSendException.get();
        if (exception != null) {
            throw new KafkaConnectorException(
                    KafkaConnectorErrorCode.PRODUCE_DATA_FAILED,
                    String.format(
                            "Kafka transaction [%s] failed to send one or more of its %d record(s) "
                                    + "asynchronously.",
                            transactionId, recordNumInTransaction),
                    exception);
        }
    }

    @Override
    public void abortTransaction() {
        kafkaProducer.abortTransaction();
    }

    @Override
    public void abortTransaction(long checkpointId) {
        for (long i = checkpointId; ; i++) {
            String transactionId = generateTransactionId(this.transactionPrefix, i);

View on GitHub (pinned to cf67b549a7)