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

如何用Polars创建查询计划并行读取S3中的多个小JSON文件?

Polars读取S3存储桶中多个小JSON文件的实现方法

完全可以用Polars实现类似Spark的功能,通过构建延迟执行的查询计划批量读取S3前缀下的多个小JSON文件,以下是两种常用方式:

方式一:直接使用通配符路径

Polars支持在路径中使用通配符匹配文件,写法和Spark类似,scan_glob会自动创建并行读取的查询计划:

import polars as pl

# 创建查询计划(延迟执行)
query_plan = pl.scan_glob("s3://my-bucket/path/to/smallfiles/*.json")

# 添加其他数据处理操作(可选)
query_plan = query_plan.filter(pl.col("column_name") > 100)

# 执行查询,加载数据
df = query_plan.collect()

scan_glob会识别所有匹配的JSON文件,并且按照Polars的并行处理逻辑读取,和你查阅的文档中提到的多文件并行处理逻辑一致。

方式二:手动指定文件列表(精细控制)

如果需要先筛选特定文件再构建查询,可以先通过S3工具(如boto3)获取目标文件列表,再批量构建查询计划:

import polars as pl
import boto3

# 获取S3上的目标JSON文件列表
s3_client = boto3.client("s3")
bucket = "my-bucket"
prefix = "path/to/smallfiles/"

response = s3_client.list_objects_v2(Bucket=bucket, Prefix=prefix)
file_paths = [
    f"s3://{bucket}/{obj['Key']}" 
    for obj in response["Contents"] 
    if obj["Key"].endswith(".json")
]

# 拼接多个文件的查询计划
query_plan = pl.concat([pl.scan_json(path) for path in file_paths])

# 执行查询
df = query_plan.collect()

这种方式适合需要自定义文件筛选规则的场景,同样采用延迟执行,仅在collect()时才会实际读取和处理数据。

关键说明

  • Polars的scan_*系列方法均为延迟执行,仅创建查询计划,不会立即加载数据,这和Spark的懒执行逻辑完全一致。
  • 确保环境已配置好S3访问权限(如环境变量、AWS凭证文件等),Polars会自动适配S3文件系统的读取。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 09:10:53