Dataflow 旨在通过在托管的计算实例池中分配工作来执行大型数据处理流水线。了解 Dataflow 如何并行化处理有助于您设计高效的流水线、避免性能瓶颈并优化资源费用。
本页介绍了 Dataflow 如何并行处理数据、如何管理和扩缩执行、限制并行性的常见因素,以及可用于优化流水线吞吐量的技术。
并行处理模型:横向与纵向
Dataflow 通过两种互补的策略来实现并行处理:
横向并行处理:流水线数据被分区,并在多个工作器实例(虚拟机)上并发处理。Dataflow 可以通过横向自动扩缩功能,根据工作负载需求自动调整工作器池大小。默认情况下,Dataflow 会将每个作业的资源限制设置为 4,000 个工作器,您可以使用配额请求调整此限制。
纵向并行处理:单个工作器实例中的多个 CPU 核心和线程同时处理流水线数据。每个工作器虚拟机都会运行工作器进程和驱动程序线程,以使用可用的计算资源。借助动态线程伸缩,Dataflow 可以根据 CPU 利用率和内存余量调整批处理流水线中每个工作器的活跃线程数。在 Dataflow Prime 中,纵向自动扩缩功能可动态扩缩分配给工作器的内存和计算资源。
工作单元和执行层次结构
为了在工作器和线程之间分配处理任务,Dataflow 会将 Apache Beam 流水线分解为离散的工作单元:
- PCollection 和分区:
PCollection表示分布式数据集。对于有界数据(批处理流水线),Dataflow 会将数据集划分为多个拆分或分片。对于无界限数据(流式处理流水线),数据会持续到达,并以消息或流分区的形式提取。 - Bundle:Dataflow 会将元素分组为任意 bundle,以供
DoFn进行处理。软件包是失败和重试的单位:如果处理某个元素时引发了未处理的异常,系统会重试整个软件包。内存消耗量高的操作可能会增加工作器内存压力,并导致内存不足错误。 - 阶段和步骤融合:在图优化期间,Dataflow 会将相邻的转换合并到融合的执行阶段,以消除中间数据具体化的开销。在融合阶段内,元素会在单线程上以紧凑的执行循环进行处理,然后再传递到下一个阶段或 shuffle 边界。
如需详细了解流水线转换和图表生成,请参阅流水线生命周期。
托管式并行性和自动扩缩
默认情况下,Dataflow 会自动管理流水线并行性,而无需手动调整分区,具体方式如下:
- 横向自动扩缩:
- 批处理流水线:评估总估计剩余工作量、来源积压和 CPU 使用率,以纵向扩容或缩容工作器池,从而快速且经济高效地完成作业。
- 流处理流水线:分析系统延迟时间、积压大小和 CPU 利用率,以便在吞吐量激增时扩缩工作器,并在流量较低时缩容工作器。如需了解详情,请参阅调整流式横向自动扩缩。
- 动态工作负载再平衡 (DWR):在批量流水线中,Dataflow 会监控各个工作器任务的进度。如果某个工作器提前完成工作,或者另一个工作器因数据倾斜(落后者)而落后,Dataflow 会动态拆分慢速工作器剩余的未处理工作,并将其重新分配给空闲工作器。如需了解详情,请参阅动态工作负载再均衡。
- 动态线程伸缩:在采用可移植 Runner 的批处理流水线中,根据 CPU 利用率和内存余量自动调整每个工作器的并发处理线程数。如需了解详情,请参阅动态线程伸缩。
- 纵向自动扩缩:在 Dataflow Prime 中,Dataflow 会动态扩缩工作器内存和计算资源,以防止出现内存不足错误并优化资源利用率。如需了解详情,请参阅纵向自动扩缩。
限制并行性的因素
由于以下数据特征或流水线图设计,流水线可能无法实现预期的并行性:
无法拆分的输入源
如果输入源无法拆分为独立的范围,Dataflow 将被迫使用单个工作器线程按顺序读取该源:
- 不可拆分的文件压缩:
.gz(gzip) 或.bzip2(无索引)等格式无法从任意字节偏移量并行读取。读取单个大型压缩文件会将提取阶段限制为单线程,直到数据被解压缩并重新分发。 - 解决方案:以可拆分的文件格式(例如 Parquet、Avro 或 Snappy 压缩格式)存储数据,或者将输入数据拆分为 Cloud Storage 中的多个较小文件。
步融合和高扇出
步融合通过减少序列化开销来提高性能,但当并行性较低的步生成大量输出元素(“高扇出”操作)时,可能会无意中限制并行性并增加内存压力:
- 示例:某个来源读取了 5 个文件,并与生成 1,000,000 个输出元素的
FlatMap转换融合。如果FlatMap转换与下游转换融合,所有 1,000,000 个元素将继续在最多五个工作线程上执行,从而严重限制下游吞吐量。此外,如果中间转换在提交之前在内存中显著扩展,则大型软件包可能会耗尽可用的工作器内存。 - 解决方案:在高扇出步骤和下游转换之间插入
Redistribute(或经典Reshuffle)转换,以打破融合并跨工作器池重新分配工作。如需调试相关的内存问题,请参阅排查内存不足错误。
键倾斜和热门键
聚合操作(GroupByKey、CoGroupByKey、Combine.PerKey)会按元素的相关联键对元素进行分组。
- 热键瓶颈:Dataflow 会将具有相同键的所有元素路由到单个工作器线程以进行聚合。如果单个键包含的数据占整个数据集的很大一部分,那么相应工作器就会成为落后者,上游工作器可能会遇到反压力。例如,默认的
null键或极其热门的类别键。 - 解决方法:
- 尽可能使用 Combiner(
CombineFn或Combine.PerKey)而不是GroupByKey,以便 Dataflow 在 Shuffle 之前执行部分本地组合。 - 向热门键添加随机整数前缀或后缀(键加盐),以在工作器之间分配键空间,然后进行第二阶段的聚合以合并加盐后的结果。
- 尽可能使用 Combiner(
下游接收器节流
将流水线输出写入数据库或第三方 API 等外部服务时,高并行度可能会使目标系统饱和:
- 限制:数百个工作线程同时发出写入调用可能会导致速率限制错误、连接超时或数据库性能下降。
- 解决方法:
- 通过使用
GroupByKey对元素进行分组或使用具有受控并行性的批处理接收器来限制写入并行性。 - 在接收器
DoFn实现中实现客户端指数退避算法和重试逻辑。
- 通过使用
优化策略
如需优化 Dataflow 作业中的并行性,请考虑以下方法:
- 使用
Redistribute防止不必要的融合:Redistribute.arbitrarily():打破了步融合,并将元素均匀地重新分配给所有可用的工作器。Redistribute.byKey():在保持键局部性的同时,在工作线程之间重新平衡键值对。- 如需查看实现示例,请参阅阻止融合。
- 监控落后者和瓶颈: 使用 Google Cloud 控制台执行详情来识别落后者数量较多或进度停滞的阶段:
后续步骤
- 了解流水线生命周期。
- 探索横向自动扩缩。
- 了解动态工作负载再平衡。
- 查看 Dataflow 流水线最佳实践。
- 了解如何排查内存不足错误。