使用Apache Beam to_dataframe时遇AttributeError:BmsSchema无element_type属性
错误原因分析
AttributeError: 'BmsSchema' object has no attribute 'element_type' 是因为误用了to_dataframe()函数:
apache_beam.dataframe.convert.to_dataframe()的作用是接收整个PCollection,将其转换为Beam分布式DataFrame,而非处理单个数据元素。- 你在
PandasTransform这个DoFn的process方法中传入的是单个BmsSchema对象,单个元素没有element_type属性,因此触发错误。
解决方案
根据需求提供两种可行修改方案:
方案一:改用Beam DataFrame API直接处理(推荐)
直接在Pipeline中用Beam DataFrame处理整个数据集,无需自定义DoFn:
import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions, StandardOptions from apache_beam.io.gcp.bigquery import WriteToBigQuery from apache_beam.dataframe.convert import to_dataframe, to_pcollection import os import typing os.environ["GOOGLE_APPLICATION_CREDENTIALS"] = "...........json" # Define the schema class BmsSchema(typing.NamedTuple): ident: str beam.coders.registry.register_coder(BmsSchema, beam.coders.RowCoder) # Do function for passing in pubsub message from the subscription class ParsePubSubMessage(beam.DoFn): def process(self, message): import json all_columns = ['ident'] main_dict = dict(zip(all_columns, [None] * len(all_columns))) record = json.loads(message.decode('utf-8')) main_dict.update(record) yield {all_columns[0]: main_dict[all_columns[0]]} def run(): options = PipelineOptions( project='dw.......', runner='DirectRunner', streaming=True, temp_location='gs://............', staging_location='gs://..........', region='europe..........', job_name='pipeline-dataflow-test' ) options.view_as(StandardOptions).streaming = True input_subscription = 'projects/......./subscriptions/...........' table_schema = { "fields": [ {"name": "ident", "type": "STRING", "mode": "NULLABLE"} ] } with beam.Pipeline(options=options) as p: messages = (p | 'Read from PubSub' >> beam.io.ReadFromPubSub(subscription=input_subscription) | 'Parse PubSub Message' >> beam.ParDo(ParsePubSubMessage()) | 'Attaching the schema' >> beam.Map(lambda x: BmsSchema(**x)).with_output_types(BmsSchema) ) # 用Beam DataFrame处理整个PCollection df = to_dataframe(messages) # 这里可添加Pandas转换逻辑,例如:df['ident'] = df['ident'].str.upper() transformed_pcoll = to_pcollection(df, BmsSchema) transformed_pcoll | 'Write to BigQuery' >> WriteToBigQuery( table='project.table_name', write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, schema=table_schema, custom_gcs_temp_location='gs://........' ) if __name__ == '__main__': run()
方案二:在DoFn中处理批量元素(适合细粒度控制场景)
如果必须用DoFn结合Pandas,需先将元素批量分组,再在process中处理批量数据:
import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions, StandardOptions from apache_beam.io.gcp.bigquery import WriteToBigQuery import pandas as pd import os import typing os.environ["GOOGLE_APPLICATION_CREDENTIALS"] = "...........json" # Define the schema class BmsSchema(typing.NamedTuple): ident: str beam.coders.registry.register_coder(BmsSchema, beam.coders.RowCoder) # Do function for passing in pubsub message from the subscription class ParsePubSubMessage(beam.DoFn): def process(self, message): import json all_columns = ['ident'] main_dict = dict(zip(all_columns, [None] * len(all_columns))) record = json.loads(message.decode('utf-8')) main_dict.update(record) yield {all_columns[0]: main_dict[all_columns[0]]} # 修改后的PandasTransform,处理批量元素 class PandasTransform(beam.DoFn): def process(self, batch_elements): # 将批量元素转为DataFrame df = pd.DataFrame(batch_elements) # 这里添加Pandas转换逻辑,例如数据清洗、格式转换 # 示例:df['ident'] = df['ident'].str.strip() # 转换回字典并输出 for _, row in df.iterrows(): yield row.to_dict() def run(): options = PipelineOptions( project='dw.......', runner='DirectRunner', streaming=True, temp_location='gs://............', staging_location='gs://..........', region='europe..........', job_name='pipeline-dataflow-test' ) options.view_as(StandardOptions).streaming = True input_subscription = 'projects/......./subscriptions/...........' table_schema = { "fields": [ {"name": "ident", "type": "STRING", "mode": "NULLABLE"} ] } with beam.Pipeline(options=options) as p: messages = (p | 'Read from PubSub' >> beam.io.ReadFromPubSub(subscription=input_subscription) | 'Parse PubSub Message' >> beam.ParDo(ParsePubSubMessage()) | 'Attaching the schema' >> beam.Map(lambda x: BmsSchema(**x)).with_output_types(BmsSchema) | 'Batch elements' >> beam.BatchElements(min_batch_size=10, max_batch_size=100) # 批量分组 | 'Transforming columns with pandas' >> beam.ParDo(PandasTransform()) ) messages | 'Write to BigQuery' >> WriteToBigQuery( table='project.table_name', write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, schema=table_schema, custom_gcs_temp_location='gs://........' ) if __name__ == '__main__': run()
关键说明
- 方案一利用Beam原生DataFrame集成,更符合分布式处理模型,适合大多数场景。
- 方案二通过
BatchElements将单个元素打包成批量,再用Pandas处理,适合自定义批量逻辑的场景。
内容的提问来源于stack exchange,提问作者oyugi.collins
相关产品推荐
相关产品推荐

