{"record":{"id":"ff3dc1c9d2508f8c","repo":"apache/seatunnel","slug":"kafka-transaction-not-started","errorCode":"KAFKA_TRANSACTION_NOT_STARTED","errorMessage":"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.","messagePattern":"Kafka transaction \\[(.+?)\\] has (.+?) 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\\.","errorType":"error_code","errorClass":"KafkaConnectorException","httpStatus":null,"severity":"critical","filePath":"seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaTransactionSender.java","lineNumber":131,"sourceCode":"    @Override\n    public Optional<KafkaCommitInfo> prepareCommit() {\n        // Flush pending async sends before capturing the transaction state. Kafka only marks the\n        // transaction as started once the AddPartitionsToTxn request has been acknowledged by the\n        // broker, and that request is issued asynchronously by the producer's sender thread. If a\n        // checkpoint reaches this point before the first record's transaction registration\n        // completes, isTxnStarted() would still return false and the resulting commit info would\n        // instruct the committer to skip EndTxn, leaving the transaction to time out and its\n        // records permanently invisible to read_committed consumers.\n        kafkaProducer.flush();\n        checkAsyncSendException();\n        boolean txnStarted = kafkaProducer.isTxnStarted();\n        if (recordNumInTransaction > 0 && !txnStarted) {\n            // Records were sent in this transaction but Kafka still reports it as not started even\n            // after flushing, meaning the transaction registration never completed. Committing with\n            // txnStarted=false would make the committer skip EndTxn and drop these records, so fail\n            // fast and let the checkpoint abort this transaction instead of silently producing a\n            // lossy commit info.\n            throw new KafkaConnectorException(\n                    KafkaConnectorErrorCode.TRANSACTION_NOT_STARTED,\n                    String.format(\n                            \"Kafka transaction [%s] has %d record(s) but is still reported as not \"\n                                    + \"started after flushing pending sends. The transaction \"\n                                    + \"registration did not complete. Refusing to commit to avoid \"\n                                    + \"data loss.\",\n                            transactionId, recordNumInTransaction));\n        }\n        KafkaCommitInfo kafkaCommitInfo =\n                new KafkaCommitInfo(\n                        transactionId,\n                        kafkaProperties,\n                        this.kafkaProducer.getProducerId(),\n                        this.kafkaProducer.getEpoch(),\n                        txnStarted);\n        return Optional.of(kafkaCommitInfo);\n    }\n","sourceCodeStart":113,"sourceCodeEnd":149,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaTransactionSender.java#L113-L149","documentation":"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.","triggerScenarios":"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).","commonSituations":"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.","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"],"exampleFix":null,"handlingStrategy":"retry","validationCode":null,"typeGuard":null,"tryCatchPattern":"try {\n  writer.prepareCommit();\n} catch (KafkaConnectorException e) {\n  if (e.getCode() == KafkaConnectorErrorCode.TRANSACTION_NOT_STARTED) {\n    // abort transaction; engine will retry checkpoint\n    abortTransaction();\n  }\n}","preventionTips":["Size transaction.timeout.ms generously for slow flushes","Avoid producer fencing by unique transactional.id per writer","Monitor Kafka transaction coordinator logs and rebalances"],"tags":["kafka","transactions","exactly-once","data-loss-prevention"],"backgroundTag":"transaction-not-started","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"}