【问题标题】:KStream to KTableKStream 转 KTable
【发布时间】:2018-07-23 07:08:50
【问题描述】:
@StreamListener("input")
@SendTo("output")
public KStream<?, MyObject> process(KStream<Object, IncomingObject> input) {
KTable table = input.flatMapValues(value -> this.getMylogic(value));            
return table.toStream();
}

我正在尝试将 KStream 转换为 KTable,然后再转换为 KStream,但我得到 无法从 KStream 转换为 KTable

值是 json。请帮忙,我如何也可以使用聚合?

{
"name":"test",
address{
"localAddress":"myaddress",
"businessAddress":"testAddress"
}
}

在 mylogic 方法中,我只将地址发送到另一个主题。 请帮忙

【问题讨论】:

  • 方法 flatMapValues 返回 KStream 而不是 KTable
  • 是的,那么如何从kstream转换为ktable?
  • 你还没有写出你为什么需要KTable的目的。要获取 KTable,您可以使用以下内容:kStream.groupByKey().aggregate(..)
  • 创建 Ktable 的目的是什么?您能否在 ktable 中提供预期的输出详细信息?

标签: java apache-kafka apache-kafka-streams


【解决方案1】:

您需要应用聚合函数,例如 count 来获取 KTable 作为结果。否则没有办法做 KStream -> KTable -> KStream。

您需要并且可以做的是 KStream.count()(例如)-> KTable -> KStream。所以基本上聚合的结果会发布到 KStream,也可能会发布到 Kafka 主题中。

【讨论】:

    猜你喜欢
    • 2018-02-23
    • 2020-10-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-12-02
    相关资源
    最近更新 更多