【发布时间】:2020-09-07 09:40:30
【问题描述】:
我已经参考了几乎所有关于通过处理器 API 的 transform() 或 process() 方法在 KStreams 上记录偏移量的问题,就像这里的许多问题中提到的那样 -
How can I get the offset value in KStream
但我无法得到这些答案的解决方案,所以我问这个问题。
我想在每次消息被流消费时记录分区、消费者组 ID 和偏移量,我不知道如何将 process() 或 transform() 方法与 ProcessorContext API 集成?如果我在我的 CustomParser 类中实现处理器接口,那么我将不得不实现所有方法,但我不确定这是否会起作用,就像在记录元数据的汇合文档中提到的那样 - https://docs.confluent.io/current/streams/developer-guide/processor-api.html#streams-developer-guide-processor-api
我已经在下面的 spring-boot 应用程序中设置了 KStreams(更改变量名以供参考)
@Bean
public Set<KafkaStreams> myKStreamJson(StreamsBuilder profileBuilder) {
Serde<JsonNode> jsonSerde = Serdes.serdeFrom(jsonSerializer, jsonDeserializer);
final KStream<String, JsonNode> pStream = myBuilder.stream(inputTopic, Consumed.with(Serdes.String(), jsonSerde));
Properties props = streamsConfig.kStreamsConfigs().asProperties();
pstream
.map((key, value) -> {
try {
return CustomParser.parse(key, value);
} catch (Exception e) {
LOGGER.error("Error occurred - " + e.getMessage());
}
return new KeyValue<>(null, null);
}
)
.filter((key, value) -> {
try {
return MessageFilter.filterNonNull(key, value);
} catch (Exception e) {
LOGGER.error("Error occurred - " + e.getMessage());
}
return false;
})
.through(
outputTopic,
Produced.with(Serdes.String(), new JsonPOJOSerde<>(TransformedMessage.class)));
return Sets.newHashSet(
new KafkaStreams(profileBuilder.build(), props)
);
}
【问题讨论】:
标签: java spring-boot apache-kafka apache-kafka-streams spring-kafka