在 Dataflow 上实现灵活高效的设计

本文档介绍了设计灵活性和效率 (DFE),这是一套用于构建弹性 Dataflow 流水线的架构最佳实践。

从僵硬的基础设施限制过渡到灵活的资源定义有助于您:

  • 最大限度提高计算资源的可获取性。
  • 确保在区域需求旺盛期间实现无缝自动扩缩。
  • 防止流水线启动延迟并消除容量瓶颈。

例如,您可以设置最低资源需求(例如 4 个 vCPU 和 16 GB RAM),而不是将流水线限制为某个可用区中的一种特定机器类型(例如要求在 us-central1-a 中使用 n1-standard-4 工作器)。如果 us-central1-a 或 N1 机器系列遇到临时容量限制,Dataflow 可以在其他可用区和机器系列(例如 E2、N2 或 N2D)中自动预配兼容的工作器虚拟机。这种灵活性有助于确保流水线在启动和扩缩时无需等待单个受限的硬件池。

本文档面向管理 Dataflow 工作负载并希望优化流水线可靠性、吞吐量和基础架构可用性的数据工程师、云架构师和平台管理员。

DFE 概览

Dataflow 是一项全代管式无服务器数据处理服务,可动态预配 Compute Engine 虚拟机 (VM) 实例来执行 Apache Beam 流水线。在大型批处理和流处理流水线中,工作器池经常会扩容到数十个或数百个虚拟机实例。

如果流水线配置了严格的基础设施限制,则在需求高峰期容易出现配置延迟。刚性限制的示例包括:

  • 对单个机器类型(例如 n1-standard-4)进行硬编码。
  • 将流水线固定到特定的 Compute Engine 可用区。

如果该特定机器类型或可用区暂时需求旺盛,Dataflow 将无法分配计算资源。这可能会导致配置延迟或错误,例如 ZONE_RESOURCE_POOL_EXHAUSTED 或 RESOURCE_POOL_EXHAUSTED。

使用 DFE 原则有助于将流水线架构从僵化的静态基础设施声明转变为灵活的基于需求的资源定义。这种灵活性使 Dataflow 能够在 Google Cloud中动态地将计算分配到各种可用的硬件池中,从而帮助您最大限度地提高计算可获取性,同时最大限度地减少运营开销。

DFE 最佳实践

采用以下最佳实践可最大限度地提高计算可获得性、改进自动扩缩响应能力并构建弹性流水线。

启用自动选择虚拟机

不要使用 worker machine type 流水线选项对静态机器类型进行硬编码,而是使用 Auto VM selection 和 Apache Beam 资源提示。当您指定最低资源要求(min_ram 或 cpu_count)时,Dataflow 会自动启用实例灵活性,并从兼容的机器类型列表中预配工作器。

工作负载支持:

  • 批处理流水线:指定资源提示时,系统会自动启用适配和自动选择虚拟机功能。
  • 流处理流水线:若要启用适配功能,需要设置 --experiments=enable_streaming_rightfitting 流水线选项,同时启用横向自动扩缩(默认启用)和 Streaming Engine (--enable_streaming_engine)。

如需配置自动虚拟机选择,请使用命令行选项、SDK 流水线选项或 Flex 模板执行参数在流水线级别指定最低资源要求(min_ram 或 cpu_count)。如需查看 Java 和 Python 的详细设置说明和代码示例,请参阅使用资源提示。

使用区域级工作器放置(避免可用区级固定)

将 Dataflow 配置为在所选区域内的任何正常运行的可用区中动态调度工作器虚拟机。

指定 --region 流水线选项,并省略 --zone 和 --worker_zone。例如:

--region=us-central1

使用托管式服务分离状态和 shuffle

不使用托管后端服务的流水线直接在工作器虚拟机磁盘和内存中执行 shuffle 数据操作和流处理状态存储。这种紧密耦合需要更大的工作器磁盘,并将工作负载生存期与特定虚拟机实例绑定,从而在容量受限时更难替换工作器。

  • 对于批处理作业 - 使用 Dataflow Shuffle:对于在受支持的工作器机器类型上运行的批处理流水线,默认启用 Dataflow Shuffle,并将 shuffle 操作从工作器虚拟机分流到由 Google 管理的专用后端服务。
  • 对于流处理作业 - 使用 Streaming Engine:Streaming Engine 会将窗口状态存储和计时器管理从工作器虚拟机卸载到专用、高度响应的后端基础架构。对于使用 Apache Beam SDK 2.30.0 或更高版本的流水线,Streaming Engine 默认处于启用状态。如需明确启用此功能,请传递 --enable_streaming_engine 流水线选项。

为批处理流水线使用灵活资源调度服务 (FlexRS)

对于非时间关键型批处理工作负载(例如夜间 ETL、数据湖提取或每日汇总),请使用灵活资源调度 (FlexRS)。

如需启用 FlexRS,请设置 flexRS 目标流水线选项:

  • 对于 Python 流水线:--flexrs_goal=COST_OPTIMIZED
  • 对于 Java 流水线:--flexRSGoal=COST_OPTIMIZED

为 Flex 模板配置灵活的启动器虚拟机类型

使用 Flex 模板启动流水线时,流水线启动器虚拟机的默认值为 e2-standard-2。在大多数情况下,默认虚拟机都能正常运行,但如果您遇到容量限制,可以在运行 gcloud dataflow flex-template run 命令时使用 --launcher-machine-type 选项自定义配置:

gcloud dataflow flex-template run my-job \
    --template-file-gcs-location="gs://my-bucket/template.json" \
    --region="us-central1" \
    --launcher-machine-type="n2-standard-2"

运营注意事项和权衡

虽然采用 DFE 最佳实践可显著提高计算可获得性、自动扩缩响应速度和运营可靠性,但在设计架构时,请考虑以下运营因素和权衡取舍:

自动选择虚拟机的注意事项

  • 可靠性与峰值性能:自动选择虚拟机时,系统会优先考虑作业启动可靠性和计算可获得性,而不是峰值执行性能。由于 Dataflow 会从多个候选机器家族(例如 E2、N2、N4 和 N2D)中预配资源,因此运行时性能和吞吐量可能会因预配的机器家族而略有不同。对于执行 SLA 严格的计算密集型工作负载,请先通过自动选择虚拟机来测试流水线,以建立性能基准,然后再广泛部署。如果工作负载需要特定的硬件平台或时钟速度,并且您可以容忍容量限制,则可以继续设置特定的机器类型。
  • 候选系列中的 Compute Engine 配额:由于自动选择虚拟机功能可以从多个候选机器系列中预配工作器,因此请确保您的 Google Cloud 项目在目标区域中为每个候选系列预留了足够的 Compute Engine vCPU 和内存配额。如果主系列中出现容量不足的情况,而您的项目缺少后备系列的配额,则工作器配置会失败并显示 QUOTA_EXCEEDED 错误。
  • 流处理流水线的前提条件:对于流处理流水线,默认情况下不会启用适配和自动选择虚拟机。您必须明确指定 --experiments=enable_streaming_rightfitting,并确保 Streaming Engine (--enable_streaming_engine) 和横向自动扩缩功能均处于有效状态。
  • 配置排除项:如果您配置了下表中的任何功能或选项,系统会自动绕过或不支持自动选择虚拟机:

    功能 标志或配置选项 备注
    显式机器类型 --worker_machine_type 或 --machine_type (Python)
    --workerMachineType (Java)
    系统会绕过自动虚拟机选择,而采用指定的机器类型。
    自定义磁盘类型、预配 IOPS 或吞吐量 --disk_type、--disk_provisioned_iops 或 --disk_provisioned_throughput_mibps 系统会绕过自动选择虚拟机。支持使用 --disk_size_gb 设置自定义磁盘大小。
    满足最低 CPU 要求的平台 --min_cpu_platform (Python)
    --minCpuPlatform (Java)
    设置最低要求 CPU 平台会绕过自动虚拟机选择。
    机密虚拟机 --experiments=enable_confidential_compute 自动选择虚拟机功能不支持机密虚拟机实例。
    GPU 或 TPU 加速器 --dataflow_service_options=worker_accelerator=... 或 accelerator 资源提示 自动选择虚拟机仅适用于不含加速器的工作负载。
    Dataflow Prime --dataflow_service_options=enable_prime Dataflow Prime 使用纵向自动扩缩和动态精确调整,而不是自动选择虚拟机。
    Flexible Resource Scheduling (FlexRS) --flexrs_goal=COST_OPTIMIZED (Python)
    --flexRSGoal=COST_OPTIMIZED (Java)
    FlexRS 会管理自己的工作器池和调度缓冲区。

Flexible Resource Scheduling (FlexRS) 的权衡取舍

  • 调度延迟时间窗口:FlexRS 可以在作业执行开始前引入最长 6 小时的调度缓冲区。请勿将 FlexRS 用于具有严格完成时间 SLA 或紧密下游依赖关系的流水线。

区域放置和数据本地化

  • 托管式服务前提条件:只有使用 Dataflow Shuffle(用于批处理)或 Streaming Engine(用于流处理)的作业支持区域工作器放置。不使用这些受管后端服务的作业会使用自动可用区放置,该功能会在区域内选择单个最佳可用区。
  • 数据本地化和跨区域出站流量:区域级放置会将工作器分布在所选区域内的可用区中。为最大限度地减少网络延迟并避免产生区域间网络出站流量费用,请确保所有数据源和接收器(例如 Cloud Storage 存储桶、BigQuery 数据集和 Pub/Sub 主题)都位于 Dataflow 作业所在的同一区域中。

Compute Engine 预留

  • 预留亲和性:按需 Dataflow 作业会自动使用采用 ANY 预留亲和性的匹配 Compute Engine 预留。不过,“自动选择虚拟机”功能不支持从特定命名预留中消耗实例。
  • 是否适合临时性工作负载:Compute Engine 预留通常不建议用于容易出现高峰、需要自动扩缩或短期运行的批处理工作负载。此外,在可用区短缺期间创建新预留会失败,并且会受到与按需创建虚拟机相同的容量限制。

后续步骤