【问题标题】:Kafka Streams Hopping window top N by dimensionKafka Streams Hopping window top N by dimension
【发布时间】:2020-05-04 06:04:11
【问题描述】:

我有一个 kafka 流,我需要一个执行以下操作的处理器:

使用 45 秒跳跃窗口和 5 秒提前计算基于域对象的一维的前 5 个计数。例如,如果流包含 Clickstream 数据,我需要按域名查看的前 5 个 url,但也需要在跳跃窗口中进行窗口化。

我见过一些做窗口计数的例子,例如:

KStream<String, GenericRecord> pageViews = ...;

// Count page views per window, per user, with hopping windows of size 5 minutes that advance every 1 minute
KTable<Windowed<String>, Long> windowedPageViewCounts = pageViews
    .groupByKey(Grouped.with(Serdes.String(), genericAvroSerde))
    .windowedBy(TimeWindows.of(Duration.ofMinutes(5).advanceBy(Duration.ofMinutes(1))))
    .count()

MusicExample 上的 Top n 聚合,例如:

songPlayCounts.groupBy((song, plays) ->
            KeyValue.pair(TOP_FIVE_KEY,
                new SongPlayCount(song.getId(), plays)),
        Grouped.with(Serdes.String(), songPlayCountSerde))
        .aggregate(TopFiveSongs::new,
            (aggKey, value, aggregate) -> {
              aggregate.add(value);
              return aggregate;
            },
            (aggKey, value, aggregate) -> {
              aggregate.remove(value);
              return aggregate;
            },
            Materialized.<String, TopFiveSongs, KeyValueStore<Bytes, byte[]>>as(TOP_FIVE_SONGS_STORE)
                .withKeySerde(Serdes.String())
                .withValueSerde(topFiveSerde)
        );

我似乎无法将 2 结合起来 - 我同时获得窗口和前 n 个聚合。有什么想法吗?

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    通常是的,但是,对于非窗口 top-N 聚合,算法将总是是一个近似值(不可能得到精确的结果,因为需要缓冲 一切 无限输入是不可能的)。但是,对于跳跃窗口,您需要进行精确计算。

    对于窗口情况,实际的聚合步骤可能只是累积每个窗口的所有输入记录(例如,返回List&lt;V&gt; 或其他一些集合)。在此结果 KTable 上,您应用 mapValues() 函数获取每个窗口(和键)的输入记录 List&lt;V&gt;,并可以计算您正在寻找的实际前 N 个结果。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-10-23
      • 2010-10-19
      • 1970-01-01
      • 2019-01-18
      • 1970-01-01
      • 2020-04-01
      相关资源
      最近更新 更多