【发布时间】:2021-01-18 03:59:14
【问题描述】:
我需要使用多个转换器解析 kafka 上的复杂消息。每个转换器解析消息的一部分并通过在消息上填充一些属性来编辑消息。最后,完全解析的消息使用 Kafka 消费者存储在数据库中。目前,我正在这样做:
streamsBuilder.stream(Topic.A, someConsumer)
\\ filters messages that have unparsed parts of type X
.filter(filterX)
\\ transformer that edits the message and produces new Topic.E messages
.transform(ParseXandProduceE::new)
.to(Topic.A, someProducer)
streamsBuilder.stream(Topic.A, someConsumer)
\\ filters messages that have unparsed parts of type Y
.filter(filterY)
\\ transformer that edits the message and produces new Topic.F messages
.transform(ParseYandProduceF::new)
.to(Topic.A, someProducer)
Transformer 看起来像:
class ParseXandProduceE implements Transformer<...> {
@Override
public KeyValue<String, Message> transform (String key, Message message) {
message.x = parse(message.rawX);
context.forward(newKey, message.x, Topic.E);
return KeyValue.pair(key, message);
}
}
但是,这很麻烦,相同的消息会多次通过这些流。
此外,还有一个消费者在数据库中存储topic.A 的消息。当前,消息在每次转换之前和每次转换之后都会存储多次。每条消息都需要存储一次。
以下方法可行,但似乎不利,因为每个过滤器+转换块都可以干净地放在自己单独的类中:
streamsBuilder.stream(Topic.A, someConsumer)
\\ transformer that filters and edits the message and produces new Topic.E + Topic.F messages
.transform(someTransformer)
.to(Topic.B, someProducer)
并让持久化消费者收听Topic.B。
后一种建议的解决方案是可行的方法,还是有其他方法可以达到相同的结果?也许有源和接收器的完整拓扑配置?如果是这样,在这种情况下会是什么样子?
【问题讨论】:
标签: java apache-kafka apache-kafka-streams