如何让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

