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

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())

这两段代码的问题非常关键:

  1. 你直接获取了Queue的not_empty锁,并且在等待过程中一直持有它;
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 17:53:04