【发布时间】:2017-07-14 07:52:19
【问题描述】:
我正在使用 Kafka 流来使用跳跃时间窗口计算过去 3 分钟内发生了多少事件:
public class ViewCountAggregator {
void buildStream(KStreamBuilder builder) {
final Serde<String> stringSerde = Serdes.String();
final Serde<Long> longSerde = Serdes.Long();
KStream<String, String> views = builder.stream(stringSerde, stringSerde, "streams-view-count-input");
KStream<String, Long> viewCount = views
.groupBy((key, value) -> value)
.count(TimeWindows.of(TimeUnit.MINUTES.toMillis(3)).advanceBy(TimeUnit.MINUTES.toMillis(1)))
.toStream()
.map((key, value) -> new KeyValue<>(key.key(), value));
viewCount.to(stringSerde, longSerde, "streams-view-count-output");
}
public static void main(String[] args) throws Exception {
// some not so important initialization code
...
}
}
当运行消费者并将一些消息推送到输入主题时,随着时间的推移,它会收到以下更新:
single 1
single 1
single 1
five 1
five 4
five 5
five 4
five 1
这几乎是正确的,但它从未收到以下更新:
single 0
five 0
如果没有它,我更新计数器的消费者永远不会在较长时间没有事件时将其设置回零。我希望消费的消息看起来像这样:
single 1
single 1
single 1
single 0
five 1
five 4
five 5
five 4
five 1
five 0
是否有一些我遗漏的配置选项/参数可以帮助我实现这种行为?
【问题讨论】:
标签: java apache-kafka apache-kafka-streams