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

Python多线程脚本设置队列maxsize后卡顿的问题排查

问题分析与解决方案

卡顿原因

设置maxsize=10后,UrlConverter.load()方法在主线程中往队列塞URL,当队列被填满10个元素时,queue.put()会阻塞主线程——此时还没有任何线程消费队列内容,要等所有URL都塞进队列后才会启动消费线程,导致主线程一直卡在put()步骤无法继续。

修复方案

需要把URL加载操作放到独立线程中,让消费线程和加载线程并行运行;同时正确使用task_done()和join()管理队列任务:

修改后的client.py代码

import socket
from pathlib import Path
from threading import Thread
from queue import Queue


class UrlConverter:
    def load(self, filename: str, queue: Queue):
        urls_file_path = str(Path(__file__).parent / Path(filename))
        with open(urls_file_path, 'r', encoding="utf-8") as txt_file:
            for line in txt_file:
                line = line.strip()
                if line:  # 跳过空行
                    queue.put(line)
        # 所有URL加载完成后,塞入结束标记
        queue.put(None)


class Client:
    def execute(self, url: str, queue: Queue):
        try:
            sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
            sock.connect(('localhost', 59999))
            sock.send(url.encode())
            data = sock.recv(1024)
            sock.close()
            print(f"{url}: {data.decode('utf-8')}\n")
        except Exception as e:
            print("Error sending request to", url, ":", str(e))
        finally:
            # 任务完成后通知队列
            queue.task_done()


class ClientThreads:
    def __init__(self, n: int):
        self.n = n

    def worker(self, queue: Queue):
        # 持续从队列取任务,直到收到结束标记
        while True:
            url = queue.get()
            if url is None:
                # 收到结束标记,退出线程
                queue.task_done()
                break
            client = Client()
            client.execute(url, queue)

    def execute(self, queue):
        # 创建指定数量的工作线程
        threads = []
        for _ in range(self.n):
            t = Thread(target=self.worker, args=(queue,))
            t.start()
            threads.append(t)
        # 等待队列所有任务完成(包括结束标记)
        queue.join()
        # 等待所有线程退出
        for t in threads:
            t.join()


def main():
    # 初始化带maxsize的队列
    urls_queue = Queue(maxsize=10)
    url_converter = UrlConverter()
    clients_threads = ClientThreads(10)

    # 启动加载URL的线程
    load_thread = Thread(target=url_converter.load, args=('urls.txt', urls_queue))
    load_thread.start()

    # 启动消费线程处理任务
    clients_threads.execute(urls_queue)

    # 等待加载线程完成
    load_thread.join()


if __name__ == "__main__":
    main()

关键修改点说明

  • 分离加载与消费线程:把URL加载操作放到独立线程,主线程同时启动消费线程,队列满时put()会阻塞加载线程,但消费线程会不断取走数据,队列有空位后加载线程会继续工作。
  • task_done()与join()的正确使用:
    • 每个工作线程完成URL发送后,调用queue.task_done()告知队列该任务已完成。
    • 调用queue.join()会阻塞,直到队列中所有任务都被标记为完成,确保所有URL都被处理。
  • 结束标记机制:加载完所有URL后,往队列塞一个None作为结束标记,工作线程收到标记后自动退出,避免线程无限等待。
  • 优化线程逻辑:替换原有的批量启动/等待线程模式,改为让线程持续从队列取任务,直到收到结束信号,更符合多线程队列的高效使用规范。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 16:33:13