Dask任务内存耗尽求助:1-10^9质数计算任务优化方案
问题与解决方案
问题场景
使用3个18GB内存、40GB存储的节点搭建Dask集群,计算1到10^9的质数时触发内存耗尽,哪怕只是把全量数字列表scatter到节点都会OOM,目标是计算质数并保存到文本文件。
核心问题分析
- 本地预加载全量数据:代码中把
numbers.txt的所有内容读到本地lines列表,直接撑爆本地内存,后续scatter也会一次性把巨量数据推给集群,超出worker内存上限。 - 质数判断算法低效:
checkPrime遍历到n-1,时间复杂度高且占用不必要的计算资源,间接加剧内存压力。 - 全量结果拉取到本地:
client.gather把所有计算结果拉回本地,1到10^9的质数约5000万条,完全无法在单节点内存存储。
解决方案
1. 懒加载数据,避免本地内存爆炸
用Dask Bag直接读取文件,无需本地预加载全量数据,让Dask自动分片分发到集群:
import dask.bag as db # 直接用Dask Bag读取文件,自动分片 b = db.read_text("numbers.txt").map(lambda x: int(x.strip()))
2. 优化质数判断算法
将遍历范围缩小到√n,提前排除偶数,大幅降低计算量和内存占用:
import math def checkPrime(n): if n <= 1: return None # 处理偶数,除了2 if n == 2: return n if n % 2 == 0: return None # 只遍历奇数到sqrt(n) for i in range(3, int(math.sqrt(n)) + 1, 2): if n % i == 0: return None return n
3. 分布式计算+直接写入结果
避免全量gather,用Dask的分布式操作直接将结果写入文件,让worker各自处理分片并写入:
# 链式操作:过滤非质数,然后保存 (b.map(checkPrime) .filter(lambda x: x is not None) .to_textfiles("output_*.txt")) # Dask会自动生成分片文件 # 等待任务完成 dask.compute()
4. 调整Dask集群内存配置
启动SSHCluster时指定worker的内存限制,避免单个worker占用过多内存:
from dask.distributed import SSHCluster, Client cluster = SSHCluster( ["10.195.4.221", "10.195.6.186", "10.195.6.164"], worker_options={"memory_limit": "16GB"} # 留2GB给系统进程 ) client = Client(cluster)
完整修改后代码
import dask import dask.bag as db import math from dask.distributed import SSHCluster, Client def checkPrime(n): if n <= 1: return None if n == 2: return n if n % 2 == 0: return None for i in range(3, int(math.sqrt(n)) + 1, 2): if n % i == 0: return None return n # 启动集群并设置内存限制 cluster = SSHCluster( ["10.195.4.221", "10.195.6.186", "10.195.6.164"], worker_options={"memory_limit": "16GB"} ) client = Client(cluster) # 懒加载数据,分布式处理 b = db.read_text("numbers.txt").map(lambda x: int(x.strip())) (b.map(checkPrime) .filter(lambda x: x is not None) .to_textfiles("output_*.txt")) # 触发计算 dask.compute() print("质数计算完成,结果已保存为output_*.txt")
内容的提问来源于stack exchange,提问作者Mujtaba Faizi
相关产品推荐
相关产品推荐

