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

如何使用Apache Beam Python SDK从S3读取Parquet文件到PCollection?

使用Apache Beam Python SDK从S3读取Parquet文件为PCollection的最佳方式

核心结论

apache_beam.io.parquetio.ReadFromParquet 支持直接从S3读取Parquet文件,并且会直接返回符合需求的PCollection,这是最贴合Beam编程模型的最佳实现方式,无需单独依赖s3io模块做底层文件操作。

前置准备

  • 安装必要依赖:确保已安装包含AWS扩展的Beam包,以及Parquet解析依赖
    pip install apache-beam[aws] pyarrow
    
  • 配置AWS访问权限:可通过以下任一方式配置
    • 设置环境变量:AWS_ACCESS_KEY_ID 和 AWS_SECRET_ACCESS_KEY
    • 本地AWS配置文件(~/.aws/credentials)
    • 在Pipeline选项中显式指定凭证(不推荐在代码中硬编码)

代码示例

import apache_beam as beam
from apache_beam.io.parquetio import ReadFromParquet

def main():
    # 初始化Pipeline选项
    pipeline_options = beam.options.pipeline_options.PipelineOptions()
    aws_options = pipeline_options.view_as(beam.options.pipeline_options.AwsOptions)
    
    # 若需显式指定凭证(可选,优先用环境变量/配置文件)
    # aws_options.access_key_id = "YOUR_ACCESS_KEY"
    # aws_options.secret_access_key = "YOUR_SECRET_KEY"

    with beam.Pipeline(options=pipeline_options) as p:
        # 直接从S3路径读取Parquet,返回PCollection
        parquet_pcoll = p | "Read S3 Parquet Files" >> ReadFromParquet(
            file_pattern="s3://your-bucket/path/to/files/*.parquet",
            # 可选:提前传入PyArrow Schema,避免自动推断schema出错
            # schema=your_predefined_pyarrow_schema
        )

        # 后续可对PCollection进行任意处理,例如打印样本数据
        parquet_pcoll | "Print Sample Records" >> beam.Map(print)

if __name__ == "__main__":
    main()

为什么不用s3io?

s3io模块的定位是底层S3文件操作工具,它返回的是文件流或字节对象,而非Beam的PCollection。如果用s3io手动读取文件,还需要自行实现Parquet解析、并行化处理等逻辑,这不仅重复造轮子,还违背了Beam的分布式处理设计思想。

注意事项

  • 对于复杂Schema的Parquet文件,建议提前用PyArrow定义Schema并传入ReadFromParquet,避免自动推断出现偏差
  • 确保运行Pipeline的环境(本地/云服务)拥有目标S3桶的读写权限
  • ReadFromParquet会自动并行处理S3中的多个Parquet文件,适合大规模数据集的高效读取

内容的提问来源于stack exchange,提问作者Christos Hadjinikolis

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 18:48:38