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
- Let the checkpoint abort and retry — this exception is designed to trigger a safe transaction abort
- Check Kafka broker logs for transaction coordinator errors (fencing, timeout, coordinator migration) around the failure time
- Verify transactional.id is unique per writer and not reused by a still-live producer instance
- 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
- Size transaction.timeout.ms generously for slow flushes
- Avoid producer fencing by unique transactional.id per writer
- Monitor Kafka transaction coordinator logs and rebalances
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
- KAFKA_GET_TRANSACTIONMANAGER_FAILED
- KAFKA_PRODUCE_DATA_FAILED
- KAFKA_VERSION_INCOMPATIBLE
- Suppressed an additional asynchronous send failure of Kafka…
- COMMON-02
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)