This is an automated email from the ASF dual-hosted git repository. davsclaus pushed a commit to branch fix/CAMEL-23994-batching-manual-commit in repository https://gitbox.apache.org/repos/asf/camel.git
commit 7ea01ed918a79e1466153b4815c4b9a0c88f836c Author: Claus Ibsen <[email protected]> AuthorDate: Tue Jul 14 16:40:51 2026 +0200 CAMEL-23994: Fix batching manual async commit never reaching Kafka In batching mode with manual async commit, calling manual.commit() only recorded the offset in OffsetCache via recordOffset() but the cache was never flushed to Kafka because the batching facade never called commitManager.commit(partition). Now the batching facade flushes recorded offsets for each partition after processExchange(), mirroring how the streaming facade works. Co-Authored-By: Claude Opus 4.6 <[email protected]> Signed-off-by: Claus Ibsen <[email protected]> --- .../support/batching/KafkaRecordBatchingProcessorFacade.java | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/consumer/support/batching/KafkaRecordBatchingProcessorFacade.java b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/consumer/support/batching/KafkaRecordBatchingProcessorFacade.java index 06745b3506fc..fa26eff2466d 100644 --- a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/consumer/support/batching/KafkaRecordBatchingProcessorFacade.java +++ b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/consumer/support/batching/KafkaRecordBatchingProcessorFacade.java @@ -52,7 +52,14 @@ public class KafkaRecordBatchingProcessorFacade extends AbstractKafkaRecordProce logRecords(allRecords); Set<TopicPartition> partitions = allRecords.partitions(); LOG.debug("Poll received records on {} partitions", partitions.size()); - return kafkaRecordProcessor.processExchange(camelKafkaConsumer, allRecords); + ProcessingResult result = kafkaRecordProcessor.processExchange(camelKafkaConsumer, allRecords); + + // Flush any manually recorded offsets to Kafka for each partition + for (TopicPartition partition : partitions) { + commitManager.commit(partition); + } + + return result; } }
