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
相关产品推荐
相关产品推荐

