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

使用Polars将分区NDJSON写入AWS S3的技术问询

Polars 超大规模NDJSON分区写入AWS S3解决方案

问题1:无sink_json_cloud时直接写入S3(不落地本地)

Polars目前没有原生的sink_json_cloud方法,但可以通过流式写入S3对象流实现无本地落地的写入,核心是结合支持流式IO的S3客户端(如s3fs),搭配LazyFrame的分块处理能力,避免加载全量数据到内存:

  1. 安装依赖:pip install s3fs
  2. 流式写入实现代码:
import polars as pl
from s3fs import S3FileSystem

# 初始化S3客户端
s3 = S3FileSystem()

# 目标S3路径
s3_path = "s3://your-bucket/path/to/output.ndjson"

# 加载超大规模数据集为LazyFrame
lf = pl.scan_parquet("s3://your-bucket/large-source-dataset.parquet")

# 分批次流式写入S3
with s3.open(s3_path, "wb") as f:
    for batch in lf.collect_stream():
        batch.write_json(f, row_oriented=True)

若需按分区写入,可结合问题2的方案为每个分区单独创建S3流。

问题2:JSON分区写入的高效实现(替代多次过滤)

传统方案需提前收集分区列唯一值再循环过滤,在分区值极多的场景下内存压力大,以下两种方案可优化:

方案1:批次内分组写入

利用map_batches对每个数据批次按分区列分组,直接将每组数据写入对应S3路径,无需提前获取所有分区值,内存占用更低:

import polars as pl
from s3fs import S3FileSystem

s3 = S3FileSystem()
bucket = "your-bucket"
base_partition_path = "path/to/partitioned-data"

def write_partition_batch(batch: pl.DataFrame) -> None:
    # 按目标分区列分组
    for partition_val, group_df in batch.group_by("date"):
        # 构造分区路径(如s3://bucket/path/date=2024-01-01/data.ndjson)
        partition_s3_path = f"s3://{bucket}/{base_partition_path}/date={partition_val}/data.ndjson"
        # 追加写入对应分区文件
        with s3.open(partition_s3_path, "ab") as f:
            group_df.write_json(f, row_oriented=True)

# 加载LazyFrame并执行批次处理
lf = pl.scan_parquet("s3://your-bucket/large-source-dataset.parquet")
lf.map_batches(write_partition_batch).collect()

方案2:流式获取分区值并过滤

若需按分区单独处理全量数据,可通过collect_stream流式获取唯一分区值,避免一次性加载所有值到内存:

import polars as pl
from s3fs import S3FileSystem

s3 = S3FileSystem()
lf = pl.scan_parquet("s3://your-bucket/large-source-dataset.parquet")

# 流式迭代获取分区列唯一值(不会一次性加载所有值)
for value_batch in lf.select("date").unique().collect_stream():
    for date_val in value_batch["date"]:
        # 过滤当前分区的LazyFrame
        partition_lf = lf.filter(pl.col("date") == date_val)
        # 写入对应S3路径
        target_path = f"s3://your-bucket/path/date={date_val}/data.ndjson"
        with s3.open(target_path, "wb") as f:
            for batch in partition_lf.collect_stream():
                batch.write_json(f, row_oriented=True)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 14:03:09