Python多线程环境下使用Queue模块出现阻塞问题求助
Python多线程环境下使用Queue模块出现阻塞问题求助
嗨,我仔细看了你的代码,马上就找到问题根源了——你手动操作了Queue内部的锁和条件变量,这直接破坏了Queue原本封装好的线程安全机制,才导致生产者的put操作莫名其妙卡住,哪怕队列是空的也不行。
问题出在哪?
Queue类本身已经把所有线程安全的细节都封装好了,比如put()、get()方法内部会自动处理锁和条件变量(比如not_empty、not_full这些内部属性),完全不需要我们手动去碰这些底层的东西。
看你消费者的clean_converted_data方法里这段代码:
with self.inqueue.not_empty: self.inqueue.not_empty.wait_for(self.inqueue.full())
这两段代码的问题非常关键:
- 你直接获取了Queue的
not_empty锁,并且在等待过程中一直持有它; wait_for(self.inqueue.full())意味着消费者会一直等待队列被填满(你的队列maxsize是8),但锁没释放的情况下,Queue内部的正常逻辑被彻底打乱,导致生产者的put操作无法正常完成,哪怕队列是空的也会阻塞。
修复方案:改用Queue原生方法,别碰内部锁
你完全不需要手动处理锁和条件变量,直接用Queue提供的get()方法就好——它会自动阻塞直到队列中有数据,线程安全且可靠。重写后的clean_converted_data方法如下:
def clean_converted_data(self) -> List[Optional[float]]: data = [] last_val = None initial_setup_done = False while True: # 直接用get(),自动等待队列有数据,无需手动处理锁 record = self.inqueue.get() logger.info(f"Consumer thread record: {record}") if record is None: # 收到哨兵值,结束消费 return data if not initial_setup_done: last_val = record data.append(last_val) initial_setup_done = True else: last_val = record data.append(last_val) # 可选:调用task_done,配合queue.join()使用(好习惯,即使你现在没用到) self.inqueue.task_done()
额外的优化建议
- 批量读取数据:你的生产者每次只读取1条记录,效率比较低,可以考虑调整
fetch_surface_data的num_read参数,一次读取多条(比如一次读8条,和队列maxsize匹配),然后批量put到队列里,减少线程切换的开销。 - SQLite连接线程安全:注意SQLite的连接不能被多个线程共享,你的代码里连接只在生产者线程中使用,这没问题,但如果后续有其他线程要操作这个连接,一定要重新创建连接。
- 哨兵值处理:生产者最后put的
None是正确的结束信号,但要确保不会重复put这个值,避免消费者逻辑异常。
按照这个修改后,生产者的put操作应该就能正常工作,不会再卡住啦!
备注:内容来源于stack exchange,提问作者mraabhijit
相关产品推荐
相关产品推荐

