Python multiprocessing apply_async回调无法更新字典问题求助
问题原因
你提交完所有apply_async异步任务后,主进程没有等待所有子进程任务完成、回调函数执行完毕,就直接读取了全局字典token_counts。apply_async是非阻塞调用,主进程提交任务后会立刻继续执行后续代码,此时大部分回调还没来得及更新字典,导致count_multiprocessing拿到的是不完整的统计结果。
解决方案
在提交完所有任务后,调用pool.close()关闭进程池(不再接受新任务),再调用pool.join()等待所有子进程完成任务,确保所有回调函数都执行完毕后,再读取全局字典。
修正后的代码
import multiprocessing import random def count_tokens(document): counter = dict() for token in document: if token in counter: counter[token] += 1 else: counter[token] = 1 return counter tokens = ['tok'+str(i) for i in range(int(9))] catalog = [random.choices(tokens, k=8) for _ in range(100)] token_counts = {token: 0 for token in tokens} def callback(result): global token_counts for token, count in result.items(): token_counts[token] += count return token_counts with multiprocessing.Pool() as pool: for document in catalog: pool.apply_async(count_tokens, args=(document,), callback=callback) # 新增:关闭进程池并等待所有任务完成 pool.close() pool.join() count_multiprocessing = dict(**token_counts) print(count_multiprocessing) token_counts = {token: 0 for token in tokens} for document in catalog: callback(count_tokens(document)) count_onecpu = dict(**token_counts) print(count_onecpu) for token, count in count_onecpu.items(): assert count == count_multiprocessing[token] for token, count in count_multiprocessing.items(): assert count == count_onecpu[token]
关键改动说明
pool.close():阻止进程池接受新的任务,确保所有已提交的任务都能被处理。pool.join():让主进程暂停执行,直到进程池中所有子进程都完成任务,同时所有异步任务的回调函数也会在主进程中执行完毕,此时token_counts已经被完全更新。
内容的提问来源于stack exchange,提问作者FraSchelle
相关产品推荐
相关产品推荐

