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

