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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 01:20:09