[
https://issues.apache.org/jira/browse/CAMEL-13768?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
michael elbaz updated CAMEL-13768:
----------------------------------
Description:
# Provide a way to rewind kafka offset to specific offset (improve seekTo ?)
there is no way to do that using camel-kafka component
for example:
https://blog.sysco.no/integration/kafka-rewind-consumers-offset/
{code:java}
boolean flag = true;
while (true) {
ConsumerRecords<String, String> records = consumer.poll(100);
if(flag) {
Map<TopicPartition, Long> query = new HashMap<>();
query.put(
new TopicPartition("simple-topic-1", 0),
Instant.now().minus(10, MINUTES).toEpochMilli());
// Get offset from timestamp
Map<TopicPartition, OffsetAndTimestamp> result =
consumer.offsetsForTimes(query);
// Rewind offset to previous position using seekTo
result.entrySet()
.stream()
.forEach(entry -> consumer.seek(entry.getKey(),
entry.getValue().offset()));
flag = false;
}
for (ConsumerRecord<String, String> record : records)
System.out.printf("offset = %d, key = %s, value = %s%n",
record.offset(), record.key(), record.value());
}
{code}
# Provide a way to access to kafkaConsumer not only when manualCommit = true
For example if i want to reset rewind if past offset i would like to be able to
call .offsetsForTimes() then seekTo() is currently not possible (maybe just
when manualCommit = true by getting the consumer but is it not convenient)
was:
# Provide a way to rewind kafka offset to specific offset (improve seekTo ?)
# Provide a way to access to kafkaConsumer not only when manualCommit = true
For example if i want to reset rewind if past offset i would like to be able to
call .offsetsForTimes() then seekTo() is currently not possible (maybe just
when manualCommit = true by getting the consumer but is it not convenient)
> SeekTo specific offset and KafkaConsumer access
> ------------------------------------------------
>
> Key: CAMEL-13768
> URL: https://issues.apache.org/jira/browse/CAMEL-13768
> Project: Camel
> Issue Type: New Feature
> Components: camel-kafka
> Affects Versions: 2.24.1
> Reporter: michael elbaz
> Priority: Minor
>
> # Provide a way to rewind kafka offset to specific offset (improve seekTo ?)
> there is no way to do that using camel-kafka component
> for example:
> https://blog.sysco.no/integration/kafka-rewind-consumers-offset/
> {code:java}
> boolean flag = true;
> while (true) {
> ConsumerRecords<String, String> records = consumer.poll(100);
> if(flag) {
> Map<TopicPartition, Long> query = new HashMap<>();
> query.put(
> new TopicPartition("simple-topic-1", 0),
> Instant.now().minus(10, MINUTES).toEpochMilli());
> // Get offset from timestamp
> Map<TopicPartition, OffsetAndTimestamp> result =
> consumer.offsetsForTimes(query);
> // Rewind offset to previous position using seekTo
> result.entrySet()
> .stream()
> .forEach(entry -> consumer.seek(entry.getKey(),
> entry.getValue().offset()));
> flag = false;
> }
> for (ConsumerRecord<String, String> record : records)
> System.out.printf("offset = %d, key = %s, value = %s%n",
> record.offset(), record.key(), record.value());
> }
> {code}
> # Provide a way to access to kafkaConsumer not only when manualCommit = true
> For example if i want to reset rewind if past offset i would like to be able
> to call .offsetsForTimes() then seekTo() is currently not possible (maybe
> just when manualCommit = true by getting the consumer but is it not
> convenient)
--
This message was sent by Atlassian JIRA
(v7.6.14#76016)