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

基于可变列表创建线程遇阻塞,求持续监听解决方案

问题原因

原脚本的核心问题是**ThreadPoolExecutor上下文管理器配合executor.map()会阻塞主线程**:每次调用executor.map()时,主线程会等待所有线程里的handle_event任务全部完成才会退出上下文,直接导致main()的循环被卡住,没法按时继续获取新的元素列表。

解决方案

我们需要让线程池持久化运行,并且异步提交任务(不等待任务结束),这样主线程就能持续执行循环,按时获取新元素。具体修改如下:

修改后的代码

from concurrent.futures import ThreadPoolExecutor
from time import sleep
import asyncio
import Client

client = Client()
# 初始化一个持久化线程池,设置合理的最大线程数(可根据实际需求调整)
thread_pool = ThreadPoolExecutor(max_workers=10)

def handle_event(event):
    try:
        for i in range(10):
            client.get_info(event)
            sleep(60)
    except Exception as e:
        # 捕获异常避免单个任务崩溃影响其他线程
        print(f"处理事件 {event} 时出错: {e}")

async def main():
    while True:
        entries = client.get_new_entry()

        if entries:
            # 逐个提交任务到线程池,提交后立即返回,不阻塞主线程
            for event in entries:
                thread_pool.submit(handle_event, event)

        await asyncio.sleep(60)

if __name__ == "__main__":
    try:
        loop = asyncio.new_event_loop()
        loop.run_until_complete(main())
    finally:
        # 程序退出时等待所有任务完成后关闭线程池
        thread_pool.shutdown(wait=True)

关键改动说明

  • 持久化线程池:把线程池移到全局初始化,避免每次循环创建销毁,同时设置固定的max_workers,防止短时间内创建过多线程耗尽系统资源。
  • 用submit替代map:submit()提交任务后会立即返回Future对象,不会阻塞主线程,主线程可以继续执行后续的休眠和新元素获取逻辑。
  • 异常捕获:在handle_event里加异常处理,避免单个任务的错误导致线程崩溃,影响其他正在运行的监控任务。
  • 优雅关闭:程序退出时调用shutdown(wait=True),确保所有正在执行的监控任务完成后再关闭线程池。

额外优化建议

如果你的Client提供了异步接口(比如async def get_info()),可以完全用异步IO实现,不需要线程池,效率会更高:

import asyncio
import Client

client = Client()

async def handle_event(event):
    try:
        for i in range(10):
            await client.get_info(event)
            await asyncio.sleep(60)
    except Exception as e:
        print(f"处理事件 {event} 时出错: {e}")

async def main():
    while True:
        # 如果get_new_entry是同步方法,建议用run_in_executor包装成异步
        # entries = await asyncio.get_running_loop().run_in_executor(None, client.get_new_entry)
        entries = client.get_new_entry()

        if entries:
            for event in entries:
                # 创建异步任务,不阻塞主线程
                asyncio.create_task(handle_event(event))

        await asyncio.sleep(60)

if __name__ == "__main__":
    asyncio.run(main())

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 12:55:22