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

如何让Dataflow工作节点获取自定义Pipeline选项中的环境变量?

问题描述

我有一个基于Apache Beam、以Dataflow为运行器的Python应用,该应用使用非公开Python包uplight-telemetry,通过创建pipeline_options对象时配置extra_packages引入。此包依赖名为OTEL_SERVICE_NAME的环境变量,但由于Dataflow工作节点中不存在该变量,导致应用启动报错。

我已通过自定义Pipeline选项传入该变量,创建pipeline_options的代码如下:

pipeline_options = ProcessBillRequests.CustomOptions(
    project=gcp_project_id,
    region="us-east1",
    job_name=job_name,
    temp_location=f'gs://{TAS_GCS_BUCKET_NAME_PREFIX}{os.getenv("UP_PLATFORM_ENV")}/temp',
    staging_location=f'gs://{TAS_GCS_BUCKET_NAME_PREFIX}{os.getenv("UP_PLATFORM_ENV")}/staging',
    runner='DataflowRunner',
    save_main_session=True,
    service_account_email= service_account,
    subnetwork=os.environ.get(SUBNETWORK_URL),
    extra_packages=[uplight_telemetry_tar_file_path],
    setup_file=setup_file_path,
    OTEL_SERVICE_NAME=otel_service_name,
    OTEL_RESOURCE_ATTRIBUTES=otel_resource_attributes
    # Set values for additional custom variables as needed
)

执行流水线的代码如下:

result = (
        pipeline
        | "ReadPendingRecordsFromDB" >> read_from_db
        | "Parse input PCollection" >> beam.Map(ProcessBillRequests.parse_bill_data_requests)        
        | "Fetch bills " >> beam.ParDo(ProcessBillRequests.FetchBillInformation())
)
pipeline.run().wait_until_finish()

请问是否有办法让Dataflow工作节点获取自定义选项中的这些环境变量?


解决方案

Dataflow工作节点不会自动将自定义Pipeline选项转为环境变量,你可以通过以下几种方式解决:

1. 使用Runner v2的environment_variables参数(推荐)

如果使用Dataflow Runner v2,可以直接在Pipeline Options中配置environment_variables,该参数会自动在所有Worker进程中设置指定的环境变量,无需修改业务代码:

pipeline_options = ProcessBillRequests.CustomOptions(
    project=gcp_project_id,
    region="us-east1",
    job_name=job_name,
    temp_location=f'gs://{TAS_GCS_BUCKET_NAME_PREFIX}{os.getenv("UP_PLATFORM_ENV")}/temp',
    staging_location=f'gs://{TAS_GCS_BUCKET_NAME_PREFIX}{os.getenv("UP_PLATFORM_ENV")}/staging',
    runner='DataflowRunner',
    save_main_session=True,
    service_account_email= service_account,
    subnetwork=os.environ.get(SUBNETWORK_URL),
    extra_packages=[uplight_telemetry_tar_file_path],
    setup_file=setup_file_path,
    experiments=['use_runner_v2'],  # 必须启用Runner v2
    environment_variables={
        'OTEL_SERVICE_NAME': otel_service_name,
        'OTEL_RESOURCE_ATTRIBUTES': otel_resource_attributes
    }
)

2. 在DoFn的setup方法中设置环境变量

如果无法使用Runner v2,可以在自定义的DoFn的setup方法中,将Pipeline选项中的值注入到Worker进程的环境变量中。setup方法会在Worker进程初始化时执行,确保后续代码能读取到环境变量:

class FetchBillInformation(beam.DoFn):
    def __init__(self, otel_service_name, otel_resource_attributes):
        self.otel_service_name = otel_service_name
        self.otel_resource_attributes = otel_resource_attributes

    def setup(self):
        import os
        # 设置环境变量
        os.environ['OTEL_SERVICE_NAME'] = self.otel_service_name
        os.environ['OTEL_RESOURCE_ATTRIBUTES'] = self.otel_resource_attributes

    def process(self, element):
        # 此时可以正常使用uplight-telemetry包
        import uplight_telemetry
        # 业务逻辑处理...
        yield processed_element

然后在流水线中传递参数:

result = (
        pipeline
        | "ReadPendingRecordsFromDB" >> read_from_db
        | "Parse input PCollection" >> beam.Map(ProcessBillRequests.parse_bill_data_requests)        
        | "Fetch bills " >> beam.ParDo(ProcessBillRequests.FetchBillInformation(
            otel_service_name=pipeline_options.OTEL_SERVICE_NAME,
            otel_resource_attributes=pipeline_options.OTEL_RESOURCE_ATTRIBUTES
        ))
)

3. 通过setup.py脚本设置环境变量

如果使用setup_file配置,可以在setup.py中添加一个安装后执行的脚本,在Worker安装依赖时设置环境变量。示例如下:

from setuptools import setup

def set_env_vars():
    import os
    os.environ['OTEL_SERVICE_NAME'] = os.getenv('OTEL_SERVICE_NAME', '')
    os.environ['OTEL_RESOURCE_ATTRIBUTES'] = os.getenv('OTEL_RESOURCE_ATTRIBUTES', '')

setup(
    # 其他配置...
    cmdclass={'install': set_env_vars}
)

注意:这种方式需要确保环境变量能传递到setup脚本的执行环境中,适用性有限,仅作为备选方案。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 20:46:07