多进程含随机任务的结果可复现性问题求助
多进程随机任务结果可复现方案寻求
背景与需求
为提升收敛算法运行速度完成并行化改造,业务层面结果达标,但因审计及业务要求,需实现结果可复现。
场景说明
- 任务中使用了随机生成器
- 任务因大量条件分支导致执行时长不定,每次运行会分配至不同worker(代码用
sleep(random.random())模拟) - 实际参数:
NB_WORKER=10、NB_OBJECT=1030、NB_ITER=100 - 环境:Python 3.7、Windows
已尝试方案及问题
曾尝试为每个Object实例分配独立随机生成器,不确定该方案合理性,且发现耗时明显增加。以下是两种简化代码示例:
初始方案代码
from multiprocessing import Pool,get_context from time import sleep from numpy import random as nprng import random from datetime import datetime class Object(): def __init__(self,id_): self.id_=id_ self.value=0 def task(object_): sleep(random.random()) # task taking random time to execute (because lot of conditional) object_.value=rng.uniform() return object_ def init_worker(client_id,generators): global rng global worker_id with client_id.get_lock(): globals()['client_id'] = client_id.value worker_id=globals()['client_id'] rng=nprng.Generator(generators[globals()['client_id']]) client_id.value += 1 if __name__ == "__main__": # Init Pool and workers with random generators NB_PROCESS=5 NB_OBJECT=34 NB_ITER=3 ctx = get_context("spawn") client_ids = ctx.Value('i', 0) sequences = [nprng.SeedSequence((1209391983918, worker_id)) for worker_id in range(NB_PROCESS)] generators = [nprng.PCG64(seq) for seq in sequences] p=Pool(processes=NB_PROCESS,initializer=init_worker, initargs=(client_ids,generators,)) # Parallel task objects=[Object(i) for i in range(NB_OBJECT)] arg=[(object_,) for object_ in objects] start=datetime.now() total_sum=0 for i in range(NB_ITER): res=p.starmap(task,arg) total_sum+=sum([o.value for o in res]) print(f'Total sum is : {total_sum}') end=datetime.now() print(end-start)
为Object分配独立随机生成器的方案代码
from multiprocessing import Pool,get_context from time import sleep from numpy import random as nprng import random from datetime import datetime class Object(): def __init__(self,id_,generator): self.id_=id_ self.value=0 self.rng=nprng.Generator(generator) def task(object_): sleep(random.random()) #task taking random time to execute (because lot of conditional) object_.value=object_.rng.uniform() return object_ def init_worker(client_id): # global rng global worker_id with client_id.get_lock(): globals()['client_id'] = client_id.value client_id.value += 1 if __name__ == "__main__": # Init Pool and workers with random generators NB_PROCESS=5 NB_OBJECT=1030 NB_ITER=3 ctx = get_context("spawn") client_ids = ctx.Value('i', 0) p=Pool(processes=NB_PROCESS,initializer=init_worker, initargs=(client_ids,)) # Parallel task sequences = [nprng.SeedSequence((1209391983918, worker_id)) for worker_id in range(NB_OBJECT)] generators = [nprng.PCG64(seq) for seq in sequences] objects=[Object(i,generators[i]) for i in range(NB_OBJECT)] arg=[(object_,) for object_ in objects] start=datetime.now() total_sum=0 for i in range(NB_ITER): res=p.starmap(task,arg) total_sum+=sum([o.value for o in res]) print(f'Total sum is : {total_sum}') end=datetime.now() print(end-start)
注:该方案耗时增加,推测是对象访问开销导致。
需求
寻求其他实现多进程随机任务结果可复现的方案。
内容的提问来源于stack exchange,提问作者Progstud
相关产品推荐
相关产品推荐

