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

如何使用Dask或其他Python库过滤S3中超10亿行的Parquet表

Dask 优化可运行方案

方案1:广播Join优化(最易落地,无Shuffle开销)

你用到的df2仅需要1700万条唯一索引值,全量加载后总大小通常不超过500MB,完全可以广播到所有Worker节点,避免大表Join产生的海量Shuffle内存开销:

import dask
from dask.distributed import Client, LocalCluster

# 先调整Worker资源配置,避免单Worker内存不足
cluster = LocalCluster(
    n_workers=8,
    threads_per_worker=12,
    memory_limit="80GB",
    # 开启内存溢写阈值,提前把冷数据刷到磁盘
    memory_target_fraction=0.7,
    memory_spill_fraction=0.8,
    memory_pause_fraction=0.9
)
client = Client(cluster)
dask.config.set({"temporary_directory": "/大容量磁盘的临时目录路径"})

# 读取小表全量到本地内存
df2_local = dd.read_parquet(
    f"s3://{MY_BUCKET}/{PATH_TO_AUX_TABLE}/*",
    columns=[INDEX_COL]
).drop_duplicates().compute()

# 读取大表时控制单个分区大小,避免单分区过大
df1 = dd.read_parquet(
    f"s3://{MY_BUCKET}/{PATH_TO_HUGE_TABLE}/*",
    blocksize="128MB",
    split_row_groups=True
)

# 开启广播Join,所有Worker直接用本地的小表副本和大表分区匹配,无Shuffle
filtered_df1 = df1.merge(
    df2_local,
    on=INDEX_COL,
    how="inner",
    broadcast=True
)
filtered_df1.to_parquet(f"s3://{MY_BUCKET}/{PATH}", write_index=False)

方案2:优化isin过滤逻辑

之前的isin方案内存溢出,主要是因为大列表的isin查询效率低、多分区重复拷贝过滤集合导致,优化后逻辑如下:

# 读取小表唯一值后转成Python set,isin查询时间复杂度从O(n)降到O(1),内存占用也更低
df2 = dd.read_parquet(f"s3://{MY_BUCKET}/{PATH_TO_AUX_TABLE}/*")
filter_values = set(df2[INDEX_COL].unique().compute())

# 读取大表时控制分区大小
df1 = dd.read_parquet(
    f"s3://{MY_BUCKET}/{PATH_TO_HUGE_TABLE}/*",
    blocksize="128MB",
    split_row_groups=True
)

# 过滤集合作为参数传入map_partitions,避免全局重复拷贝
def filter_partition(partition, filter_vals):
    return partition[partition[INDEX_COL].isin(filter_vals)]
filtered_df1 = df1.map_partitions(filter_partition, filter_vals=filter_values, meta=df1.meta)

filtered_df1.to_parquet(f"s3://{MY_BUCKET}/{PATH}", write_index=False)

方案3:谓词下推提前过滤(性能最优)

如果你的Parquet文件带行组统计信息,可直接在读取阶段用PyArrow谓词下推过滤,不符合条件的行根本不会加载到内存:

import pyarrow as pa

filter_arr = pa.array(df2[INDEX_COL].unique().compute())
df1 = dd.read_parquet(
    f"s3://{MY_BUCKET}/{PATH_TO_HUGE_TABLE}/*",
    filters=[(INDEX_COL, "in", filter_arr)],
    blocksize="128MB"
)
df1.to_parquet(f"s3://{MY_BUCKET}/{PATH}", write_index=False)
其他可选Python工具方案

如果不想调优Dask参数,也可以用更适配单节点大内存场景的工具实现需求:

1. DuckDB

单节点OLAP引擎,原生支持直接读写S3上的Parquet文件,不需要全量加载数据到内存,代码更简洁性能更高:

import duckdb
con = duckdb.connect()
# 加载S3扩展
con.install_extension("httpfs")
con.load_extension("httpfs")
# 配置S3认证信息(公开桶可跳过)
con.execute("""
SET s3_region='你的S3区域';
SET s3_access_key_id='你的AK';
SET s3_secret_access_key='你的SK';
""")
# 直接执行SQL关联过滤并写入结果
con.execute(f"""
COPY (
    SELECT t1.* FROM read_parquet('s3://{MY_BUCKET}/{PATH_TO_HUGE_TABLE}/*') t1
    INNER JOIN (SELECT DISTINCT {INDEX_COL} FROM read_parquet('s3://{MY_BUCKET}/{PATH_TO_AUX_TABLE}/*')) t2
    ON t1.{INDEX_COL} = t2.{INDEX_COL}
) TO 's3://{MY_BUCKET}/{PATH}' (FORMAT PARQUET);
""")

2. PySpark

Shuffle溢写机制成熟,处理大表关联的稳定性比Dask更高,700GB内存单节点跑15亿行过滤完全没有压力。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 09:45:03