如何用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
相关产品推荐
相关产品推荐

