{"record":{"id":"fb2489b45c71e91c","repo":"apache/seatunnel","slug":"kafka-produce-data-failed","errorCode":"KAFKA_PRODUCE_DATA_FAILED","errorMessage":"Kafka transaction [%s] failed to send one or more of its %d record(s) asynchronously.","messagePattern":"Kafka transaction \\[(.+?)\\] failed to send one or more of its (.+?) record\\(s\\) asynchronously\\.","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":153,"sourceCode":"                                    + \"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\n    private void checkAsyncSendException() {\n        Exception exception = asyncSendException.get();\n        if (exception != null) {\n            throw new KafkaConnectorException(\n                    KafkaConnectorErrorCode.PRODUCE_DATA_FAILED,\n                    String.format(\n                            \"Kafka transaction [%s] failed to send one or more of its %d record(s) \"\n                                    + \"asynchronously.\",\n                            transactionId, recordNumInTransaction),\n                    exception);\n        }\n    }\n\n    @Override\n    public void abortTransaction() {\n        kafkaProducer.abortTransaction();\n    }\n\n    @Override\n    public void abortTransaction(long checkpointId) {\n        for (long i = checkpointId; ; i++) {\n            String transactionId = generateTransactionId(this.transactionPrefix, i);","sourceCodeStart":135,"sourceCodeEnd":171,"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#L135-L171","documentation":"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.","triggerScenarios":"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.","commonSituations":"Broker unavailable or record-too-large / message-size-limit exceeded; serialization failures; TimeoutException under load; topic authorization failures discovered asynchronously.","solutions":["Read the 'Caused by' exception to identify the producer-level failure (timeout, size, auth, serialization)","Check broker availability, network, and topic ACLs for the producing user","Reduce batch/record size or increase max.request.size / delivery timeouts if record-too-large or timeout","After fixing the cause, restart the pipeline from the last successful checkpoint (the transaction will be aborted)"],"exampleFix":null,"handlingStrategy":"try-catch","validationCode":null,"typeGuard":null,"tryCatchPattern":"try {\n  writer.write(row);\n} catch (KafkaConnectorException e) {\n  if (e.getCode() == KafkaConnectorErrorCode.PRODUCE_DATA_FAILED) {\n    Throwable cause = e.getCause(); // inspect producer failure: timeout, size, auth\n    log.error(\"Async Kafka send failed\", cause);\n  }\n}","preventionTips":["Pre-validate record size against max.request.size and broker message.max.bytes","Keep brokers reachable and monitor producer error-rate metrics","Configure sensible delivery.timeout.ms and retries"],"tags":["kafka","async-send","producer","exactly-once"],"backgroundTag":"produce-failed","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"}