ProcessPoolExecutor结合共享内存计数时在255处停滞的问题求解
问题原因及解决方案
核心问题1:共享内存存储类型溢出
你分配的共享内存仅1字节大小,默认对应无符号8位整数,取值范围为0-255。当计数器递增到255后,再加1会溢出为0,后续任务读取到0后又从1开始循环,自然无法达到500的目标值。
核心问题2:全局Lock未真正共享
你定义的全局Lock在多进程环境下(尤其是Windows系统),每个子进程会重新初始化该锁,导致互斥机制完全失效,可能引发多个进程同时读写计数器,出现数据覆盖或异常。
修改后的代码
from multiprocessing import shared_memory, Manager from concurrent.futures import ProcessPoolExecutor as Executor, as_completed import time, random import struct def counter(lock): existing_shm = shared_memory.SharedMemory(name='shm') # 读取共享内存中的4字节整数 old_state = struct.unpack('i', existing_shm.buf[:4])[0] with lock: time.sleep(random.random()/10) new_state = old_state + 1 # 将新值打包为4字节写入共享内存 existing_shm.buf[:4] = struct.pack('i', new_state) print(new_state) existing_shm.close() if __name__=='__main__': with Manager() as manager: # 创建跨进程共享的锁 lock = manager.Lock() with Executor(12) as p: # 分配4字节内存存储int32类型整数 shm = shared_memory.SharedMemory(create=True, size=4, name='shm') # 初始化计数器为0 shm.buf[:4] = struct.pack('i', 0) futures = [p.submit(counter, lock) for i in range(500)] for future in as_completed(futures): pass shm.close() shm.unlink()
关键修改点说明
- 共享内存大小调整:从1字节改为4字节,适配int32类型,支持更大的数值范围(最高可达2^31-1)
- 共享锁实现:使用
Manager().Lock()创建跨进程共享的锁,确保所有子进程使用同一锁实现互斥 - 整数读写处理:通过
struct模块的pack/unpack方法完成4字节整数的序列化与反序列化,彻底避免类型溢出问题
运行修改后的代码,计数器会正常递增到500,不会再出现停在255的情况。
内容的提问来源于stack exchange,提问作者nya_ko314
相关产品推荐
相关产品推荐

