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

如何解决Asyncio中‘Task was destroyed but it is pending’警告?

搞定asyncio的"Task was destroyed but it is pending!"警告

问题出在哪?

  1. 写文件逻辑太折腾:你把async with aiofiles.open塞到了while True循环里,等于每处理一条数据就打开、关闭一次文件,不仅效率低,还容易引发异常。
  2. 画蛇添足的任务取消:main函数里你在await write_q.join()之后主动调用consumer.cancel(),但这时候write_task已经拿到哨兵值None,正要自行退出,强行取消会打断任务收尾流程,导致任务处于pending状态被销毁,触发警告。
  3. 队列与任务配合混乱:哨兵值本就是用来让write_task自动退出的,await write_q.join()会等待所有队列项处理完成,之后write_task会自然结束,完全不需要手动取消。

修正步骤

1. 把文件打开移到循环外部

全程只打开一次文件,避免反复IO操作,既高效又稳定。

2. 移除不必要的cancel操作

让consumer任务在收到哨兵值后自然退出,不要强行终止。

3. 给task_done加保险

用finally块确保无论数据处理是否出错,都标记队列项完成,避免队列join一直等待。

改好的完整代码

import asyncio
import aiofiles
import aiocsv
import json

async def task(l: list, write_q: asyncio.Queue) -> None:
    # 从数据源拿数据,包装成请求任务丢进队列
    for i in l:
        req: dict = {
            "headers": {"Accept": "application/json"},
            "url": "https://httpbin.org/post",
            "data": i
        }
        await write_q.put(req)
    
    # 丢个哨兵值,告诉消费者没数据了
    await write_q.put(None)

async def write_task(write_q: asyncio.Queue) -> None:
    headers: bool = True
    # 把文件打开移到循环外面,全程只开一次
    async with aiofiles.open("file.csv", mode="a+", newline='') as f:
        w = aiocsv.AsyncWriter(f)
        while True:
            data = await write_q.get()
            try:
                if not data:
                    # 拿到哨兵值,直接退出循环
                    break
                
                if headers:
                    # 写表头
                    await w.writerow(["status", "data"])
                    headers = False
                
                # 写数据行
                await w.writerow(["200", json.dumps(data)])
            finally:
                # 不管成功失败,都标记这个队列项处理完了
                write_q.task_done()
        
        # 退出前刷新文件,确保数据都写进去
        await f.flush()

async def main() -> None:
    # 造点测试用的假数据
    items: list[list[str]] = [["hello", "world"], ["asyncio", "test"]] * 5 

    write_q = asyncio.Queue()
    producer = asyncio.create_task(task(items, write_q))
    consumer = asyncio.create_task(write_task(write_q))

    # 等生产者把数据都丢进队列,顺便抓异常
    errors = await asyncio.gather(producer, return_exceptions=True)
    print(f"INFO: 生产者完工! 异常信息: {errors}")

    # 等队列里所有数据都处理完
    await write_q.join()
    # 等消费者自己跑完,别瞎取消
    await consumer
    print("INFO: 写入任务完工! ")
    print("INFO: 全部搞定!")

if __name__ == "__main__":
    loop = asyncio.new_event_loop()
    loop.run_until_complete(main())

关键点说明

  • 文件操作优化:一次打开文件全程复用,减少IO开销,避免文件指针异常移动。
  • finally保task_done:确保队列的join能正常结束,不会因异常卡住。
  • 取消操作移除:让消费任务自然退出,彻底解决pending警告。
  • 等待任务结束:用await consumer确保消费任务完全收尾,流程更严谨。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 22:55:54