apache/iceberg · warning

Error aborting producer transaction

Error message

Error aborting producer transaction

What it means

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.

Solutions

  1. 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.
  2. Check for duplicate workers sharing the same transactional.id (connector task restarts or misconfigured max tasks) and ensure only one instance is active.
  3. Increase transaction.timeout.ms if long commit cycles cause coordinator-side transaction expiry.
  4. Verify broker reachability and that the transaction coordinator is available; retry the Connect task after the broker recovers.
Defensive patterns

Strategy: try-catch

Try / catch

try {
  producer.commitTransaction();
} catch (Exception e) {
  try {
    producer.abortTransaction();
  } catch (Exception ex) {
    LOG.warn("Error aborting producer transaction", ex); // inspect 'ex' for fencing/coordinator errors
  }
  throw e; // recreate the producer / restart the task
}

Prevention

When it happens

Trigger: 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).

Common situations: 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.

Related errors


AI-assisted analysis of apache/iceberg@86d9c8fc54 (2026-09-12). Data as JSON: /api/errors/d812b9fd06ee2011. Report an issue: GitHub.

Appendix: source

Thrown at kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/Channel.java:111

                })
            .collect(Collectors.toList());

    synchronized (producer) {
      producer.beginTransaction();
      try {
        // NOTE: we shouldn't call get() on the future in a transactional context,
        // see docs for org.apache.kafka.clients.producer.KafkaProducer
        recordList.forEach(producer::send);
        if (!sourceOffsets.isEmpty()) {
          producer.sendOffsetsToTransaction(
              offsetsToCommit, KafkaUtils.consumerGroupMetadata(context));
        }
        producer.commitTransaction();
      } catch (Exception e) {
        try {
          producer.abortTransaction();
        } catch (Exception ex) {
          LOG.warn("Error aborting producer transaction", ex);
        }
        throw e;
      }
    }
  }

  protected abstract boolean receive(Envelope envelope);

  protected void consumeAvailable(Duration pollDuration) {
    ConsumerRecords<String, byte[]> records = consumer.poll(pollDuration);
    while (!records.isEmpty()) {
      records.forEach(
          record -> {
            // the consumer stores the offsets that corresponds to the next record to consume,
            // so increment the record offset by one. Keep the highest position seen for the
            // partition: a re-read of the control topic, e.g. after a rebalance resumes from the
            // last committed offsets, would otherwise move the tracked position backwards and
            // commit a consumer offset behind records that were already handled.

View on GitHub (pinned to 86d9c8fc54)