【问题标题】:How can I get the offset value in KStream如何在 KStream 中获取偏移值
【发布时间】:2017-04-09 23:51:12
【问题描述】:

我正在使用 Kafka Streams 开发 PoC。现在我需要获取流消费者中的偏移值,并使用它为每条消息生成一个唯一键(topic-offset)->hash。原因是:生产者是 syslog,只有少数有 ID。我无法在消费者中生成 UUID,因为在重新处理的情况下我需要重新生成相同的密钥。

我的问题是:org.apache.kafka.streams.processor.ProcessorContext 类公开了一个返回值的 .offset() 方法,但我使用的是 KStream 而不是处理器,我找不到返回相同内容的方法。

有人知道如何从 KStream 中提取每一行的消费者值吗? 提前致谢。

【问题讨论】:

标签: java apache-kafka apache-kafka-streams


【解决方案1】:

您可以通过process(...)transform(...)transformValues(...) 使用混合匹配DSL 和处理器API。

它允许您访问类似于普通处理器 API 的当前记录偏移量。在您的情况下,您似乎想使用KStream#transform(...)

【讨论】:

  • 即使使用处理器 API,我们也只能访问键和值,而不能访问 ConsumerRecord,只有时间戳提取器似乎有 ConsumerRecord。能否请您添加有关访问 partitionId、时间戳、偏移量和其他字段的其他详细信息。
  • 如问题中所述,Processor API 通过Processor#init(...) 方法提供了一个ProcessorContext 对象。 ProcessorContext 会在每次调用 process() 之前使用将要处理的下一条记录的元数据进行更新。因此,当process() 被调用时,您可以通过调用相应的ProcessorContext 方法来获取记录偏移量等。只需在使用init() 中提供的ProcessorContext 初始化的类成员变量中保留对它的引用。
  • 谢谢! ProcessorContext 拥有所需的一切。
猜你喜欢
  • 1970-01-01
  • 2021-08-16
  • 2023-03-26
  • 1970-01-01
  • 2015-06-16
  • 2011-07-26
  • 2021-08-30
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多