如何调试Polars批量处理S3 Parquet文件时的挂起问题?
回答
一、定位挂起位置的方法
启用Polars调试级日志
除了set_verbose,配置更详细的日志级别,输出Polars内部执行细节:import logging logging.basicConfig(level=logging.DEBUG) pl.Config.set_log_level("debug") pl.Config.set_fmt_str_lengths(100)这会记录读取文件、join、聚合等阶段的进度,帮助锁定挂起前的最后一个操作。
系统级进程跟踪
- Linux/macOS:用
strace -p <进程PID>跟踪系统调用,查看进程是否卡在IO(S3网络/本地磁盘)或锁竞争;用htop监控CPU、内存、IO使用率,判断是资源耗尽还是无响应。 - Windows:用Process Explorer查看线程状态,或
procmon跟踪文件/网络操作。
- Linux/macOS:用
手动添加阶段日志
在关键步骤前后插入日志,明确执行进度:import logging logging.basicConfig(level=logging.INFO, format="%(asctime)s - %(message)s") def join_files(basedir): logging.info(f"Start processing dir: {basedir}") df_a = pl.scan_parquet(os.path.join(basedir, "file_a.parquet")) logging.info(f"Loaded file_a for {basedir}") df_a = df_a.join(ref, on="id_1") logging.info(f"Joined ref with file_a for {basedir}") df_b = pl.scan_parquet(os.path.join(basedir, "file_b.parquet")) logging.info(f"Loaded file_b for {basedir}") result = df_a.join(df_b, on="id_a").select(["id_1", "id_b"]).unique() logging.info(f"Completed dir: {basedir}") return result def process_batch(dirs): logging.info(f"Processing batch with {len(dirs)} dirs") dfs = [join_files(d) for d in dirs] logging.info(f"Generated {len(dfs)} lazy DFs") concat_df = pl.concat(dfs, rechunk=False) logging.info(f"Concatenated DFs") unique_df = concat_df.unique() logging.info(f"Applied unique") logging.info(f"Starting collect()") result = unique_df.collect() logging.info(f"Completed collect()") return result隔离测试缩小范围
将挂起过的批次单独重复运行,看是否能复现;逐步减小批次大小,定位触发挂起的特定目录或组合。检查依赖版本
旧版本Polars或s3fs可能存在死锁/挂起bug,升级到最新稳定版:pip install --upgrade polars s3fs
二、更优的实现方式
优化参考表加载
若reference.parquet是小数据集,直接加载为内存DataFrame,避免惰性缓存的重复IO:ref = pl.read_parquet("reference.parquet").select("id_1").unique()后续join时直接使用内存中的
ref,无需每次读取缓存。调整批次处理逻辑
- 逐目录处理并追加结果:避免大量惰性DF concat后一次性collect,减少内存压力:
# 初始化空结果文件 pl.DataFrame({"id_1": [], "id_b": []}).write_parquet("final_result.parquet") def process_dir(basedir): df_a = pl.read_parquet(os.path.join(basedir, "file_a.parquet")) df_a = df_a.join(ref, on="id_1") df_b = pl.read_parquet(os.path.join(basedir, "file_b.parquet")) return df_a.join(df_b, on="id_a").select(["id_1", "id_b"]).unique() for basedir in all_dirs: dir_result = process_dir(basedir) # 追加到结果文件 current_result = pl.read_parquet("final_result.parquet") pl.concat([current_result, dir_result]).unique().write_parquet("final_result.parquet") - 分批生成临时文件再合并:若内存仍紧张,每处理N个目录生成一个临时文件,最后统一合并:
batch_size = 15 temp_files = [] all_dirs = [...] # 你的S3目录列表 for i in range(0, len(all_dirs), batch_size): batch_dirs = all_dirs[i:i+batch_size] batch_result = process_batch(batch_dirs) temp_path = f"temp_batch_{i}.parquet" batch_result.write_parquet(temp_path) temp_files.append(temp_path) # 合并所有临时文件 final_df = pl.concat([pl.read_parquet(f) for f in temp_files]).unique() final_df.write_parquet("final_result.parquet") # 清理临时文件 import os for f in temp_files: os.remove(f)
- 逐目录处理并追加结果:避免大量惰性DF concat后一次性collect,减少内存压力:
优化S3读取配置
配置s3fs的超时与重试策略,避免因网络问题导致挂起:import s3fs fs = s3fs.S3FileSystem( client_kwargs={ "connect_timeout": 15, "read_timeout": 45, "retries": {"max_attempts": 6} } ) # 读取时指定自定义文件系统 df_a = pl.scan_parquet(os.path.join(basedir, "file_a.parquet"), filesystem=fs)减少重复去重操作
当前目录级和全局级各做一次unique,可去掉目录级的unique,仅在最终合并时执行一次,降低计算开销:def join_files(basedir): df_a = pl.scan_parquet(os.path.join(basedir, "file_a.parquet")).join(ref, on="id_1") df_b = pl.scan_parquet(os.path.join(basedir, "file_b.parquet")) return df_a.join(df_b, on="id_a").select(["id_1", "id_b"])
内容的提问来源于stack exchange,提问作者jsc
相关产品推荐
相关产品推荐

