Ray远程函数报FUNCTION_SIZE_ERROR_THRESHOLD超限错误排查
问题根因
报错和你提到的5MB大小的Tweets数据集无关,核心原因是PROB_SCORES远程函数隐式捕获了全局作用域下的预训练model和tokenizer对象。
Ray在分发远程任务到Worker节点前,会序列化函数本身及其闭包引用的所有变量。预训练NLP模型体积通常在数百MB级别,正好匹配报错中476MiB的函数体积,直接触发了默认95MiB的函数大小阈值限制。
解决方案
优先使用Ray Actor封装模型加载逻辑,从根源避免大对象随任务反复序列化,性能最优。
- 方案1:Ray Actor封装(推荐)
模型和分词器只会在每个Worker进程启动时加载一次,常驻Worker内存,不会随任务重复序列化传输,并行效率最高。
参考实现代码:
import ray import torch from scipy.special import softmax # 补充你自己需要的transformers导入逻辑 ray.init(num_cpus=16, ignore_reinit_error=True, object_store_memory=10**10) @ray.remote(num_cpus=1) class SentimentScorer: def __init__(self, model_path, tokenizer_path): # 初始化时加载模型和分词器,全程不参与任务序列化 self.tokenizer = # 替换为你的分词器加载代码,如AutoTokenizer.from_pretrained(tokenizer_path) self.model = # 替换为你的模型加载代码,如AutoModelForSequenceClassification.from_pretrained(model_path) self.model.eval() # 开启推理模式,降低内存占用 def predict_batch(self, text_list): # 批量推理逻辑,建议单批次大小控制在16-64条平衡调度开销和吞吐 with torch.no_grad(): # 关闭梯度计算,减少内存占用、提升推理速度 encoded = self.tokenizer(text_list, return_tensors='pt', padding=True, truncation=True) output = self.model(**encoded) scores = output[0].detach().numpy() scores = softmax(scores, axis=1) # 批量计算情感得分 prob_scores = scores[:,0] * (-1) + scores[:,-1] * 1 return prob_scores.tolist() # 初始化和CPU核数对应的Actor实例,每个Actor独立持有一份模型副本 NUM_WORKERS = 16 scorers = [SentimentScorer.remote(你的模型路径, 你的分词器路径) for _ in range(NUM_WORKERS)] # 将Tweets内容拆分为对应数量的批次,分发给不同Actor并行计算 content_list = Tweets['Content'].tolist() batch_size = len(content_list) // NUM_WORKERS + 1 batches = [content_list[i:i+batch_size] for i in range(0, len(content_list), batch_size)] tasks = [scorer.predict_batch.remote(batch) for scorer, batch in zip(scorers, batches)] results = ray.get(tasks) # 最终结果按批次顺序拼接即可 all_scores = [score for batch_res in results for score in batch_res]
- 方案2:
ray.put()预存大对象(适合临时调试)
如果不想重构为Actor结构,可以提前将大对象存入Ray对象存储,函数只传入对象引用,避免闭包打包大对象。注意该方案每次调用任务仍有对象引用的拉取开销,性能弱于Actor方案。
参考实现代码:
import ray import torch from scipy.special import softmax ray.init(num_cpus=16, ignore_reinit_error=True, object_store_memory=10**10) # 提前将大对象放入对象存储,拿到引用ID tokenizer_ref = ray.put(tokenizer) model_ref = ray.put(model) @ray.remote def PROB_SCORES(words, tokenizer_id, model_id): tokenizer = ray.get(tokenizer_id) model = ray.get(model_id) with torch.no_grad(): encoded_tweet = tokenizer(words, return_tensors='pt', padding=True, truncation=True) output = model(**encoded_tweet) scores = output[0][0].detach().numpy() scores = softmax(scores) prob_scores = scores[0] * (-1) + scores[-1] * 1 return prob_scores # 调用时传入对象引用,不要直接传大对象本身 ray.get([PROB_SCORES.remote(Tweets['Content'], tokenizer_ref, model_ref)])
- 避坑提示
- 不要直接调大
FUNCTION_SIZE_ERROR_THRESHOLD参数绕过报错,序列化数百MB的函数分发会带来极大的网络和序列化开销,并行效率极低。 - 你当前的写法只提交了单个远程任务,完全无法利用16核的并行能力,必须将数据拆分为多批次分发给多个任务/Worker才能实现并行加速。
- 推理时务必加
torch.no_grad(),否则会保存不必要的梯度张量,浪费内存拖慢速度。
- 不要直接调大
内容的提问来源于stack exchange,提问作者Inundata
相关产品推荐
相关产品推荐

