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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 18:01:16