使用Modin+Ray处理超大CSV时Ray Workers因OOM被杀死的解决方法
问题描述
使用Modin + Ray读取56GB、15亿行的预排序CSV文件(已通过Linux sort完成排序),执行groupby('1').sum()时出现Worker内存不足(OOM)被杀死,且计算效率低下。机器配置为48核、126GB内存,因防火墙限制无法访问Ray Web界面排查问题。
原始代码
import modin.pandas as pd import ray ray.init() df = pd.read_csv("./file", index_col=0, header=None, names=["1", "2"]) df.groupby('1').sum()
Ray环境信息
RayContext(dashboard_url='127.0.0.1:8265', python_version='3.8.10', ray_version='2.3.0', ray_commit='cf7a56b4b0b648c324722df7c99c168e92ff0b45', address_info={'node_ip_address': 'XXXXXXX', 'raylet_ip_address': 'XXXXXXX', 'redis_address': None, 'object_store_address': '/tmp/ray/session_2023-04-18_10-24-54_203554_3133/sockets/plasma_store', 'raylet_socket_name': '/tmp/ray/session_2023-04-18_10-24-54_203554_3133/sockets/raylet', 'webui_url': '127.0.0.1:8265', 'session_dir': '/tmp/ray/session_2023-04-18_10-24-54_203554_3133', 'metrics_export_port': XXXXXX, 'gcs_address': 'XXXX', 'address': 'XXXXXXX', 'dashboard_agent_listen_port': XXXXXX, 'node_id': 'XXXXXXXXXXXXXXXXXXXXXXXXXX'})
解决方案
1. 精准配置Ray资源限制
默认Ray会尽可能占用系统资源,容易导致内存过载。手动指定资源参数,预留系统内存,避免Worker抢占资源:
import ray # 根据48核、126G内存适配:分配40核给Ray,留8核给系统;总内存限制100G,留26G给系统 ray.init( num_cpus=40, object_store_memory=20 * 1024**3, # 20G共享对象存储池,存储中间计算结果 _memory=100 * 1024**3, # Ray整体内存上限 worker_process_memory_gb=2.5 # 单个Worker内存上限,防止单进程OOM )
2. 优化Modin CSV读取参数
利用文件已排序特性,配合读取参数减少内存占用,提升读取效率:
import modin.pandas as pd df = pd.read_csv( "./file", header=None, names=["1", "2"], index_col=0, dtype={"1": str, "2": float}, # 显式指定数据类型,避免自动推断占用额外内存 engine="pyarrow", # 使用pyarrow引擎,读取速度更快、内存开销更低 low_memory=False # 禁用低内存模式,适合大文件批量读取 )
- 若
"1"是整数类型,可指定int64进一步压缩内存占用;显式dtype是避免内存突增的关键。
3. 利用预排序特性优化GroupBy
文件已按"1"排序,可跳过全局洗牌(shuffle)步骤,大幅降低内存开销:
方案A:用Ray Data做局部聚合(最可靠)
from ray.data import read_csv # 用Ray Data读取已排序CSV ds = read_csv( "./file", header=None, column_names=["1", "2"], dtype={"1": str, "2": float} ) # 分块局部聚合→全局聚合,避免全局洗牌 def partial_sum(partition): return partition.groupby("1")["2"].sum().reset_index() partial_results = ds.map_groups(partial_sum).to_pandas() final_result = partial_results.groupby("1")["2"].sum()
- 每个分块先做局部聚合,中间数据量大幅减少,彻底避免全局洗牌的内存压力。
方案B:Modin GroupBy参数优化
# 因数据已排序,设置sort=False跳过排序步骤,节省计算资源 result = df.groupby('1', sort=False).sum()
4. 无Dashboard下的内存监控
用Ray命令行工具替代Dashboard排查问题:
- 查看Ray节点资源状态:
ray status - 查看内存使用详情:
ray memory - 查看Worker日志:前往Ray会话目录(即
/tmp/ray/session_XXXXXX)下的logs文件夹,检查Worker日志中的OOM报错,定位内存溢出环节。
5. 极端情况:分阶段拆分处理
若上述方案仍OOM,按"1"的前缀拆分文件,分块处理后合并结果:
# 按第一列前2个字符拆分文件(适配字符串类型的key) awk -F ',' '{print >> $1"_part.csv"}' ./file
import modin.pandas as pd import glob all_results = [] for f in glob.glob("*_part.csv"): df_part = pd.read_csv(f, header=None, names=["1", "2"], dtype={"1": str, "2": float}) part_sum = df_part.groupby("1")["2"].sum() all_results.append(part_sum) final_result = pd.concat(all_results).groupby("1")["2"].sum()
内容的提问来源于stack exchange,提问作者Ranger
相关产品推荐
相关产品推荐

