【发布时间】:2017-11-29 14:26:34
【问题描述】:
我正在编写 kafka Streams 中的 Hopping Window Code,其中 minMaxCalculator() 在流按键分组后计算流内的最小值和最大值。
KTable<Windowed<String>, aggrTest> WinMinMax = Records.groupByKey().aggregate(new aggrTestInitilizer(),
new minMaxCalculator()
, TimeWindows.of(TimeUnit.SECONDS.toMillis(5)).advanceBy(TimeUnit.SECONDS.toMillis(1)),aggrMessageSerde,"aggr-test");
一旦我按键分组,我想并行处理为所有键生成的窗口,即使有一个 kafka 分区。我们应该怎么做?我在哪里可以设置与窗口对应的并行度?
【问题讨论】:
标签: apache-kafka apache-kafka-streams