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

如何让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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 22:35:51