【发布时间】:2017-11-15 13:51:11
【问题描述】:
我正在尝试读取一个 Kafka 主题,在此基础上进行一些处理,并将结果存储在另一个主题中。
我的代码如下所示:
builder
.stream(settings.Streams.inputTopic)
.mapValues[Seq[Product]]((e: EventRecord) ⇒ fx(e))
// Something needs to be done here...
.to(settings.Streams.outputTopic)
fx(e) 函数进行一些处理并返回一个Seq[Product]。我想将所有产品作为单独的条目存储在主题中。问题是从主题读取的消息包含多个产品,因此返回值为 fx(e)。
是否可以将此行为嵌入到流中?
【问题讨论】:
标签: scala apache-kafka-streams