Apache Beam使用窗口实现慢更新侧输入运行无输出问题求助
问题根因
- 时间线完全不匹配:主输入元素设置的时间戳为1~99,对应Unix时间戳的1970年区间,而
PeriodicImpulse默认从程序运行的当前系统时间开始生成数据,二者的时间窗口没有任何交集,导致主输入处理时无法匹配到对应的side input窗口数据,cross_join拿到的rights为空,无输出内容。 PeriodicImpulse默认生成无限流,未指定停止时间,程序会一直挂起运行不会自动结束。- side input未显式配置和主输入对齐的窗口策略,窗口匹配逻辑不符合预期。
修正方案
- 给
PeriodicImpulse显式指定匹配主输入时间范围的start_time和stop_time,对齐时间线。 - 给side input显式设置和主输入一致的窗口策略,保证窗口匹配正常。
- 按需调整触发策略保证数据正常输出。
修正后可运行代码
import apache_beam as beam from apache_beam.transforms.periodicsequence import PeriodicImpulse from apache_beam.transforms.window import TimestampedValue from apache_beam.transforms import window def cross_join(left, rights): for x in rights: yield left, x if __name__ == '__main__': data = list(range(1, 100)) pattern = 'pat' main_interval = 10 side_interval = 5 pipeline = beam.Pipeline() side_input = ( pipeline # 显式指定start和stop时间,对齐主输入的1~99时间范围 | 'PeriodicImpulse' >> PeriodicImpulse( start_time=0, stop_time=100, fire_interval=side_interval, apply_windowing=True) | 'MapToFileName' >> beam.Map(lambda x: pattern + str(x)) # 显式设置和主输入一致的窗口策略 | 'WindowSideInto' >> beam.WindowInto(window.FixedWindows(main_interval)) ) main_input = ( pipeline | 'MpImpulse' >> beam.Create(data) | 'MapMpToTimestamped' >> beam.Map(lambda src: TimestampedValue(src, src)) | 'WindowMpInto' >> beam.WindowInto(window.FixedWindows(main_interval))) result = ( main_input | 'ApplyCrossJoin' >> beam.FlatMap(cross_join, rights=beam.pvalue.AsIter(side_input)) | 'log' >> beam.Map(print)) res = pipeline.run() res.wait_until_finish()
内容的提问来源于stack exchange,提问作者Bảo Gia
相关产品推荐
相关产品推荐

