PyFlink 1.15滑动窗口报错:CountSlidingWindowAssigner无get_default_trigger属性
解决PyFlink 1.15中滑动计数窗口的
AttributeError问题 这个错误是因为在Flink 1.15版本中,直接实例化CountSlidingWindowAssigner并传入window()方法时,缺少了必要的触发器配置——而1.16版本对这个API做了封装,会自动帮你处理触发器的绑定。
修正方案:使用count_window()快捷方法
在1.15的PyFlink DataStream API中,创建滑动计数窗口最简洁且正确的方式是使用count_window(size, slide)这个封装好的方法,它会自动为窗口配置合适的CountTrigger,避免手动实例化分配器带来的问题。
修正后的代码如下:
input_stream = self.table_env.to_data_stream(flat_table) result_stream = input_stream.key_by(get_key) \ .count_window(self.window_size, self.window_slide) \ .apply(CTATThresholds(), result_type=Types.TUPLE([Types.STRING(), Types.STRING(), Types.STRING(), Types.STRING(), Types.STRING(), Types.FLOAT(), Types.FLOAT(), Types.FLOAT()]))
为什么之前的代码会报错?
在Flink 1.15中,CountSlidingWindowAssigner只是窗口分配器,它负责定义窗口的范围,但还需要搭配触发器(Trigger)来决定何时触发窗口计算。count_window()方法内部已经帮你完成了窗口分配器和CountTrigger的绑定,而你手动实例化CountSlidingWindowAssigner时,没有指定触发器,就会抛出get_default_trigger的属性错误。
如果你确实需要手动使用CountSlidingWindowAssigner(不推荐,除非有特殊需求),可以补充触发器的配置,代码如下:
from pyflink.datastream.trigger import CountTrigger input_stream = self.table_env.to_data_stream(flat_table) result_stream = input_stream.key_by(get_key) \ .window(CountSlidingWindowAssigner(self.window_size, self.window_slide)) \ .trigger(CountTrigger.of(self.window_size)) \ .apply(CTATThresholds(), result_type=Types.TUPLE([Types.STRING(), Types.STRING(), Types.STRING(), Types.STRING(), Types.STRING(), Types.FLOAT(), Types.FLOAT(), Types.FLOAT()]))
不过显然第一种方案更简洁,也符合Flink API的设计意图。
内容的提问来源于stack exchange,提问作者Michael Boesl
相关产品推荐
相关产品推荐

