You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.21 08:50:13