【发布时间】:2020-01-13 18:30:46
【问题描述】:
我需要在 kafka 流中对每个用户进行不同的计数。这是我的初始实现,但聚合有错误
Required: [Seq[String], mutableHashSet[String]]
Found: mutable.HashSet[String]
我也不确定如何为 mutable.HashSet 提供自定义 serde ...
val totalUniqueCategoriesCounts: KTable[String, Int] = inputStream
.filter((_ , ev) => ev.eventData.evData.pageType.isDefined)
.groupBy((_, ev) => ev.eventData.custData.customerUid.get)
.aggregate(initializer = mutable.HashSet[String])(
(aggKey: String, newValue: Event, aggValue: mutable.HashSet[String]) => {
val cat: String = newValue.eventData.cntData.contentCategory.get
aggValue += cat
aggValue
}, **Serde Here?**)
.mapValues((set: mutable.HashSet[String]) => set.size)
//.count()
totalUniqueCategoriesCounts.toStream.to("total_unique_categories")
任何帮助将不胜感激。
我也关心性能。这是在 kafka 流中进行不同计数的最佳方法吗?
更新修复了代码问题,但仍然担心此(如果有)或任何更好的方法来做同样的事情的性能影响。
【问题讨论】:
标签: apache-kafka apache-kafka-streams