如何为Flink窗口函数设置UID?(关联Checkpoint/Savepoint测试)
Flink窗口算子的UID设置方法
要为窗口本身(WindowOperator)设置UID,需在window()方法调用后、process()方法调用前添加.uid()方法,以此指定窗口算子的唯一标识。修改后的代码如下:
DataStream<Alarm> complexAlarmStream = alarmStream.keyBy(Alarm::getOntId) .window(TumblingProcessingTimeWindows.of(Time.minutes(config.get(CheckPointTestingHighlevelFunctionParams.ONT_POWEROFF_DETECTION_WINDOW_TIME_IN_MIN)))) .uid("TumblingWindow-ONT-PowerOff") // 为窗口算子设置UID .process(new CheckPointTestingHighlevelFunction()) .uid("CheckPointTestingFunction") .name("CheckPointTestingFunction");
说明
window()后的.uid():作用于窗口算子,该算子负责窗口的创建、数据归集与生命周期管理,是窗口状态的核心载体。process()后的.uid():作用于你自定义的CheckPointTestingHighlevelFunction对应的业务处理算子,标识业务逻辑相关的状态。
分开设置两类UID,可避免作业迭代时状态ID冲突,保障Checkpoint/Savepoint的兼容性与可恢复性。
内容的提问来源于stack exchange,提问作者Jobs
相关产品推荐
相关产品推荐

