用RuntimeValueProvider调用Apache Beam read_csv创建Dataflow模板遇WontImplementError
尝试创建可在GCP Dataflow上运行的Apache Beam管道模板,使用Apache Beam DataFrame模块的read_csv读取文件,希望将文件名作为模板参数传入,因此采用RuntimeValueProvider实现,编写代码如下:
import apache_beam as beam from apache_beam.dataframe.io import read_csv from apache_beam.options.pipeline_options import PipelineOptions class MyOptions(PipelineOptions): @classmethod def _add_argparse_args(cls, parser): parser.add_value_provider_argument('--file_name', type=str, default= 'gs://default-bucket/default-file.csv') pipeline_options = PipelineOptions( runner='DataflowRunner', project='my-project', job_name='read-csv', temp_location='gs://dataflow-test-bucket/temp', region='us-central1') p = beam.Pipeline(options=pipeline_options) my_options = pipeline_options.view_as(MyOptions) # Hardcoding the file works fine: df = p | read_csv('gs://default-bucket/default-file.csv') df = p | read_csv(my_options.file_name) beam.dataframe.convert.to_pcollection(df) | beam.Map(print) p.run().wait_until_finish()
运行时出现错误:
Exception has occurred: WontImplementError
non-deferred
File "D:\WorkArea\dataflow_args_test_projects\read_csv.py", line 37, in
df = p | read_csv(my_options.file_name)
询问使用read_csv时正确调用RuntimeValueProvider的方式。
问题原因
Apache Beam DataFrame的read_csv方法目前不支持直接传入RuntimeValueProvider类型参数,因为该方法的底层实现未做延迟执行(deferred)适配,导致触发WontImplementError。
正确实现方式
需要通过原生Beam IO组件读取文件(这类组件支持RuntimeValueProvider),再将读取结果转换为Beam DataFrame。以下提供两种可行方案:
方案1:小文件场景 - 用ReadFromText读取后解析为DataFrame
适合文件体积较小的场景,通过收集所有文本行后一次性解析为CSV:
import apache_beam as beam from apache_beam.dataframe.convert import to_dataframe from apache_beam.options.pipeline_options import PipelineOptions import pandas as pd from io import StringIO class MyOptions(PipelineOptions): @classmethod def _add_argparse_args(cls, parser): parser.add_value_provider_argument( '--file_name', type=str, default='gs://default-bucket/default-file.csv' ) pipeline_options = PipelineOptions( runner='DataflowRunner', project='my-project', job_name='read-csv', temp_location='gs://dataflow-test-bucket/temp', region='us-central1' ) p = beam.Pipeline(options=pipeline_options) my_options = pipeline_options.view_as(MyOptions) # 用原生ReadFromText接收RuntimeValueProvider参数 lines = p | beam.io.ReadFromText(my_options.file_name) # 合并所有行并解析为DataFrame def parse_csv_content(line_collection): csv_str = '\n'.join(line_collection) return pd.read_csv(StringIO(csv_str)) df = lines | beam.CombineGlobally(parse_csv_content) | to_dataframe() # 后续处理逻辑 beam.dataframe.convert.to_pcollection(df) | beam.Map(print) p.run().wait_until_finish()
方案2:大文件场景 - 用ReadFromCsv读取后转换为DataFrame
适合大文件场景,利用Beam原生的ReadFromCsv组件直接解析CSV,再转换为DataFrame:
import apache_beam as beam from apache_beam.dataframe.convert import to_dataframe from apache_beam.options.pipeline_options import PipelineOptions from apache_beam.io.csv import ReadFromCsv class MyOptions(PipelineOptions): @classmethod def _add_argparse_args(cls, parser): parser.add_value_provider_argument( '--file_name', type=str, default='gs://default-bucket/default-file.csv' ) pipeline_options = PipelineOptions( runner='DataflowRunner', project='my-project', job_name='read-csv', temp_location='gs://dataflow-test-bucket/temp', region='us-central1' ) p = beam.Pipeline(options=pipeline_options) my_options = pipeline_options.view_as(MyOptions) # 用原生ReadFromCsv接收RuntimeValueProvider参数,直接解析CSV csv_pcoll = p | ReadFromCsv(my_options.file_name) # 将PCollection转换为Beam DataFrame df = to_dataframe(csv_pcoll) # 后续处理逻辑 beam.dataframe.convert.to_pcollection(df) | beam.Map(print) p.run().wait_until_finish()
核心逻辑说明
原生Beam IO组件(如ReadFromText、ReadFromCsv)的实现支持延迟执行,能够正确处理RuntimeValueProvider类型的参数。通过这类组件读取数据后,再通过to_dataframe()转换为Beam DataFrame,即可继续使用DataFrame的相关操作。
内容的提问来源于stack exchange,提问作者Manish

