Apache Beam流式处理:如何基于事件码启停交易数据聚合?
问题描述
问题背景
我有来自单个PubSub的流式数据,将其拆分为两组:Group A为单事件码组,Group B为交易数据组。我需要在Group A出现特定事件码后聚合交易数据,例如在Group A出现事件码U时开始聚合,出现事件码B时结束聚合。
疑问
如何让Apache Beam管道根据Group A中的事件码判断聚合的启停时机?
我已编写一个示例管道,但未得到预期结果。假设我有两个PCollection,分别对应事件码和交易数据:
from datetime import datetime import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions beam_options = PipelineOptions() with beam.Pipeline(options=beam_options) as p: event = (p | "Create1" >> beam.Create([ {"eventCode": "A", "dates": datetime(2022, 12, 1, 7, 0, 0)}, {"eventCode": "C", "dates": datetime(2022, 12, 1, 7, 1, 0)}, {"eventCode": "U", "dates": datetime(2022, 12, 1, 7, 2, 0)}, {"eventCode": "D", "dates": datetime(2022, 12, 1, 7, 3, 0)}, {"eventCode": "E", "dates": datetime(2022, 12, 1, 7, 4, 0)}, {"eventCode": "F", "dates": datetime(2022, 12, 1, 7, 5, 0)}, {"eventCode": "G", "dates": datetime(2022, 12, 1, 7, 6, 0)}, {"eventCode": "B", "dates": datetime(2022, 12, 1, 7, 7, 0)}, {"eventCode": "T", "dates": datetime(2022, 12, 1, 7, 8, 0)}, {"eventCode": "H", "dates": datetime(2022, 12, 1, 7, 9, 0)}, {"eventCode": "I", "dates": datetime(2022, 12, 1, 7, 10, 0)}, {"eventCode": "J", "dates": datetime(2022, 12, 1, 7, 11, 0)}, {"eventCode": "M", "dates": datetime(2022, 12, 1, 7, 12, 0)}, {"eventCode": "B", "dates": datetime(2022, 12, 1, 7, 14, 0)}, {"eventCode": "Y", "dates": datetime(2022, 12, 1, 7, 15, 0)}, {"eventCode": "X", "dates": datetime(2022, 12, 1, 7, 16, 0)}, ])) trx_data = (p | "Create2" >> beam.Create([ {"trxCode": "TRX001", "price": 156, "dates": datetime(2022, 12, 1, 7, 1, 0)}, {"trxCode": "TRX002", "price": 157, "dates": datetime(2022, 12, 1, 7, 2, 0)}, {"trxCode": "TRX003", "price": 158, "dates": datetime(2022, 12, 1, 7, 3, 0)}, {"trxCode": "TRX004", "price": 159, "dates": datetime(2022, 12, 1, 7, 4, 0)}, {"trxCode": "TRX005", "price": 160, "dates": datetime(2022, 12, 1, 7, 5, 0)}, {"trxCode": "TRX001", "price": 161, "dates": datetime(2022, 12, 1, 7, 6, 0)}, {"trxCode": "TRX002", "price": 162, "dates": datetime(2022, 12, 1, 7, 7, 0)}, {"trxCode": "TRX006", "price": 163, "dates": datetime(2022, 12, 1, 7, 8, 0)}, {"trxCode": "TRX001", "price": 164, "dates": datetime(2022, 12, 1, 7, 9, 0)}, {"trxCode": "TRX007", "price": 165, "dates": datetime(2022, 12, 1, 7, 10, 0)}, {"trxCode": "TRX008", "price": 166, "dates": datetime(2022, 12, 1, 7, 11, 0)}, {"trxCode": "TRX003", "price": 167, "dates": datetime(2022, 12, 1, 7, 12, 0)}, {"trxCode": "TRX005", "price": 168, "dates": datetime(2022, 12, 1, 7, 13, 0)}, {"trxCode": "TRX009", "price": 169, "dates": datetime(2022, 12, 1, 7, 14, 0)}, {"trxCode": "TRX010", "price": 170, "dates": datetime(2022, 12, 1, 7, 15, 0)}, ])) # 创建窗口模拟流处理过程 event_window = (event | beam.Map(lambda d: beam.window.TimestampedValue(d, d["dates"].timestamp())) | beam.WindowInto(beam.window.FixedWindows(1)) | beam.Map(lambda d: (d["eventCode"], d)) | beam.GroupByKey() | beam.MapTuple(lambda e, d: d) ) trx_window = (trx_data | "TrxData1" >> beam.Map(lambda d: beam.window.TimestampedValue(d, d["dates"].timestamp())) | "TrxData2" >> beam.WindowInto(beam.window.FixedWindows(10*60)) | "TrxData3" >> beam.Map(lambda d: (d["secCode"], d)) | "TrxData4" >> beam.GroupByKey() ) process = (trx_window | beam.ParDo(TransformWithSideInput(), stock_u=beam.pvalue.List(event_window)) | beam.Map(print) )
针对TransformWithSideInput(),我编写了如下代码:
class TransformWithSideInput(beam.DoFn): def start_bundle(self): self.bucket = [] def process(self, element, event): for ethe in event: event_code = e.get("eventCode") event_code_dates = e.get("dates") queue = [] if event_code == "U": queue = ["U"] elif event_code == "B": queue = ["B"] if queue == ["U"]: # 如果存在U码,则开始将数据存入bucket # 逻辑上,需要处理的交易数据为事件码到达时间之后的所有交易 if element.get("dates") >= event_code_dates: self.bucket.append(element) elif queue == ["B"]: yield self.bucket def finish_bundle(self): if self.bucket: self.flush() def flush(self): yield self.bucket
我将事件管道作为side input传入TransformWithSideInput,目的是检查事件条件,在满足条件时返回对应数据。但目前交易管道无任何输出。
需求为:当event_code == "U"时启动交易管道,直到event_code == "B"出现时停止并返回结果。
解决方案
现有代码的核心问题
- 侧输入使用错误:
beam.pvalue.List侧输入是静态的,只会在管道启动时加载一次事件数据,无法处理流式到达的新事件,导致无法实时响应U/B事件的启停信号。 - 数据字段错误:交易数据中没有
secCode字段,trx_window中的GroupByKey会导致所有交易数据丢失,无输出。 - DoFn逻辑错误:变量名拼写错误(
ethe应为e),且element是GroupByKey后的键值对,不是单个交易数据,无法直接获取dates字段。 - 状态管理错误:
start_bundle中的bucket仅在当前bundle有效,流式处理中需要跨bundle维护聚合状态,必须使用Beam的State API。 - 窗口设置不合理:FixedWindows不适合这种基于事件触发的聚合场景,应该使用全局窗口结合触发器。
正确实现思路
由于事件和交易数据来自同一个PubSub,我们可以将两者合并到同一个PCollection,按时间戳排序后,使用Beam的State和Timer API来维护聚合状态:
- 用布尔状态标记是否处于聚合模式(
is_aggregating) - 用列表状态存储当前聚合的交易数据(
aggregated_trxs) - 处理每个元素时:
- 若为
U事件,开启聚合模式 - 若为
B事件,关闭聚合模式并输出当前聚合的交易数据,随后清空状态 - 若为交易数据且处于聚合模式,将其加入聚合状态
- 若为
完整代码实现
from datetime import datetime import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions from apache_beam.transforms.userstate import ReadModifyWriteStateSpec, BagStateSpec class AggregateTrxsOnEvents(beam.DoFn): # 定义状态:标记是否开启聚合 IS_AGGREGATING = ReadModifyWriteStateSpec('is_aggregating', bool) # 定义状态:存储聚合的交易数据 AGGREGATED_TRXS = BagStateSpec('aggregated_trxs', beam.coders.StrUtf8Coder()) def process(self, element, is_aggregating=beam.DoFn.StateParam(IS_AGGREGATING), aggregated_trxs=beam.DoFn.StateParam(AGGREGATED_TRXS)): # 区分事件和交易数据 if 'eventCode' in element: event_code = element['eventCode'] if event_code == 'U': # 开启聚合 is_aggregating.write(True) elif event_code == 'B': # 关闭聚合并输出结果 is_aggregating.write(False) trxs_list = list(aggregated_trxs.read()) if trxs_list: yield { 'start_time': element['dates'].strftime('%Y-%m-%d %H:%M:%S'), 'aggregated_trxs': trxs_list } # 清空聚合状态 aggregated_trxs.clear() else: # 处理交易数据:仅当聚合开启时存储 if is_aggregating.read(): # 将交易数据序列化为字符串存储(可根据需求使用更高效的编码) trx_str = f"{element['trxCode']},{element['price']},{element['dates'].strftime('%Y-%m-%d %H:%M:%S')}" aggregated_trxs.add(trx_str) def run(): beam_options = PipelineOptions(streaming=True) # 启用流式模式 with beam.Pipeline(options=beam_options) as p: # 合并事件和交易数据,添加时间戳并标记类型 events_with_ts = ( p | "CreateEvents" >> beam.Create([ {"eventCode": "A", "dates": datetime(2022, 12, 1, 7, 0, 0)}, {"eventCode": "C", "dates": datetime(2022, 12, 1, 7, 1, 0)}, {"eventCode": "U", "dates": datetime(2022, 12, 1, 7, 2, 0)}, {"eventCode": "D", "dates": datetime(2022, 12, 1, 7, 3, 0)}, {"eventCode": "E", "dates": datetime(2022, 12, 1, 7, 4, 0)}, {"eventCode": "F", "dates": datetime(2022, 12, 1, 7, 5, 0)}, {"eventCode": "G", "dates": datetime(2022, 12, 1, 7, 6, 0)}, {"eventCode": "B", "dates": datetime(2022, 12, 1, 7, 7, 0)}, {"eventCode": "T", "dates": datetime(2022, 12, 1, 7, 8, 0)}, {"eventCode": "H", "dates": datetime(2022, 12, 1, 7, 9, 0)}, {"eventCode": "I", "dates": datetime(2022, 12, 1, 7, 10, 0)}, {"eventCode": "J", "dates": datetime(2022, 12, 1, 7, 11, 0)}, {"eventCode": "M", "dates": datetime(2022, 12, 1, 7, 12, 0)}, {"eventCode": "B", "dates": datetime(2022, 12, 1, 7, 14, 0)}, {"eventCode": "Y", "dates": datetime(2022, 12, 1, 7, 15, 0)}, {"eventCode": "X", "dates": datetime(2022, 12, 1, 7, 16, 0)}, ]) | "EventTimestamp" >> beam.Map(lambda x: beam.window.TimestampedValue(x, x["dates"].timestamp())) ) trxs_with_ts = ( p | "CreateTrxs" >> beam.Create([ {"trxCode": "TRX001", "price": 156, "dates": datetime(2022, 12, 1, 7, 1, 0)}, {"trxCode": "TRX002", "price": 157, "dates": datetime(2022, 12, 1, 7, 2, 0)}, {"trxCode": "TRX003", "price": 158, "dates": datetime(2022, 12, 1, 7, 3, 0)}, {"trxCode": "TRX004", "price": 159, "dates": datetime(2022, 12, 1, 7, 4, 0)}, {"trxCode": "TRX005", "price": 160, "dates": datetime(2022, 12, 1, 7, 5, 0)}, {"trxCode": "TRX001", "price": 161, "dates": datetime(2022, 12, 1, 7, 6, 0)}, {"trxCode": "TRX002", "price": 162, "dates": datetime(2022, 12, 1, 7, 7, 0)}, {"trxCode": "TRX006", "price": 163, "dates": datetime(2022, 12, 1, 7, 8, 0)}, {"trxCode": "TRX001", "price": 164, "dates": datetime(2022, 12, 1, 7, 9, 0)}, {"trxCode": "TRX007", "price": 165, "dates": datetime(2022, 12, 1, 7, 10, 0)}, {"trxCode": "TRX008", "price": 166, "dates": datetime(2022, 12, 1, 7, 11, 0
相关产品推荐
相关产品推荐

