{"record":{"id":"7e96eaba08734e6d","repo":"apache/seatunnel","slug":"suppressed-an-additional-asynchronous-send-failure","errorCode":null,"errorMessage":"Suppressed an additional asynchronous send failure of Kafka transaction [{}]","messagePattern":"Suppressed an additional asynchronous send failure of Kafka transaction \\[(.+?)\\]","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaTransactionSender.java","lineNumber":94,"sourceCode":"        checkAsyncSendException();\n        kafkaProducer.send(producerRecord, this::onSendCompleted);\n        recordNumInTransaction++;\n    }\n\n    /**\n     * Records the first asynchronous send failure of the current transaction so that {@link\n     * #prepareCommit()} can fail the checkpoint instead of committing a partial transaction.\n     *\n     * <p>Invoked on the producer's sender thread.\n     */\n    private void onSendCompleted(RecordMetadata metadata, Exception exception) {\n        if (exception == null) {\n            return;\n        }\n        if (!asyncSendException.compareAndSet(null, exception)) {\n            // Only the first failure becomes the checkpoint failure cause. Log the later ones so a\n            // broker-side incident affecting several partitions can still be diagnosed.\n            log.warn(\n                    \"Suppressed an additional asynchronous send failure of Kafka transaction [{}]\",\n                    transactionId,\n                    exception);\n        }\n    }\n\n    @Override\n    public void beginTransaction(String transactionId) {\n        this.transactionId = transactionId;\n        this.kafkaProducer = getTransactionProducer(transactionId);\n        kafkaProducer.beginTransaction();\n        // Reset the per-transaction state. A new transaction always runs on a newly created\n        // producer, so a failure recorded for the previous transaction no longer applies. Keeping\n        // it would turn a single transient send error into a permanent checkpoint failure loop.\n        recordNumInTransaction = 0;\n        asyncSendException.set(null);\n    }\n","sourceCodeStart":76,"sourceCodeEnd":112,"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#L76-L112","documentation":"Kafka's producer callback for an ongoing transaction reported an exception after a previous send failure was already captured. Only the first exception becomes the checkpoint failure cause; subsequent ones are logged and suppressed to avoid masking the primary error.","triggerScenarios":"onSendCompleted (producer callback) receives a non-null exception while asyncSendException is already set — e.g. multiple partitions' sends fail during the same broker outage, and only the first failure is retained.","commonSituations":"Broker outage or leader election affecting several partitions at once; transactional.id timeouts; authorization failures on produce after an earlier failure; network partition during transaction commit window.","solutions":["Inspect this suppressed log for additional affected partitions and root causes","Check the first asyncSendException / checkpoint failure exception, which is the primary cause","Verify broker health, replication, and that the transactional producer can reach the cluster","Review ACLs (IDEVENT/WRITE) for the transactional producer user"],"exampleFix":null,"handlingStrategy":"try-catch","validationCode":"// Check transactional producer health before/at checkpoint time\nproducer.partitionsFor(topic); // fails fast on broker/ACL issues","typeGuard":null,"tryCatchPattern":"// In callback: retain first failure, log the rest\nAtomicReference<Exception> first = new AtomicReference<>();\ncallback -> {\n    if (e != null && !first.compareAndSet(null, e)) {\n        LOG.warn(\"Suppressed additional send failure\", e);\n    }\n}","preventionTips":["Monitor broker availability and under-replicated partitions","Ensure transactional.id user has WRITE ACLs on all target partitions","Alert on the first (non-suppressed) async send failure — it is the checkpoint cause"],"tags":["kafka","transactions","producer"],"backgroundTag":"kafka-send-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"}