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 连接器:
ReadFromText和WriteToText是内置的文件源和汇。
运行此脚本,默认使用 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")