【问题标题】:Extract Seq using Kafka Streams使用 Kafka Streams 提取 Seq
【发布时间】: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


    【解决方案1】:

    使用flatMapValues 代替mapValues

    import scala.collection.JavaConverters.asJavaIterableConverter
    
    builder
      .stream(settings.Streams.inputTopic)
      .flatMapValues(e => fx(e).toIterable.asJava)
      .to(settings.Streams.outputTopic)
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-07-30
      • 2018-12-16
      • 2019-03-30
      • 2020-04-25
      • 2018-02-22
      • 2020-01-30
      • 1970-01-01
      相关资源
      最近更新 更多