{"record":{"id":"c58532a585b9d082","repo":"apache/seatunnel","slug":"failed-to-commit-consumer-offsets-for-checkpoint","errorCode":null,"errorMessage":"Failed to commit consumer offsets for checkpoint {}","messagePattern":"Failed to commit consumer offsets for checkpoint (.+?)","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/source/KafkaSourceReader.java","lineNumber":154,"sourceCode":"                checkpointOffsetMap.get(checkpointId);\n\n        if (committedPartitions == null) {\n            logger.debug(\"Offsets for checkpoint {} have already been committed.\", checkpointId);\n            return;\n        }\n\n        if (committedPartitions.isEmpty()) {\n            logger.debug(\"There are no offsets to commit for checkpoint {}.\", checkpointId);\n            removeAllOffsetsToCommitUpToCheckpoint(checkpointId);\n            return;\n        }\n\n        ((KafkaSourceFetcherManager) splitFetcherManager)\n                .commitOffsets(\n                        committedPartitions,\n                        (ignored, e) -> {\n                            if (e != null) {\n                                logger.warn(\n                                        \"Failed to commit consumer offsets for checkpoint {}\",\n                                        checkpointId,\n                                        e);\n                                return;\n                            }\n                            offsetsOfFinishedSplits\n                                    .keySet()\n                                    .removeIf(committedPartitions::containsKey);\n                            removeAllOffsetsToCommitUpToCheckpoint(checkpointId);\n                        });\n    }\n\n    private void removeAllOffsetsToCommitUpToCheckpoint(long checkpointId) {\n        while (!checkpointOffsetMap.isEmpty() && checkpointOffsetMap.firstKey() <= checkpointId) {\n            checkpointOffsetMap.remove(checkpointOffsetMap.firstKey());\n        }\n    }\n}","sourceCodeStart":136,"sourceCodeEnd":172,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/source/KafkaSourceReader.java#L136-L172","documentation":"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.","triggerScenarios":"notifyCheckpointComplete invokes KafkaSourceFetcherManager.commitOffsets; the async callback receives a non-null exception (broker rebalance in progress, offsets metadata changed, broker unavailable).","commonSituations":"Consumer group rebalancing when the checkpoint completes; offsets expired or partition ownership changed; Kafka cluster briefly unavailable; group coordinator restarts during commit.","solutions":["Check the chained exception in this log for the exact commit error","Verify Kafka broker and consumer group health (kafka-consumer-groups.sh --describe)","Monitor for repeated occurrences — occasional rebalance-driven warnings are benign","Ensure checkpoint interval is not extremely small relative to rebalance windows"],"exampleFix":null,"handlingStrategy":"retry","validationCode":"// Check consumer group stability before/around commits\nkafka-consumer-groups.sh --bootstrap-server broker:9092 --describe --group <group>","typeGuard":null,"tryCatchPattern":"consumer.commitAsync(offsets, (offsets, e) -> {\n    if (e != null) {\n        LOG.warn(\"Offset commit failed for checkpoint {}\", cpId, e);\n        // offsets are still recovered from checkpoint on restore\n    }\n});","preventionTips":["Accept occasional commit warnings during rebalances; alert on persistent ones","Keep checkpoint intervals reasonably long vs rebalance frequency","Monitor consumer group lag to detect chronic commit failures"],"tags":["kafka","offsets","checkpoint"],"backgroundTag":"offset-commit-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"}