如何用Apache Beam Python Pipeline读取云存储中的Excel文件?
Apache Beam读取GCS中Excel文件并转为PCollection的解决方案
问题场景
尝试用Apache Beam Python Pipeline读取云存储(GCS)中的Excel文件,使用Pandas读取后无法转换为PCollection,运行代码时报错:
TypeError: unsupported operand type(s) for >>: 'str' and 'NoneType'
原代码如下:
def read_data_from_excel_file(): bucket_name = "nidec-ga-transient" blob_name = "ConcessoesERestricoes.xlsx" storage_client = storage.Client() bucket = storage_client.bucket(bucket_name) blob = bucket.blob(blob_name) data_bytes = blob.download_as_bytes() df = pd.read_excel(data_bytes, 'Lista de Gargalos') return df Pipeline = ( pipeline_load_data | "Importar Dados CloudStorage" >> read_data_from_excel_file() # | "Write_to_BQ" >> beam.io.WriteToBigQuery( # tabela, # schema=table_schema, # write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, # create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, # custom_gcs_temp_location = 'gs://ddc-test-262213-staging/henrique.klock@dojo.technology/temp' ) )
错误原因
Beam的>>操作符要求右侧必须是PTransform类型(如beam.Create、beam.ParDo等),但原代码直接调用read_data_from_excel_file()返回的是Pandas DataFrame,并非合法的PTransform,导致类型不匹配触发错误。同时原代码未正确使用Beam Pipeline的上下文管理(with beam.Pipeline()),进一步加剧了类型问题。
可行解决方法
方法1:小文件场景 - 转为字典列表生成PCollection
直接读取Excel后转为字典列表,通过beam.Create转换为PCollection,适合文件体积较小的情况:
import apache_beam as beam import pandas as pd from google.cloud import storage def read_excel_from_gcs(): bucket_name = "nidec-ga-transient" blob_name = "ConcessoesERestricoes.xlsx" storage_client = storage.Client() bucket = storage_client.bucket(bucket_name) blob = bucket.blob(blob_name) data_bytes = blob.download_as_bytes() # 读取指定sheet并转为字典列表(每行对应一个字典) df = pd.read_excel(data_bytes, sheet_name='Lista de Gargalos') return df.to_dict('records') # 正确使用Beam Pipeline上下文 with beam.Pipeline() as pipeline_load_data: pcollection = ( pipeline_load_data | "加载Excel数据为PCollection" >> beam.Create(read_excel_from_gcs()) | "写入BigQuery" >> beam.io.WriteToBigQuery( tabela, schema=table_schema, write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, custom_gcs_temp_location='gs://ddc-test-262213-staging/henrique.klock@dojo.technology/temp' ) )
方法2:大文件场景 - 自定义ParDo分块处理
如果Excel文件体积较大,一次性读取可能导致内存溢出,可通过ParDo分块处理:
import apache_beam as beam import pandas as pd from google.cloud import storage from io import BytesIO class ReadExcelDoFn(beam.DoFn): def process(self, element): bucket_name, blob_name = element storage_client = storage.Client() bucket = storage_client.bucket(bucket_name) blob = bucket.blob(blob_name) # 用BytesIO包装字节流,支持分块读取(可根据需求扩展分块逻辑) with BytesIO(blob.download_as_bytes()) as f: df = pd.read_excel(f, sheet_name='Lista de Gargalos') for _, row in df.iterrows(): yield row.to_dict() with beam.Pipeline() as pipeline_load_data: pcollection = ( pipeline_load_data | "传入文件路径" >> beam.Create([("nidec-ga-transient", "ConcessoesERestricoes.xlsx")]) | "读取并解析Excel" >> beam.ParDo(ReadExcelDoFn()) | "写入BigQuery" >> beam.io.WriteToBigQuery( tabela, schema=table_schema, write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, custom_gcs_temp_location='gs://ddc-test-262213-staging/henrique.klock@dojo.technology/temp' ) )
关键注意事项
- 必须确保Beam Pipeline操作链中,
>>右侧是合法的PTransform,不能直接返回DataFrame或其他非Transform类型数据。 - 处理大文件时优先选择
ParDo分块读取,避免内存压力。 - 确保安装依赖库:
pip install apache-beam pandas google-cloud-storage openpyxl(openpyxl用于读取.xlsx格式文件)
内容的提问来源于stack exchange,提问作者Henrique Klock
相关产品推荐
相关产品推荐

