This is an automated email from the ASF dual-hosted git repository. davsclaus pushed a commit to branch fix/CAMEL-23994-offset-repo-off-by-one in repository https://gitbox.apache.org/repos/asf/camel.git
commit eaa19b7500d8be1784b2812a0f49ed07b56a35d1 Author: Claus Ibsen <[email protected]> AuthorDate: Tue Jul 14 16:39:33 2026 +0200 CAMEL-23994: Fix off-by-one when writing to offset repository SyncCommitManager and AsyncCommitManager were storing offset+1 in the offset repository, but OffsetPartitionAssignmentAdapter.resumeFromOffset adds +1 when seeking (state+1). This double increment caused one message per partition to be lost on every restart. Now both managers store the raw record offset, matching the contract established by CommitToOffsetManager and the resumeFromOffset comment: "The state contains the last read offset, so seek from the next one." Co-Authored-By: Claude Opus 4.6 <[email protected]> Signed-off-by: Claus Ibsen <[email protected]> --- .../apache/camel/component/kafka/consumer/AsyncCommitManager.java | 5 ++++- .../org/apache/camel/component/kafka/consumer/SyncCommitManager.java | 2 +- .../kafka/integration/commit/BaseManualCommitTestSupport.java | 4 ++-- 3 files changed, 7 insertions(+), 4 deletions(-) diff --git a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/consumer/AsyncCommitManager.java b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/consumer/AsyncCommitManager.java index b7f776ac56b4..d1d927e982a0 100644 --- a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/consumer/AsyncCommitManager.java +++ b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/consumer/AsyncCommitManager.java @@ -94,7 +94,10 @@ public class AsyncCommitManager extends AbstractCommitManager { if (exception == null) { if (offsetRepository != null) { for (var entry : committed.entrySet()) { - saveStateToOffsetRepository(entry.getKey(), entry.getValue().offset(), offsetRepository); + // Store the last read offset (Kafka committed offset - 1) to match + // the contract in OffsetPartitionAssignmentAdapter.resumeFromOffset + // which seeks to state + 1 + saveStateToOffsetRepository(entry.getKey(), entry.getValue().offset() - 1, offsetRepository); } } } diff --git a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/consumer/SyncCommitManager.java b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/consumer/SyncCommitManager.java index 9850e45d1569..3d2e3d2a6e15 100644 --- a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/consumer/SyncCommitManager.java +++ b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/consumer/SyncCommitManager.java @@ -75,7 +75,7 @@ public class SyncCommitManager extends AbstractCommitManager { consumer.commitSync(offsets, Duration.ofMillis(timeout)); if (offsetRepository != null) { - saveStateToOffsetRepository(partition, lastOffset, offsetRepository); + saveStateToOffsetRepository(partition, offset, offsetRepository); } offsetCache.removeCommittedEntries(offsets, null); diff --git a/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/integration/commit/BaseManualCommitTestSupport.java b/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/integration/commit/BaseManualCommitTestSupport.java index 1935fa7567ef..b6a71e22172e 100644 --- a/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/integration/commit/BaseManualCommitTestSupport.java +++ b/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/integration/commit/BaseManualCommitTestSupport.java @@ -97,8 +97,8 @@ abstract class BaseManualCommitTestSupport extends BaseKafkaTestSupport { final String state = Awaitility.await().until(() -> stateRepository.getState(topic + "/0"), Matchers.notNullValue()); - // We send 5 records initially, so we expect the offset to be 5 after first step execution - assertEquals("5", state, "5 messages were sent in the first step, therefore the offset should be 5"); + // We send 5 records (offsets 0-4), so we expect the last read offset (4) to be stored + assertEquals("4", state, "5 messages were sent in the first step, therefore the last read offset should be 4"); // Second step: We shut down our route, we expect nothing will be recovered by our route contextExtension.getContext().getRouteController().stopRoute("foo");
