【问题标题】:Kafka Streams : Flushing intermediate Windowed results as commit interval and window time are not in syncKafka Streams:刷新中间窗口结果,因为提交间隔和窗口时间不同步
【发布时间】:2020-01-21 05:54:26
【问题描述】:

kafka 流的配置:

threads = 1;
replicationFactor = 1;
ktableCommitInterval= 10000;
ktableMemory=72000000;
timeDuration=10;

拓扑:

KStream<Windowed<String>,String> windowedStringKStream =
                         streamsBuilder.stream(inputTopic, Consumed.with(Serdes.String(),Serdes.String()))
.groupByKey(Grouped.with(Serdes.String(),Serdes.String()))
.windowedBy(TimeWindows.of(Duration.ofSeconds(timeDuration)).grace(Duration.ofSeconds(0)))
                        .reduce(Numners::append,Materialized.<String, String, WindowStore<Bytes,byte[]>>as(storeName).withCachingEnabled().withRetention(Duration.ofSeconds(timeDuration)).withKeySerde(Serdes.String()).withValueSerde(Serdes.String()))
.toStream();
       
 

代码说明:

Code appends numbers in a 10 Second window. Incremental Number records are sent exactly at a interval of 1 second into input topic.

输出:

问题:

提交间隔设置为 10 秒。缓存大小设置为 72 MB。数据以字节为单位。状态存储已启用缓存。文档指出,将数据推送到下游的 kafka 流的操作语义取决于缓存大小或提交间隔,无论首先发生什么。但根据实验,提交在一分钟内发生两次。观察是提交间隔在应用程序启动时开始,但在数据开始时开始窗口。如图所示,中间窗口结果被推送,最终窗口结果也被推送。

对于我正在处理的用例,无法使用 Suppress(),因为如果没有关于该主题的新数据,它将不会刷新数据。

任何帮助将不胜感激。如果有人遇到这种情况或想重现这种情况,请告诉我。

【问题讨论】:

    标签: java apache-kafka apache-kafka-streams windowing


    【解决方案1】:

    唯一的解决方案是构建一个自定义的transform(),它实现了一个自定义版本的“抑制”,即使没有新的输入数据到达,您也可以发出数据(例如,挂钟时间标点可能有助于实现它)。

    目前(从 Apache Kafka 2.4 版本开始),没有内置支持。如果您不使用suppress(),则窗口聚合可能总是在窗口实际关闭之前发出一些中间结果,并且无法以不同方式配置Kafka Streams。

    【讨论】:

    • 感谢您的回复。我有一个改进抑制融合文档的请求。 [链接] (confluent.io/blog/kafka-streams-take-on-watermarks-and-triggers) 文章没有强调流时间是每个分区而不是每个主题。给新手造成误解。由于 Stream 线程不共享任何状态,因此可以隐式理解,但我们可以使其显式。我还使用高级 dsl 实现了转换。面临的问题 -> link
    猜你喜欢
    • 1970-01-01
    • 2019-10-15
    • 2018-08-21
    • 2016-12-20
    • 1970-01-01
    • 2019-07-29
    • 2022-12-08
    • 1970-01-01
    相关资源
    最近更新 更多