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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 07:18:15