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

如何在多个Spark节点上运行线程池?

问题

我之前用Python的multiprocessing和concurrent.futures模块实现过多进程与多线程结合的解决方案,但multiprocessing模块在Spark中仅能在driver节点运行,因此不得不改用sc.parallelize将工作负载分发至worker节点,现在遇到了一些问题,希望获得排查建议。

我目前仅尝试将包含100个值的列表分配至10个worker节点,并在每个节点上实现线程处理,代码如下:

from concurrent.futures import ThreadPoolExecutor
import time

def multithread(task, l):
   with ThreadPoolExecutor() as executor:
      results = list(executor.map(task, l))
   return results

def square(x):
   time.sleep(1)
   return x**2

def partition(l, n):
   # 该函数将输入列表分割为'n'个块
   for i in range(0, len(l), n):
      yield l[i:i +n]

num = list(range(100))
workers = 10
chunks = list(partition(num, workers))
rdd = sc.parallelize(chunks, numSlices=workers)
results = rdd.map(lambda x: multithread(x)).collect()

更新说明

除Ahmad提供的解决方案外,若要使workers变量动态化,请不要使用os.cpu_count(),而是需要通过Spark的getConf()方法获取该值,同时建议将其传入ThreadPoolExecutor()的max_workers参数。


内容的提问来源于stack exchange,提问作者chicagobeast12

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 07:47:17