PyFlink窗口聚合不触发:结果流为空问题排查及解决
PyFlink滑动EventTime窗口聚合无输出问题解决
我在使用PyFlink做滑动EventTime窗口聚合时遇到了问题:窗口能正常累积数据(能看到AverageAggregate里的ADDING日志),但始终不会输出聚合结果,结果流为空。我怀疑和窗口触发机制有关,但找不到具体原因。
相关代码如下:
env.set_stream_time_characteristic(TimeCharacteristic.EventTime) class MyTimestampAssigner(TimestampAssigner): def extract_timestamp(self, element, previous_element_timestamp): date_str = element[0].strftime('%Y-%m-%d') # print(element,datetime.strptime(date_str, '%Y-%m-%d').timestamp()) return datetime.strptime(date_str, '%Y-%m-%d').timestamp() joined_data_stream = joined_data_stream.assign_timestamps_and_watermarks( WatermarkStrategy .for_bounded_out_of_orderness(Duration.of_days(1)) .with_timestamp_assigner(MyTimestampAssigner()) ) keyed_stream = joined_data_stream.key_by(lambda x: x[1]) class AverageAggregate(functions.AggregateFunction): def create_accumulator(self) -> [int, int]: return 0, 0 def add(self, value, accumulator): print('ADDING') return accumulator[0]+value[3] , accumulator[1] + 1 def merge(self, a: [int, int], b: [int, int]) -> [int, int]: print('MERGING') return a[0] + b[0], a[1] + b[1] def get_result(self, accumulator): print('GETTING',accumulator[0] / accumulator[1]) return accumulator[0] / accumulator[1] avg_stream = joined_data_stream.key_by(lambda x: x[1]).window(SlidingEventTimeWindows.of(Time.days(2),Time.days(1))) \ .aggregate(AverageAggregate(),accumulator_type=Types.TUPLE([Types.LONG(), Types.LONG()]),output_type=Types.DOUBLE())
问题解决(更新)
问题根源在于时间戳单位不匹配:PyFlink的EventTime时间戳要求是毫秒级整数,但原代码中datetime.strptime(date_str, '%Y-%m-%d').timestamp()返回的是秒级浮点数,导致水位线(Watermark)无法正确推进,窗口永远不会触发关闭和结果输出。
修改时间戳赋值代码,将秒级时间戳转为毫秒级整数即可解决:
原代码片段:
return datetime.strptime(date_str, '%Y-%m-%d').timestamp()
修改后代码片段:
return int(datetime.strptime(date_str, '%Y-%m-%d').timestamp()*1000)
修改后窗口正常触发,聚合结果能正常输出。
内容的提问来源于stack exchange,提问作者Qwetroman
相关产品推荐
相关产品推荐

