【问题标题】:Kafka stream count distinct in scala?卡夫卡流在scala中计数不同?
【发布时间】: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


    【解决方案1】:

    似乎我忘记的只是mutable.HashSet[String]之后的(),所以应该是

    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
          })
        .mapValues((set: mutable.HashSet[String]) => set.size)
        //.count()
      totalUniqueCategoriesCounts.toStream.to("total_unique_categories")
    

    Intellij 完全让我失望了:/

    性能问题仍然存在,非常感谢任何意见。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2017-10-02
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-05-11
      • 2017-02-08
      相关资源
      最近更新 更多