如何让Dask多Worker并行运行单个自定义任务?
如何让Dask集群协作完成单个哈希挖矿任务?
背景
已在云端部署带有scheduler和worker的Dask集群,配置如下(5个worker,每个1CPU/1线程/1GB内存):
cluster = DropletCluster(n_workers=5, region="sgp1", image="ubuntu-22-04-x64", size="s-1vcpu-1gb-amd") client = Client(cluster)
自定义哈希挖矿函数,目标是找到哈希前缀为000000的数值:
def Tester_3(): found = False i = 0 start = time.time() while not found: result_sha256 = sha256_blockchain(i) found = bool(re.search(r"^000000", result_sha256)) if found: break i+=1 return f"------Running for {time.time() - start} sekon------"
之前的尝试问题
尝试1:client.run - 每个worker独立执行完整任务
调用代码:
xx = client.run(Tester_3) xx
结果是每个worker各自从头开始找哈希,完成独立的任务,而非协作:
{'tls://143.198.x.x:33785': '------Running for 32.62452673912048 sekon------', 'tls://159.223.x.x:43919': '------Running for 33.88207650184631 sekon------', 'tls://159.65.x.x:44949': '------Running for 37.02353572845459 sekon------', 'tls://159.65.x.x:43181': '------Running for 31.585288524627686 sekon------', 'tls://206.189.x.x:38597': '------Running for 42.807809591293335 sekon------'}
尝试2:dask.delayed/client.submit - 仅单个worker执行任务
调用代码:
z = dask.delayed(Tester_3)() print(z.compute())
结果仅用单个worker完成任务,无法利用集群资源:
------Running for 37.95407176017761 sekon------ CPU times: user 50.2 ms, sys: 23 ms, total: 73.2 ms Wall time: 38 s
尝试3:multiprocessing/PoolExecutor - 仅本地机器执行
任务完全在本地运行,没有用到集群的worker节点。
解决方案:拆分任务实现协作
要让所有worker协作完成同一个哈希查找任务,核心是拆分任务区间,让每个worker处理不重叠的数值范围,同时监听结果并及时终止未完成任务。
步骤1:重构任务函数
修改函数,让它接受起始值和步长,每个worker处理独立的数值序列(比如worker0处理0,5,10...,worker1处理1,6,11...),避免重复计算:
import time import hashlib import re from dask.distributed import Client, as_completed # 补全原代码中未给出的哈希计算函数 def sha256_blockchain(i): return hashlib.sha256(str(i).encode()).hexdigest() # 重构后的子任务函数:处理从start开始,步长为step的数值序列 def find_target_hash(start, step, start_time): i = start while True: current_hash = sha256_blockchain(i) if re.match(r"^000000", current_hash): return { "found_by": start, "target_i": i, "target_hash": current_hash, "total_time": round(time.time() - start_time, 2) } i += step
步骤2:提交并行任务并监听结果
将拆分后的子任务提交给集群,一旦有worker找到目标哈希,立即取消其他所有未完成的任务:
# 假设已连接到Dask集群 # cluster = DropletCluster(n_workers=5, region="sgp1", image="ubuntu-22-04-x64", size="s-1vcpu-1gb-amd") # client = Client(cluster) n_workers = 5 global_start_time = time.time() # 给每个worker分配唯一的起始值,步长等于worker数量 futures = [ client.submit(find_target_hash, worker_id, n_workers, global_start_time) for worker_id in range(n_workers) ] # 监听任务完成状态,拿到第一个结果后终止其他任务 for future in as_completed(futures): result = future.result() print(f"找到目标结果!耗时{result['total_time']}秒") print(f"Worker起始ID:{result['found_by']},对应数值i:{result['target_i']},哈希:{result['target_hash']}") # 取消所有未完成的任务 for f in futures: if not f.done(): f.cancel() break
方案说明
- 任务拆分:每个worker处理不重叠的数值序列,避免重复计算,最大化利用集群算力
- 结果监听:通过
as_completed实时获取第一个完成的任务结果,及时终止其他任务,避免无效计算 - 资源利用:所有集群worker都参与同一个哈希查找任务,真正实现协作
内容的提问来源于stack exchange,提问作者Mayer Reflino
相关产品推荐
相关产品推荐

