Flink KeyedBroadcastProcessFunction的EventTime定时器在TestHarness中不触发
问题原因分析
Flink 多输入算子的当前水印由所有输入流的水印最小值决定,你使用的KeyedBroadcastProcessFunction包含普通数据流和广播流两个输入,你只推进了普通数据流的水印,广播流的水印仍保持初始值Long.MIN_VALUE,因此整个算子的有效水印始终无法上涨,自然不会触发事件时间定时器。
修复方案
按以下步骤调整测试代码即可:
- 取消测试代码中广播流水印推进逻辑的注释,将广播流的水印推进到和主流水印相同的高度,保证两者都大于你注册的定时器触发时间
- 如果不需要在测试中验证广播流的水印逻辑,也可以直接给广播流设置最大值
Long.MAX_VALUE,让算子水印完全由普通数据流决定
修改后的测试代码示例
@Test public void evaluateFormular_ShouldSumOnTimer() throws Exception { long minutesToWait = 1; var definition = createTestCondition("Test", String.format("%s(_var)", method), "_var", "Test", "1", minutesToWait); var message = new CalculationControlMessage(); message.setAction(ControlMessageAction.Create); message.setCalculationDefinition(definition); harness.processBroadcastElement(message, 100l); this.processValues(harness, values); // 同时推进广播流和普通流的水印 long triggerWatermark = Time.minutes(minutesToWait).toMilliseconds() + 1; harness.processBroadcastWatermark(triggerWatermark); harness.processWatermark(triggerWatermark); assertEquals(harness.numEventTimeTimers(), 1); assertEquals("there should be a formular evaluated", 1, harness.extractOutputValues().size()); harness.extractOutputValues().forEach(datapoint -> { assertEquals(datapoint.getValue(), expected, 0d); }); }
额外检查项
如果调整后仍有问题,可以再排查两个点:
- 调用
processValues处理普通流元素时,是否为每个DataPointEvent正确设置了早于触发水印的事件时间戳 - 你在业务代码中注册的事件时间定时器的时间确实小于最终推进的水印值
内容的提问来源于stack exchange,提问作者emilio
相关产品推荐
相关产品推荐

