Apache Beam 统一批流处理

FreeGuideOnline 最新 2026-07-10

bash pip install apache-beam


创建一个简单的单词计数流水线,体验 Beam 的代码结构。

### 第一条 Pipeline:WordCount

```python
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions

# 创建流水线选项,此处使用默认配置
options = PipelineOptions()

# 构建流水线
with beam.Pipeline(options=options) as p:
    # 步骤1:读取文本文件(有界数据)
    lines = p | 'ReadLines' >> beam.io.ReadFromText('input.txt')
    
    # 步骤2:拆分单词
    words = lines | 'ExtractWords' >> beam.FlatMap(lambda line: line.split())
    
    # 步骤3:键值对计数
    counts = (
        words
        | 'PairWithOne' >> beam.Map(lambda word: (word, 1))
        | 'CountPerWord' >> beam.CombinePerKey(sum)
    )
    
    # 步骤4:格式化输出
    output = counts | 'FormatResult' >> beam.Map(lambda kv: f'{kv[0]}: {kv[1]}')
    
    # 步骤5:写回文件
    output | 'WriteResults' >> beam.io.WriteToText('output')

代码解析

  • 管道构建with beam.Pipeline() as p 创建一个流水线对象,退出上下文时自动调用运行。
  • PTransform 链条:每一步用 | 操作符连接,并通常给出一个标签(如 'ReadLines')。
  • 核心变换
    • beam.FlatMap 类似 Map + Flatten,将一行拆分为多个单词。
    • beam.Map 将单词映射为 (word, 1) 元组。
    • beam.CombinePerKey(sum) 按 key 聚合求和,是 GroupByKey + Combine 的优化组合。
  • IO 连接器ReadFromTextWriteToText 是内置的文件源和汇。

运行此脚本,默认使用 DirectRunner(本地执行),你会得到每个单词的出现次数。

批流统一的窗口模型

有界数据(批处理)在读取完成时数据集就确定了,但无界数据(流处理)是持续到达的。为了对无界数据进行聚合,必须引入窗口的概念。

窗口类型

Beam 提供多种内建窗口策略:

  • 固定窗口(Fixed Windows) :按固定时间长度切片,例如每 30 秒一个窗口。
  • 滑动窗口(Sliding Windows) :窗口长度和滑动步长均可定义,允许窗口重叠。
  • 会话窗口(Session Windows) :根据数据活动的间隙动态划分,用于用户行为分析等。

为流式 Pipeline 添加窗口

以实时点击流分析为例,统计每 10 分钟的用户点击量。

with beam.Pipeline() as p:
    clicks = (
        p
        | 'ReadFromPubSub' >> beam.io.ReadFromPubSub(topic='projects/myproj/topics/clicks')
        | 'ParseJSON' >> beam.Map(json.loads)
        | 'WithEventTime' >> beam.Map(lambda e: beam.window.TimestampedValue(
            (e['user_id'], 1), e['timestamp']))
    )
    # 应用固定窗口
    windowed = clicks | 'WindowInto' >> beam.WindowInto(
        beam.window.FixedWindows(600)  # 10分钟固定窗口
    )
    # 按用户聚合
    user_counts = (
        windowed
        | 'GroupByUser' >> beam.CombinePerKey(sum)
    )
    # 输出结果
    user_counts | 'WriteToBigQuery' >> beam.io.WriteToBigQuery(...)

关键点

  • TimestampedValue:为每个元素绑定事件时间,Beam 以此作为窗口分配的依据。
  • WindowInto:将 PCollection 分配到指定的窗口。

事件时间与水印、触发器

事件时间 vs 处理时间

  • 事件时间:数据实际发生的时间,存储在元素内。
  • 处理时间:数据进入流水线的时间。 默认情况下,Beam 优先使用事件时间进行窗口聚合。需要显式指定时间戳来源。

水印(Watermark)

是一个系统估计值,表示“在此时间之前的事件全部到达”的界限。Beam 通过水印决定何时触发窗口计算。当水印越过窗口末尾时,系统认为该窗口数据已完整。

触发器

控制何时将窗口内的中间结果物化输出。在水印触发之外,还可以设置:

  • 处理时间触发器:例如每隔 5 分钟触发一次早期结果。
  • 数据量触发器:例如每达到 100 个元素触发。
  • 复合触发器:组合多种条件。

示例:在窗口结束前,每 2 分钟输出一次中间计数。

trigger = (
    beam.trigger.Repeatedly(
        beam.trigger.AfterProcessingTime(120)  # 每2分钟处理时间触发
    )
    | beam.trigger.AfterWatermark()            # 水印通过后触发最终结果
)
windowed = clicks | beam.WindowInto(
    beam.window.FixedWindows(600),
    trigger=trigger,
    accumulation_mode=beam.trigger.AccumulationMode.ACCUMULATING
)

ACCUMULATING 模式表示后续触发会包含之前的结果,形成累加更新。

Beam Runner 选择指南

Runner 适用场景 执行模式
DirectRunner 本地开发、调试、小规模测试 批 / 流(有限)
FlinkRunner 流处理主力,支持精确一次语义 批 / 流
SparkRunner 批量处理,或 Spark 生态集成 批(微批流)
DataflowRunner 全托管服务,自动扩缩容,强一致性 批 / 流
Twister2Runner 早期研究型 Runner,学术/实验用 批 / 流

切换 Runner 只需修改 PipelineOptions,无需变更业务代码。例如:

from apache_beam.options.pipeline_options import FlinkRunnerOptions, PipelineOptions
options = PipelineOptions(runner='FlinkRunner', flink_master='localhost:8081')

高级特性简述

状态与计时器(State & Timers)

ParDo 中可以按 key 和窗口维护持久化状态,并设置计时器实现复杂的有状态处理(如检测会话超时)。这是实现精确流式算法的基础。

侧输入(Side Inputs)

PTransform 除了主输入 PCollection 外,还可以接收辅助数据集(如字典、配置),用于数据过滤或关联。侧输入会在管道优化时广播或缓存。

动态数据处理(Schema & SQL)

Beam 支持 Schema 化 PCollection(结构化的命名字段),并可直接用 SQL 查询。

from apache_beam.transforms.sql import SqlTransform
result = data | SqlTransform("SELECT user, COUNT(*) FROM PCOLLECTION GROUP BY user")