apache/seatunnel · critical · KafkaConnectorException

KAFKA_TRANSACTION_NOT_STARTED

KAFKA_TRANSACTION_NOT_STARTED

Error message

Kafka transaction [%s] has %d record(s) but is still reported as not started after flushing pending sends. The transaction registration did not complete. Refusing to commit to avoid data loss.

What it means

KafkaTransactionSender.prepareCommit() detects that records were buffered and sent in the current transaction (recordNumInTransaction > 0) but the Kafka producer still reports the transaction as not started even after a flush. Committing in that state would cause the committer to skip EndTxn and silently drop the records, so it throws KAFKA_TRANSACTION_NOT_STARTED to fail the checkpoint instead of losing data.

Solutions

  1. Let the checkpoint abort and retry — this exception is designed to trigger a safe transaction abort
  2. Check Kafka broker logs for transaction coordinator errors (fencing, timeout, coordinator migration) around the failure time
  3. Verify transactional.id is unique per writer and not reused by a still-live producer instance
  4. Increase transaction.timeout.ms on the producer if sends are slow and the transaction is expiring
Defensive patterns

Strategy: retry

Try / catch

try {
  writer.prepareCommit();
} catch (KafkaConnectorException e) {
  if (e.getCode() == KafkaConnectorErrorCode.TRANSACTION_NOT_STARTED) {
    // abort transaction; engine will retry checkpoint
    abortTransaction();
  }
}

Prevention

When it happens

Trigger: Transactional produce where txnStarted remains false despite recordNumInTransaction > 0 after flushPendingSends() in prepareCommit(); indicates send() never triggered/observed the transaction registration (e.g. initTransactions/partition callback path failed silently).

Common situations: Broker-side transaction coordinator issues (transaction timeout, coordinator rebalance); producer fenced between beginTransaction and first send; async send callbacks not yet run so the 'started' flag never flipped.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/ff3dc1c9d2508f8c. 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:131

    @Override
    public Optional<KafkaCommitInfo> prepareCommit() {
        // Flush pending async sends before capturing the transaction state. Kafka only marks the
        // transaction as started once the AddPartitionsToTxn request has been acknowledged by the
        // broker, and that request is issued asynchronously by the producer's sender thread. If a
        // checkpoint reaches this point before the first record's transaction registration
        // completes, isTxnStarted() would still return false and the resulting commit info would
        // instruct the committer to skip EndTxn, leaving the transaction to time out and its
        // records permanently invisible to read_committed consumers.
        kafkaProducer.flush();
        checkAsyncSendException();
        boolean txnStarted = kafkaProducer.isTxnStarted();
        if (recordNumInTransaction > 0 && !txnStarted) {
            // Records were sent in this transaction but Kafka still reports it as not started even
            // after flushing, meaning the transaction registration never completed. Committing with
            // txnStarted=false would make the committer skip EndTxn and drop these records, so fail
            // fast and let the checkpoint abort this transaction instead of silently producing a
            // lossy commit info.
            throw new KafkaConnectorException(
                    KafkaConnectorErrorCode.TRANSACTION_NOT_STARTED,
                    String.format(
                            "Kafka transaction [%s] has %d record(s) but is still reported as not "
                                    + "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);
    }

View on GitHub (pinned to cf67b549a7)