【问题标题】:How to send OffsetCommitRequest for a consumer in Kafka?如何在 Kafka 中为消费者发送 OffsetCommitRequest?
【发布时间】:2016-07-15 06:18:10
【问题描述】:

我正在尝试在 0.9 版的 Kafka 消费者 API 中使用 OffsetCommitRequest,它包含在以下包中: org.apache.kafka.common.requests.OffsetCommitRequest

如何发送此请求?使用它的理想方法是什么? 我想在 Kafka 本身中提交偏移量。我没有找到任何与 0.9 版相关的文档。其中大部分适用于 0.8.x

此外,此请求的构造函数需要生成 ID、成员 ID 和保留时间。这些字段是什么?

【问题讨论】:

标签: java apache-kafka offset kafka-consumer-api


【解决方案1】:

如果你想手动提交偏移量,也许你应该设置消费者属性

enable.auto.commit=false

并使用 kafka 消费者的 commitSync() 或 commitAsync() 方法。 例如,您可以在处理完所有 ConsumerRecords 后调用 commitSync()。 或者,即使在收到每个 ConsumerRecord 之后,您也可以只提交您想要的 TopicPartition。像这样:

Map<TopicPartition, OffsetAndMetadata> offsetMap = new HashMap<>();
offsetMap.put(new TopicPartition(someTopic, somePartition), new OffsetAndMetadata(someOffset));
kafkaConsumer.commitSync(offsetMap);

【讨论】:

  • 这是我采取的最终解决方案。无论如何,谢谢:) 但我仍然想回答我原来的问题。我认为这个请求对象有时会非常方便。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-08-03
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-09-07
相关资源
最近更新 更多