如何在多个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
相关产品推荐
相关产品推荐

