【发布时间】:2018-07-27 10:56:40
【问题描述】:
在你的 kafka 流上应用 groupBy() 或 groupByKey() 时,你会得到一个 KGroupedStream 对象。是否可以在时间窗口内实现该对象,并将分组数据写入主题?
【问题讨论】:
标签: apache-kafka apache-kafka-streams
在你的 kafka 流上应用 groupBy() 或 groupByKey() 时,你会得到一个 KGroupedStream 对象。是否可以在时间窗口内实现该对象,并将分组数据写入主题?
【问题讨论】:
标签: apache-kafka apache-kafka-streams
您可以编写自定义聚合器以将 groupedStream 实现为 KTable。它将具有列表格式的分组记录。以后可以发布到kafka topic。
KTable<K, ArrayList<V>> groupedTable = streamObject.groupByKey().aggregate(
// Custom Initializer
ArrayList::new,
// aggregator
(key, value, list) -> {
list.add(value);
return list;
}, new ArrayListSerde<V>(serdeType), storageName);
groupedTable.through(newTopic);
// Or convert into stream and publish
groupedTable.toStream().to(newTopic);
【讨论】: