Google PDE: 使用 Dataflow 和 Apache Beam 进行流处理 — 学习指南
属于 Google Professional Data Engineer — 学习指南. 使用经过验证的答案练习: Google 考试中心, 或参加限时模拟考试: ExamRoll.io.
概述
Google Cloud 上的流处理以 Apache Beam 的统一编程模型为核心,并由 Dataflow runner 执行。Beam 提供了一个逻辑抽象——对 PCollection 应用一系列转换 (transform) 的流水线 (pipeline)——它将您的代码与并行度、自动扩缩容和容错等执行细节解耦。在流处理中,正确性取决于时间语义(事件时间 vs 处理时间)、窗口化(固定、滑动、会话、全局)、水印、触发器以及对迟到数据的处理。在 Dataflow 上实现卓越运营需要正确的 worker 规格、自动扩缩容策略、流处理引擎、shuffle 选择、幂等接收器设计、死信处理和强大的可观测性。
Apache Beam 模型和时间语义
流水线、转换、PCollection、runner:
- Beam 流水线将一个由 PTransform 组成的有向无环图应用于 PCollection(有界或无界)。
- Runner(如 Dataflow、Spark、Flink、Direct)负责执行流水线;Dataflow 提供托管的自动扩缩容、检查点和运营可见性。
- 转换包括逐元素操作 (ParDo)、分组与组合 (GroupByKey, Combine)、连接 (CoGroupByKey) 以及 IO 操作 (PubSubIO, BigQueryIO, FileIO)。
窗口:
- 固定窗口:无重叠的时间片(例如,1 分钟的滚动窗口),用于周期性聚合。
- 滑动窗口:有重叠的窗口,用于平滑的滚动指标(例如,窗口大小为 5 分钟,每 1 分钟滑动一次)。
- 会话窗口:在一段不活跃间隙后关闭的动态窗口,非常适合用户会话或设备突发活动场景。
- 全局窗口:整个无界流的默认非窗口化视图;通常与触发器配合使用,以实现周期性物化。
事件时间 vs 处理时间:
- 事件时间:事件在源头发生的时间;即使传输延迟可变,也能实现逻辑上一致的聚合。
- 处理时间:流水线观察到事件的时间;对操作性触发器有用,但不能保证语义上的正确性。
水印:
- 水印估算事件时间的完整性(即 runner 推测它已看到截至时间 T 的所有事件)。
- 在反压或源延迟的情况下,水印可能不规律地推进或停滞;时间戳小于水印的任何到达数据都属于迟到数据。
触发器和迟到:
- 默认:AfterWatermark 触发器,在水印通过窗口末端时触发;如果允许的迟到时间 = 0,则迟到数据会被丢弃。
- 提前触发(基于处理时间或计数)可以提供低延迟的初步结果。
- 延迟触发允许在迟到数据到达时进行修正;累积模式决定了窗格是累积结果还是丢弃先前的输出。
- 根据业务容忍度以及存储/计算的权衡来选择允许的迟到时间;更长的迟到时间会增加状态保留时间和成本。
有状态处理、计时器、会话化、去重:
- 有状态的 DoFn 会为每个键保留状态(例如,最后一次看到的事件、运行中的聚合),并设置计时器来发出或清除状态。
- 会话化可以通过 SessionWindows 自然地表达;对于自定义逻辑,请使用键控状态和处理时间/事件时间计时器。
- 去重:为每个事件使用一个稳定的 ID,并结合使用窗口内的 Distinct/Combine,或使用每个键的状态(例如,带 TTL 的布隆过滤器或集合)。需要在内存、误报率与严格准确性之间进行权衡。
故障模式和权衡:
- 使用处理时间窗口计算业务指标会在流量尖峰或重试期间导致结果漂移;应优先使用事件时间窗口。
- 过小的窗口加上频繁的提前触发会导致过多的窗格产出和接收器写入放大。
- 无限制的允许迟到时间会导致状态膨胀;务必为状态设置 TTL 边界,并设置计时器来清除休眠键。
运营流处理工作负载的 Dataflow
工作器规模调整与自动扩缩容:
- 水平自动扩缩容会根据积压、水印延迟、CPU 和吞吐量来增减工作器;设置合理的 maxWorkers 以吸收流量尖峰。
- 针对瓶颈选择机器类型:CPU 密集型(更多 vCPU)、内存密集型(高内存类型)、网络密集型(更大的虚拟机可减少 Shuffle 开销)。
- 对于繁重的 Shuffle 或基于文件的接收器,增加启动磁盘。监控系统延迟和积压秒数。
Streaming Engine 和 Shuffle:
- Streaming Engine 将状态和 Shuffle 操作外部化到服务后端,从而提高弹性、减轻工作器内存压力并实现更快的更新。
- 对于批处理繁重的阶段或大规模的键分组,使用 Dataflow Shuffle 将 Shuffle I/O 从工作器中卸载。两者都能减少热点工作器故障和磁盘抖动。
反压、热键和数据倾斜:
- Dataflow 通过动态工作重平衡来管理反压;尽管如此,在适用时仍需调整源的流控制(例如,Pub/Sub 中未处理的消息/字节数)。
- 热键(例如,热门 ID)会产生拖后腿的任务。通过键分片 (key#N)、部分预聚合后重新分区键或基于 sketch 的近似计算来缓解。
- 由异常记录(巨大负载)或突发性发布者引起的倾斜可能需要按发布者分区、批处理或压缩。
Pub/Sub 集成:
- 使用 Pub/Sub 主题进行数据注入;启用消息属性以携带元数据(例如,deviceId、事件时间戳)。
- 使用 PubSubIO 进行注入;从属性或负载中提取事件时间戳,否则回退到使用发布时间。
- 排序键提供按键排序;由于“至少一次”交付的特性,Dataflow 仍需要下游行为具有幂等性。
流式传输到 BigQuery 的模式:
- 优先使用 BigQueryIO 配合 Storage Write API,通过流偏移量和自动重试,在流内部实现高吞吐、低延迟的“精确一次”语义。
- 对于速率较低的简单管道,流式插入是可接受的;设置 insertId 以便为客户端重试去重。
- 对流式缓冲区的查询是最终一致的;对于时间关键型分析,可在缓冲区延迟后进行查询(例如,等待约 2 倍的观测可用性延迟),或通过微批次窗口和 Storage Write API 的 committed 模式进行物化。
精确一次效果、幂等性、重放和接收器:
- Beam 保证“至少一次”处理;“精确一次”必须在接收器端通过幂等写入、事务或去重键来实现。
- BigQuery:使用 Storage Write API 的默认流或 committed 流,以在流内部实现精确一次;对于流式插入,设置一个稳定的 insertId。
- 文件:写入具有唯一名称的临时文件,在窗口完成时最终确定,并确保原子性重命名;避免覆盖以防止部分重复。
- 外部数据库:使用基于稳定 ID 的更新插入 (upsert) 或实现去重窗口。
- 为重放而设计:保持转换的确定性;确保接收器在重试时能进行去重。
死信处理、错误路由、可观测性:
- 在 ParDo 中将有风险的解析/扩充操作包装在 try/catch 块中,并通过 TupleTag 将失败记录发送到一个死信 PCollection;附带上负载、错误代码和上下文。
- 将死信队列 (DLQ) 路由到 BigQuery 或 Cloud Storage 进行分析;考虑使用一个单独的 Pub/Sub 主题进行重新处理。
- 可观测性:使用 Dataflow 作业指标(水印延迟、系统延迟、吞吐量)、自定义计数器、分布指标以及 Cloud Logging 中的各步骤日志。在 Cloud Monitoring 中针对延迟和错误率创建警报。使用 Error Reporting 聚合异常。
性能调优模式:
- 高效读取:对于 BigQuery 源,优先使用 Storage Read API 或基于查询的读取,仅选择所需字段并应用筛选。
- 组合器提升:在 GroupByKey 之前使用组合器以减少 Shuffle 数据量。
- 旁路输入:将小型参考数据缓存在内存中;注意扇出和更新频率。
- 序列化:使用紧凑的模式(Avro/Proto),并避免在热路径上进行过多的 JSON 解析。
← BigQuery 分析与数仓工程 · 所有领域 · 消息传递、事件注入与实时服务 →
练习这些题目 → · 在 ExamRoll.io 上限时练习 →
Pass the whole exam — not just this question
You found this answer. Get every verified question and explanation in one place, and save hours of prep. Free to start.
通过考试 →