【发布时间】:2021-12-16 11:31:51
【问题描述】:
我正在尝试测量具有窗口操作的 Flink 应用程序的延迟,如下所示:
SingleOutputStreamOperator<String> branch = stream
.getSideOutput(outputTag2)
.keyBy(MetricObject::getRootAssetId)
.window(TumblingEventTimeWindows.of(Time.seconds(60)))
.trigger(ContinuousEventTimeTrigger.of(Time.seconds(15)))
.aggregate(new CountDistinctAggregate(), new CountDistinctProcess())
.name("windowed-count-distinct")
.uid("windowed-count-distinct")
.map((value)->String.valueOf(value.getTimestamp().toEpochMilli()))
.name("send-timestamp");
我正在考虑事件时间并使用此水印策略提取时间戳:
.<SingleRecord>forBoundedOutOfOrderness(Duration.ofSeconds(15))
.withTimestampAssigner((event, timestamp) -> event.getTimestamp().toEpochMilli()))
聚合函数将特定对象保存为累加器,其中还包含提取的时间戳;这些时间戳写在一个 kafka 主题中。问题是返回的时间戳是这些:
1639651859988
1639651890163
1639651904900
1639651919728
1639651919728
1639651949973
1639651965085
1639651979870
返回的时间戳并不像我预期的那样等距,第四个和第五个是相等的,但是它们以 15 秒的间隔返回,这是不可能的,因为应用程序记录的输入是每秒连续生成的(每秒 10 个)。在其他测试中,我也遇到了更糟糕的情况:
1639651979870
1639651992771
1639651992771
1639651992771
1639651992771
1639652189791
1639652205001
1639652219876
奇怪的事实是,当我使用没有触发器的简单翻转窗口时:
.window(TumblingEventTimeWindows.of(Time.seconds(15)))
返回的时间戳与预期的一样间隔:
1639652429766
1639652444930
1639652459900
1639652474609
1639652489746
1639652504862
1639652519734
1639652534847
我真的不明白问题出在哪里,聚合函数中的累加器似乎没有正确升级。
【问题讨论】:
标签: java apache-flink flink-streaming