【发布时间】: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