Flink KeyedProcessFunction跨key共用普通变量是否存在并发修改问题
问题1解答
对,该判断成立。你代码中定义的timestamp是KeyedProcessFunction实例的普通成员变量,Flink每个算子子任务(TaskManager上运行的算子并发实例)对应一个该类的实例,该实例的普通成员变量会被当前子任务处理的所有key共享。只要任意一个被处理的事件触发了currentHour.after(timestamp)的判断逻辑,timestamp就会被更新为新的小时值,所有后续处理的key都会共用这个更新后的值。
注意:该逻辑仅在你使用处理时间做小时判断时符合预期,如果你基于事件时间做业务逻辑判断,直接取系统时间截断小时的逻辑会和事件时间的进度不匹配,存在逻辑错误。
问题2解答
不会出现并发修改问题。Flink的算子实例处理元素的逻辑是严格单线程执行的,同一个算子子任务的processElement方法不会被多线程同时调用,因此不存在多线程同时读写timestamp变量的场景,不会有并发安全问题。
额外注意事项
- 你当前定义的
timestamp是普通堆变量,没有纳入Flink的状态管理,如果作业发生故障重启、扩缩容,该变量会被重新初始化,不会保留故障前的取值,如果你需要故障恢复后该值的逻辑一致性,可以改用**算子状态(Operator State)**存储该变量。
内容的提问来源于stack exchange,提问作者sparkless
相关产品推荐
相关产品推荐

