apache/seatunnel · warning

Failed to commit consumer offsets for checkpoint

Error message

Failed to commit consumer offsets for checkpoint {}

What it means

KafkaSourceReader attempted to commit consumer offsets to Kafka as part of checkpoint completion and the commit callback returned an exception. The reader logs a warning and continues — checkpointing itself has already succeeded, so data is not lost, but the committed offsets in Kafka may lag.

Solutions

  1. Check the chained exception in this log for the exact commit error
  2. Verify Kafka broker and consumer group health (kafka-consumer-groups.sh --describe)
  3. Monitor for repeated occurrences — occasional rebalance-driven warnings are benign
  4. Ensure checkpoint interval is not extremely small relative to rebalance windows
Defensive patterns

Strategy: retry

Validate before calling

// Check consumer group stability before/around commits
kafka-consumer-groups.sh --bootstrap-server broker:9092 --describe --group <group>

Try / catch

consumer.commitAsync(offsets, (offsets, e) -> {
    if (e != null) {
        LOG.warn("Offset commit failed for checkpoint {}", cpId, e);
        // offsets are still recovered from checkpoint on restore
    }
});

Prevention

When it happens

Trigger: notifyCheckpointComplete invokes KafkaSourceFetcherManager.commitOffsets; the async callback receives a non-null exception (broker rebalance in progress, offsets metadata changed, broker unavailable).

Common situations: Consumer group rebalancing when the checkpoint completes; offsets expired or partition ownership changed; Kafka cluster briefly unavailable; group coordinator restarts during commit.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/c58532a585b9d082. Report an issue: GitHub.

Appendix: source

Thrown at seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/source/KafkaSourceReader.java:154

                checkpointOffsetMap.get(checkpointId);

        if (committedPartitions == null) {
            logger.debug("Offsets for checkpoint {} have already been committed.", checkpointId);
            return;
        }

        if (committedPartitions.isEmpty()) {
            logger.debug("There are no offsets to commit for checkpoint {}.", checkpointId);
            removeAllOffsetsToCommitUpToCheckpoint(checkpointId);
            return;
        }

        ((KafkaSourceFetcherManager) splitFetcherManager)
                .commitOffsets(
                        committedPartitions,
                        (ignored, e) -> {
                            if (e != null) {
                                logger.warn(
                                        "Failed to commit consumer offsets for checkpoint {}",
                                        checkpointId,
                                        e);
                                return;
                            }
                            offsetsOfFinishedSplits
                                    .keySet()
                                    .removeIf(committedPartitions::containsKey);
                            removeAllOffsetsToCommitUpToCheckpoint(checkpointId);
                        });
    }

    private void removeAllOffsetsToCommitUpToCheckpoint(long checkpointId) {
        while (!checkpointOffsetMap.isEmpty() && checkpointOffsetMap.firstKey() <= checkpointId) {
            checkpointOffsetMap.remove(checkpointOffsetMap.firstKey());
        }
    }
}

View on GitHub (pinned to cf67b549a7)