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

Python BoundedSemaphore获取后进程卡住问题排查求助

子进程获取信号量后卡住问题排查与解决

我在排查无法通过调试暂停Python进程的问题时,通过逐行记录执行日志定位子进程卡住的原因,发现子进程在获取BoundedSemaphore信号量后立即卡住。

我的队列包含1000-2000条记录,execute_function单独运行2000次完全正常,但批量运行时偶尔会在执行若干次后卡住。从日志看,卡住的进程最后一条日志是'sema acquired',也就是成功获取信号量后,后续代码完全没执行。

相关代码如下:

import queue
from multiprocessing import Process, BoundedSemaphore

my_queue = queue.Queue()
my_semaphore = BoundedSemaphore(10)

def execute_function(obj):
    log('execute function started')
    # 业务逻辑代码
    do the thing...

def start_function(obj, my_semaphore):
    log('start_function started')
    my_semaphore.acquire()
    log('sema acquired')
    execute_function(obj)
    log('function executed')
    my_semaphore.release()
    log('sema released')

for object in objects:
    my_queue.put(object)

threads = []
while not my_queue.empty():
    obj = my_queue.get()
    thread = Process(target=start_function, args=[obj, my_semaphore])
    threads.append(thread)
    thread.start()

for t in threads:
    t.join()

核心问题

这里的关键错误是跨进程传递BoundedSemaphore对象:

  • Python的multiprocessing中,信号量是进程间同步原语,但直接将父进程的信号量对象传给子进程,会导致每个子进程拿到的是该对象的副本,而非共享实例。
  • 这会导致信号量的计数完全失效,子进程间无法真正实现并发控制,甚至可能出现死锁或卡住的情况。

解决方法

改用multiprocessing.Manager()创建可跨进程共享的信号量:

from multiprocessing import Process, Manager, BoundedSemaphore
import queue

def execute_function(obj):
    log('execute function started')
    # 业务逻辑代码
    do the thing...

def start_function(obj, my_semaphore):
    log('start_function started')
    my_semaphore.acquire()
    log('sema acquired')
    execute_function(obj)
    log('function executed')
    my_semaphore.release()
    log('sema released')

if __name__ == '__main__':
    manager = Manager()
    my_semaphore = manager.BoundedSemaphore(10)
    my_queue = queue.Queue()
    objects = [...]  # 替换为你的对象列表

    for object in objects:
        my_queue.put(object)

    threads = []
    while not my_queue.empty():
        obj = my_queue.get()
        thread = Process(target=start_function, args=(obj, my_semaphore))
        threads.append(thread)
        thread.start()

    for t in threads:
        t.join()

额外优化建议

  • 避免使用queue.Queue搭配多进程,改用multiprocessing.Queue,它是专门为跨进程通信设计的,更安全可靠。
  • 可以用Pool来替代手动管理进程和信号量,Pool自带并发控制,代码更简洁:
from multiprocessing import Pool

def execute_function(obj):
    log('execute function started')
    # 业务逻辑代码
    do the thing...

if __name__ == '__main__':
    objects = [...]  # 替换为你的对象列表
    # 限制并发数为10
    with Pool(10) as pool:
        pool.map(execute_function, objects)

内容的提问来源于stack exchange,提问作者user2396640

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 22:30:49