{"record":{"id":"d812b9fd06ee2011","repo":"apache/iceberg","slug":"error-aborting-producer-transaction","errorCode":null,"errorMessage":"Error aborting producer transaction","messagePattern":"Error aborting producer transaction","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/Channel.java","lineNumber":111,"sourceCode":"                })\n            .collect(Collectors.toList());\n\n    synchronized (producer) {\n      producer.beginTransaction();\n      try {\n        // NOTE: we shouldn't call get() on the future in a transactional context,\n        // see docs for org.apache.kafka.clients.producer.KafkaProducer\n        recordList.forEach(producer::send);\n        if (!sourceOffsets.isEmpty()) {\n          producer.sendOffsetsToTransaction(\n              offsetsToCommit, KafkaUtils.consumerGroupMetadata(context));\n        }\n        producer.commitTransaction();\n      } catch (Exception e) {\n        try {\n          producer.abortTransaction();\n        } catch (Exception ex) {\n          LOG.warn(\"Error aborting producer transaction\", ex);\n        }\n        throw e;\n      }\n    }\n  }\n\n  protected abstract boolean receive(Envelope envelope);\n\n  protected void consumeAvailable(Duration pollDuration) {\n    ConsumerRecords<String, byte[]> records = consumer.poll(pollDuration);\n    while (!records.isEmpty()) {\n      records.forEach(\n          record -> {\n            // the consumer stores the offsets that corresponds to the next record to consume,\n            // so increment the record offset by one. Keep the highest position seen for the\n            // partition: a re-read of the control topic, e.g. after a rebalance resumes from the\n            // last committed offsets, would otherwise move the tracked position backwards and\n            // commit a consumer offset behind records that were already handled.","sourceCodeStart":93,"sourceCodeEnd":129,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/Channel.java#L93-L129","documentation":"This is a warning logged in Channel.send when the Kafka producer's commitTransaction() fails and the subsequent abortTransaction() also throws. The original commit exception is rethrown; this warning records that the transactional abort itself could not be completed cleanly, leaving the producer's transactional state uncertain.","triggerScenarios":"producer.commitTransaction() throws (e.g. broker unreachable, transactional.id fenced by a newer producer instance, unknown producer epoch) and the immediately following producer.abortTransaction() also throws (e.g. producer already fenced into an invalid state, coordinator unavailable).","commonSituations":"Kafka broker restart or network partition during a Connect sink commit; duplicate connector workers with the same transactional.id fencing each other; long-running transactions exceeding transaction.timeout.ms so the coordinator expires the producer.","solutions":["Inspect the suppressed 'ex' warning to determine why the abort failed; in most fencing cases the producer is unusable and the task must be restarted so a fresh producer with a bumped epoch is created.","Check for duplicate workers sharing the same transactional.id (connector task restarts or misconfigured max tasks) and ensure only one instance is active.","Increase transaction.timeout.ms if long commit cycles cause coordinator-side transaction expiry.","Verify broker reachability and that the transaction coordinator is available; retry the Connect task after the broker recovers."],"exampleFix":null,"handlingStrategy":"try-catch","validationCode":null,"typeGuard":null,"tryCatchPattern":"try {\n  producer.commitTransaction();\n} catch (Exception e) {\n  try {\n    producer.abortTransaction();\n  } catch (Exception ex) {\n    LOG.warn(\"Error aborting producer transaction\", ex); // inspect 'ex' for fencing/coordinator errors\n  }\n  throw e; // recreate the producer / restart the task\n}","preventionTips":["Ensure a unique transactional.id per connector task instance so workers do not fence each other.","Size transaction.timeout.ms generously relative to your commit interval.","Monitor broker availability and transaction coordinator health.","Restart tasks on transactional errors rather than reusing a possibly-fenced producer."],"tags":["kafka","transactions","producer","connect"],"backgroundTag":"transaction-abort-failed","analyzedSha":"86d9c8fc543e7c56c9f624eb725f76c9baff9570","analyzedAt":"2026-09-12T00:46:39.097Z","contentChangedAt":"2026-09-12T00:46:39.097Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}