如何测试含ContinuousProcessingTimeTrigger的Flink滚动窗口管道?测试无输出
问题:Flink窗口测试无输出排查
尝试编写一个简单的Flink测试,但测试从未触发输出,预期管道会每隔1秒在1秒滚动处理时间窗口内输出聚合后的字符串结果,不清楚代码遗漏了什么,测试中CollectSink.collectedResults始终为空。
测试代码
object CollectSink { val collectedResults = mutable.ListBuffer[(Int, String)]() } class CollectSink extends SinkFunction[(Int, String)] { override def invoke(value: (Int, String), context: SinkFunction.Context): Unit = { CollectSink.collectedResults += value } } class MyTestSuite extends BaseSuite with BeforeAndAfterAll with BeforeAndAfterEach { val flinkCluster = new MiniClusterWithClientResource( new MiniClusterResourceConfiguration.Builder() .setNumberSlotsPerTaskManager(2) .setNumberTaskManagers(1) .build ) override def beforeAll() = { flinkCluster.before() } override def beforeEach() = { CollectSink.collectedResults.clear() } override def afterAll() = { flinkCluster.after() } it should "Simple pipeline emits the expected output" in { val env = StreamExecutionEnvironment.getExecutionEnvironment env.setParallelism(2) val sourceDataStream = env.fromCollection( Seq( (1, "a"), (1, "b"), (2, "c"), (1, "d"), (2, "e") ) ) val collectSink = new CollectSink() // Sample pipeline that emits aggregated strings // It emits the result every second for every window of 1 second sourceDataStream .keyBy(_._1) .window(TumblingProcessingTimeWindows.of(Time.seconds(1))) .trigger(ContinuousProcessingTimeTrigger.of[TimeWindow](Time.seconds(1))) .reduce(new ReduceFunction[(Int, String)] { override def reduce(value1: (Int, String), value2: (Int, String)): (Int, String) = (value1._1, value1._2 + value2._2) }) .addSink(collectSink) env.execute() // Issue here is that it never collects anything somehow, size is always 0 // It seems like the the pipeline never emits anything in this test CollectSink.collectedResults should have size 1 } }
解决方案
核心问题分析
- 处理时间在测试环境无自动推进:Flink的处理时间依赖系统时钟,但MiniCluster测试环境不会自动推进时间。
fromCollection生成的数据流瞬间处理完成,窗口的1秒时间阈值永远无法到达,触发器自然不会触发输出。 - 预期结果数量错误:按key分组后,key=1有3条数据,key=2有2条数据,最终应该生成2个聚合结果,而非测试中预期的1个。
修复步骤
手动推进测试时间
改用异步执行任务后,通过MiniCluster的时间服务手动推进时间,确保窗口和触发器被触发:// 替换原有的env.execute() val jobClient = env.executeAsync() // 推进时间到窗口周期结束后,确保所有计算完成 flinkCluster.getMiniCluster().advanceTime(Time.seconds(2)) jobClient.get()修正结果数量断言
将断言改为匹配实际应该生成的结果数量:CollectSink.collectedResults should have size 2可选:改用事件时间测试(更稳定)
如果希望测试不依赖系统时钟,可切换为事件时间,通过分配时间戳和推进水位线来触发窗口:sourceDataStream .assignTimestampsAndWatermarks(WatermarkStrategy.forMonotonousTimestamps() .withTimestampAssigner((_, timestamp) => timestamp)) .keyBy(_._1) .window(TumblingEventTimeWindows.of(Time.seconds(1))) .trigger(ContinuousEventTimeTrigger.of[TimeWindow](Time.seconds(1))) ...
内容的提问来源于stack exchange,提问作者YHDEV
相关产品推荐
相关产品推荐

