You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何基于单个Google Cloud Dataflow Flex模板镜像,通过多命令CLI运行多管道?

基于单个Flex模板镜像运行多个Dataflow Python管道的实现方案

一、重构Python脚本结构

将多个管道封装为独立函数,通过自定义PipelineOptions接收管道名称参数,实现单脚本多管道的调用逻辑:

import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions, GoogleCloudOptions

# 自定义管道选项,用于接收管道名称参数
class MultiPipelineOptions(PipelineOptions):
    @classmethod
    def _add_argparse_args(cls, parser):
        parser.add_argument(
            '--pipeline_name',
            required=True,
            help='指定要运行的管道名称,可选值:pipeline1、pipeline2'
        )

# 管道1的业务逻辑
def run_pipeline1(options):
    with beam.Pipeline(options=options) as p:
        (p
         | "读取管道1输入" >> beam.io.ReadFromText(options.view_as(GoogleCloudOptions).input)
         | "处理管道1数据" >> beam.Map(lambda x: x.upper())
         | "写入管道1输出" >> beam.io.WriteToText(options.view_as(GoogleCloudOptions).output)
        )

# 管道2的业务逻辑
def run_pipeline2(options):
    with beam.Pipeline(options=options) as p:
        (p
         | "读取管道2输入" >> beam.io.ReadFromText(options.view_as(GoogleCloudOptions).input)
         | "处理管道2数据" >> beam.Map(lambda x: x.lower())
         | "写入管道2输出" >> beam.io.WriteToText(options.view_as(GoogleCloudOptions).output)
        )

if __name__ == '__main__':
    # 解析所有管道选项
    pipeline_options = PipelineOptions()
    specific_options = pipeline_options.view_as(MultiPipelineOptions)
    
    # 根据参数调用对应管道
    if specific_options.pipeline_name == "pipeline1":
        run_pipeline1(pipeline_options)
    elif specific_options.pipeline_name == "pipeline2":
        run_pipeline2(pipeline_options)
    else:
        raise ValueError(f"不支持的管道名称:{specific_options.pipeline_name}")

二、配置Flex模板Dockerfile

遵循Dataflow Flex模板规范,确保镜像能正确加载脚本和依赖:

# 使用官方Dataflow Python基础镜像
FROM gcr.io/dataflow-templates-base/python39-template-launcher-base

WORKDIR /dataflow

# 复制依赖文件并安装
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt

# 复制多管道脚本
COPY pipeline_script.py .

# 设置Flex模板必要环境变量
ENV FLEX_TEMPLATE_PYTHON_PY_FILE="/dataflow/pipeline_script.py"
ENV FLEX_TEMPLATE_PYTHON_REQUIREMENTS_FILE="/dataflow/requirements.txt"
# PY_OPTIONS用于传递Python解释器参数(如-u开启无缓冲输出),非脚本业务参数
ENV FLEX_TEMPLATE_PYTHON_PY_OPTIONS="-u"

三、定义模板元数据(metadata.json)

通过元数据文件暴露可配置参数,方便启动作业时指定管道及其他参数:

{
  "name": "多管道Flex模板",
  "description": "从单个镜像运行多个Dataflow管道",
  "parameters": [
    {
      "name": "pipeline_name",
      "label": "管道名称",
      "helpText": "选择要运行的管道:pipeline1或pipeline2",
      "isOptional": false
    },
    {
      "name": "input",
      "label": "输入路径",
      "helpText": "GCS输入文件路径",
      "isOptional": false
    },
    {
      "name": "output",
      "label": "输出路径",
      "helpText": "GCS输出文件路径",
      "isOptional": false
    }
  ]
}

四、构建并部署Flex模板

  1. 构建并推送镜像到GCR:
docker build -t gcr.io/你的项目ID/multi-pipeline-flex-template .
docker push gcr.io/你的项目ID/multi-pipeline-flex-template
  1. 创建Flex模板并上传到GCS:
gcloud dataflow flex-template build gs://你的存储桶/templates/multi-pipeline-template.json \
  --image-gcr-path gcr.io/你的项目ID/multi-pipeline-flex-template \
  --sdk-language PYTHON \
  --metadata-file metadata.json

五、启动指定管道的作业

通过--parameters参数指定要运行的管道及输入输出:

# 启动pipeline1
gcloud dataflow flex-template run "pipeline1-job-$(date +%Y%m%d-%H%M%S)" \
  --template-file-gcs-location gs://你的存储桶/templates/multi-pipeline-template.json \
  --region us-central1 \
  --parameters pipeline_name=pipeline1 \
  --parameters input=gs://你的存储桶/input/pipeline1.txt \
  --parameters output=gs://你的存储桶/output/pipeline1_result

# 启动pipeline2
gcloud dataflow flex-template run "pipeline2-job-$(date +%Y%m%d-%H%M%S)" \
  --template-file-gcs-location gs://你的存储桶/templates/multi-pipeline-template.json \
  --region us-central1 \
  --parameters pipeline_name=pipeline2 \
  --parameters input=gs://你的存储桶/input/pipeline2.txt \
  --parameters output=gs://你的存储桶/output/pipeline2_result

关于FLEX_TEMPLATE_PYTHON_PY_OPTIONS的说明

该环境变量仅用于传递Python解释器级别的参数(如-u开启无缓冲输出、-O开启优化模式),不能直接传递脚本的业务参数。业务参数需通过Dataflow的--parameters机制或自定义PipelineOptions传递,这也是Flex模板的标准用法。

内容的提问来源于stack exchange,提问作者seandavi

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.21 20:48:45