Python多进程内存优化:ECS任务OOM与进程挂起问题求助
JSON转Parquet内存优化与进程挂起解决方案
问题根源梳理
- 单进程Pandas处理:DataFrame占用内存无法及时释放,内存碎片累积导致OOM。
- 逐个创建Process:频繁创建销毁进程会耗尽系统内核资源(如进程描述符),导致任务挂起。
- Pool复用进程:子进程内Pandas的内存碎片、未释放对象随进程复用持续累积,重现内存问题。
具体解决方案
1. 带资源清理的限制并发进程池
使用multiprocessing.Pool控制并发数,同时在子进程内显式清理内存,避免进程复用导致的内存累积:
import gc import pandas as pd from multiprocessing import Pool def process_file(s3_json_file_key): # 读取JSON并执行扁平化逻辑 df = pd.read_json(s3_json_file_key) flattened_df = df.explode("nested_column") # 替换为你的实际扁平化代码 # 写入Parquet文件 flattened_df.to_parquet(f"{s3_json_file_key}.parquet") # 强制清理内存资源 del df, flattened_df gc.collect() if __name__ == "__main__": # 根据ECS配置设置并发数(如4核CPU设为4,内存不足则设为2) with Pool(processes=4) as pool: pool.map(process_file, s3_json_file_keys)
2. 周期性重启的进程池
用ProcessPoolExecutor结合内存监控,每处理一定数量文件或进程内存超标时重启进程池,彻底清除累积的内存碎片:
import gc import psutil import pandas as pd from concurrent.futures import ProcessPoolExecutor def get_process_memory_mb(): return psutil.Process().memory_info().rss / (1024 ** 2) def process_file(s3_json_file_key): df = pd.read_json(s3_json_file_key) flattened_df = df.explode("nested_column") flattened_df.to_parquet(f"{s3_json_file_key}.parquet") del df, flattened_df gc.collect() return get_process_memory_mb() if __name__ == "__main__": max_workers = 4 memory_threshold = 512 # 单个进程内存超过512MB时重启池 batch_size = 100 # 每处理100个文件重启池 executor = ProcessPoolExecutor(max_workers=max_workers) for idx, s3_key in enumerate(s3_json_file_keys): mem_usage = executor.submit(process_file, s3_key).result() # 触发重启条件 if (idx + 1) % batch_size == 0 or mem_usage > memory_threshold: executor.shutdown(wait=True) executor = ProcessPoolExecutor(max_workers=max_workers)
3. Pandas底层内存优化
从数据读取阶段减少内存占用,从源头降低压力:
- 读取JSON时指定
dtype,避免自动推断为大内存类型:dtype_spec = {"id": "int32", "category": "category"} df = pd.read_json(s3_json_file_key, dtype=dtype_spec) - 分块读取大JSON文件,避免一次性加载全量数据:
chunk_iter = pd.read_json(s3_json_file_key, chunksize=10000) for chunk in chunk_iter: flattened_chunk = chunk.explode("nested_column") flattened_chunk.to_parquet(f"{s3_json_file_key}.parquet", mode="append") - 扁平化后删除无用列,减少DataFrame体积:
flattened_df = flattened_df.drop(["unused_col1", "unused_col2"], axis=1)
4. 替代方案:用Dask实现并行处理
Dask DataFrame兼容Pandas语法,自动分块处理数据,内存占用可控,适合大规模文件处理:
import dask.dataframe as dd # 批量读取JSON文件(自动分块) ddf = dd.read_json(s3_json_file_keys) # 扁平化处理(与Pandas语法一致) flattened_ddf = ddf.explode("nested_column") # 写入Parquet(支持S3路径) flattened_ddf.to_parquet("s3://your-target-bucket/output-path/", write_index=False)
内容的提问来源于stack exchange,提问作者Rayeesahmed Shaikh
相关产品推荐
相关产品推荐

