confluentinc / confluentinc/confluent-kafka-javascript

Bug: #clearCacheAndResetPositions seeks to wrong offset causing message duplication

Open
#417 2 comments 0 reactions 0 assignees View on GitHub
Dominant language
TypeScript
Stars
304
Forks
45
Avg merge
11h 47m
Merged PRs (30d)
5

Description

### Description
There's a bug in _consumer.js where the **#clearCacheAndResetPositions** method seeks to the wrong offset, causing already-processed messages to be reprocessed.

### Bug Location
#### File: consumer.js
#### Method: **#clearCacheAndResetPositions**
The comment on the method clearly states the intended behavior:
```
/* Seek to stored offset for each topic partition. It's possible that we've
* consumed messages upto N from the internalClient, but the user has stale'd the cache
* after consuming just k (< N) messages. We seek back to last consumed offset + 1. */
```
However, the actual implementation does not add +1 to the offset:
```
const lastConsumedOffsets = this.#lastConsumedOffsets.get(key);
const topicPartitionOffsets = [
{
topic: topicPartition.topic,
partition: topicPartition.partition,
offset: lastConsumedOffsets.offset, // Bug: should be lastConsumedOffsets.offset + 1
leaderEpoch: lastConsumedOffsets.leaderEpoch,
}
];
seeks.push(this.#seekInternal(topicPartitionOffsets));
```

### Expected Behavior
When cache expiration triggers #clearCacheAndResetPositions, the consumer should seek to lastConsumedOffset + 1 to avoid reprocessing the last successfully processed message.

### Actual Behavior
The consumer seeks to lastConsumedOffset, causing the last successfully processed message to be consumed and processed again.

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.