You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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"出现时停止并返回结果。


解决方案

现有代码的核心问题

  1. 侧输入使用错误:beam.pvalue.List侧输入是静态的,只会在管道启动时加载一次事件数据,无法处理流式到达的新事件,导致无法实时响应U/B事件的启停信号。
  2. 数据字段错误:交易数据中没有secCode字段,trx_window中的GroupByKey会导致所有交易数据丢失,无输出。
  3. DoFn逻辑错误:变量名拼写错误(ethe应为e),且element是GroupByKey后的键值对,不是单个交易数据,无法直接获取dates字段。
  4. 状态管理错误:start_bundle中的bucket仅在当前bundle有效,流式处理中需要跨bundle维护聚合状态,必须使用Beam的State API。
  5. 窗口设置不合理: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
相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.08 00:30:54