如何通过PipelineOptions配置Dataflow的use_unsupported_python_version参数
问题:Dataflow配置忽略Python版本限制不生效
我正在尝试使用Google Dataflow在两个BigQuery表之间传输数据,代码如下:
import apache_beam as beam from apache_beam.io.gcp.internal.clients import bigquery from apache_beam.options.pipeline_options import PipelineOptions import argparse def parseArgs(): parser = argparse.ArgumentParser() parser.add_argument( '--experiment', default='use_unsupported_python_version', help='This does not seem to do anything.') args, beam_args = parser.parse_known_args() return beam_args def beamer(rows=[]): if len(rows) == 0: return project = 'myproject-474601' gcs_temp_location = 'gs://my_temp_bucket/tmp' gcs_staging_location = 'gs://my_temp_bucket/staging' table_spec = bigquery.TableReference( projectId=project, datasetId='mydataset', tableId='test') beam_options = PipelineOptions( parseArgs(), # 该配置不生效 project=project, runner='DataflowRunner', job_name='unique-job-name', temp_location=gcs_temp_location, staging_location=gcs_staging_location, use_unsupported_python_version=True, # 该配置也不生效 experiment='use_unsupported_python_version' # 该配置同样不生效 ) with beam.Pipeline(options=beam_options) as p: quotes = p | beam.Create(rows) quotes | beam.io.WriteToBigQuery( table_spec, # custom_gcs_temp_location = gcs_temp_location, # 不需要该配置? method='FILE_LOADS', write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED) return if __name__ == '__main__': beamer(rows=[{'id': 'ein', 'value': None, 'year': None, 'valueHistory': [{'year': 2021, 'amount': 900}]}])
但运行时收到报错,提示Dataflow不支持我使用的Python版本,报错信息如下:
Exception: Dataflow runner currently supports Python versions ['3.6', '3.7', '3.8'], got 3.9.7 (default, Sep 16 2021, 08:50:36) [Clang 10.0.0 ]. To ignore this requirement and start a job using an unsupported version of Python interpreter, pass --experiment use_unsupported_python_version pipeline option.
根据报错提示,我尝试了三种配置方式:通过argparse传入实验参数、直接向PipelineOptions传递use_unsupported_python_version=True参数、传递experiment='use_unsupported_python_version'参数,还参考了官方管道选项文档中的参数合并方法,但均未生效,仍然报相同的版本不支持错误。请问如何正确配置才能让Dataflow使用我当前的Python版本运行作业?
解决方案
配置不生效的核心原因是参数名称和格式错误,按以下方式修改即可:
- 修正实验参数格式:Apache Beam中实验参数的键为复数形式
experiments,取值为列表类型,不能直接传字符串。同时不需要单独传use_unsupported_python_version参数。 - 修正argparse逻辑(如果保留命令行传参的话):你之前的parseArgs函数只返回了系统解析到的beam参数,自定义的experiment参数没有被合并到最终配置里。
最简修改方案(无需保留命令行传参)
直接修改PipelineOptions的配置,把错误的experiment参数替换为正确的格式即可:
beam_options = PipelineOptions( project=project, runner='DataflowRunner', job_name='unique-job-name', temp_location=gcs_temp_location, staging_location=gcs_staging_location, experiments=['use_unsupported_python_version'] )
同时可以把之前无效的parseArgs()、use_unsupported_python_version=True、experiment='use_unsupported_python_version'这几行配置删掉。
如果需要保留命令行传参的方案
修改parseArgs函数,把自定义参数合并到返回的beam参数列表中:
def parseArgs(): parser = argparse.ArgumentParser() parser.add_argument( '--experiment', default='use_unsupported_python_version', help='Experiment flag to allow unsupported Python version') args, beam_args = parser.parse_known_args() # 把自定义实验参数加到beam参数列表中 beam_args.extend(['--experiments', args.experiment]) return beam_args
此时PipelineOptions里不需要再额外传experiments参数,直接传入parseArgs()的返回值即可。
注意事项
- 该配置仅允许你本地用高版本Python提交作业,Dataflow worker的默认运行环境还是使用官方支持的Python版本。如果你用到了Python3.9的新语法,需要额外配置自定义容器镜像匹配本地Python版本,否则运行时可能出现兼容性问题。
- 条件允许的话优先使用Dataflow官方支持的Python版本开发作业,可避免很多不必要的兼容性风险。
内容的提问来源于stack exchange,提问作者zakdances
相关产品推荐
相关产品推荐

