Joseph Tam created CAMEL-24208:
----------------------------------

             Summary: pollOnError=ERROR_HANDLER skips valid messages on other 
partitions due to SeekUtil bug
                 Key: CAMEL-24208
                 URL: https://issues.apache.org/jira/browse/CAMEL-24208
             Project: Camel
          Issue Type: Bug
          Components: camel-kafka
    Affects Versions: 4.21.0
            Reporter: Joseph Tam


*{*}Bug Summary:{*}*

When `pollOnError` is set to `ERROR_HANDLER` (the default) and a consumer is 
assigned multiple partitions, a `RecordDeserializationException` on one 
partition causes `SeekUtil` to improperly increment the offset of _every_ 
partition assigned to that consumer, silently skipping valid messages.

*{*}Technical Details:{*}*
If a `RecordDeserializationException` is thrown inside the `poll()` loop, 
`KafkaFetchRecords` catches it globally as an `Exception` and hardcodes the 
`partitionLastOffset` to `-1L`:

```java
// KafkaFetchRecords.java
catch (Exception e) {
long partitionLastOffset = -1L;
this.pollExceptionStrategy.handle(partitionLastOffset, e);
}
```

It then delegates to `BridgeErrorStrategy.handle(long partitionLastOffset, 
Exception exception)`. Because `partitionLastOffset` is `-1L`, 
`BridgeErrorStrategy` passes `-1L` down to `SeekUtil.seekToNextOffset`.

When `SeekUtil` receives `-1L`, it loses track of which partition actually 
failed. To recover, it iterates over `consumer.assignment()` (every single 
partition assigned to the consumer) and arbitrarily advances their fetch 
positions by 1:

```java
// SeekUtil.java
if (tps != null) {
for (TopicPartition tp : tps)

{ long next = consumer.position(tp) + 1L; // ... consumer.seek(tp, next); }

}
```

*{*}Impact:{*}*
If a consumer thread owns multiple partitions (e.g., 4 consumers processing a 
topic with 16 partitions), a bad message on `Partition-0` will cause the 
consumer to seek past the bad message on `Partition-0`, but it will _also_ 
advance the offset of `Partition-1`, `Partition-2`, and `Partition-3` by 1.

If there were valid messages staged at the current offset of those other 
partitions, they are permanently skipped and lost without being processed.

*{*}Proposed Fix:{*}*
`BridgeErrorStrategy` (or `SeekUtil`) should intercept 
`RecordDeserializationException` (or similar exceptions that expose the 
partition) to extract the exact `TopicPartition` and offset provided by the 
Kafka client, and strictly seek _only_ that specific partition forward.

For example, modifying `BridgeErrorStrategy`:

```java
if (partitionLastOffset == -1L && exception instanceof 
RecordDeserializationException) {
RecordDeserializationException rde = (RecordDeserializationException) exception;
TopicPartition tp = rde.topicPartition();
long offset = rde.offset();
if (tp != null && offset >= 0)

{ this.consumer.seek(tp, offset + 1L); return; // skip calling 
SeekUtil.seekToNextOffset }

}
```



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to