FastAPI+asyncio报错:Task was destroyed but it is pending求助
问题:asyncio + FastAPI WebSocket 断开后出现未销毁的Pending任务错误
我已花费48小时调试asyncio + FastAPI相关问题,恳请帮助。我注册了一个on_gift处理器,解析目标礼物并将其添加至gift_queue,随后通过循环从队列取出数据发送至WebSocket。我的需求是阻塞主线程直至客户端断开连接,再执行清理操作。目前清理流程已执行,但首次请求后出现如下错误:
Task was destroyed but it is pending! task: <Task pending name='Task-12' coro=<WebSocketProtocol13._receive_frame_loop() running at /Users/zane/miniconda3/lib/python3.9/site-packages/tornado/websocket.py:1106> wait_for=<Future pending cb=[IOLoop.add_future.<locals>.<lambda>() at /Users/zane/miniconda3/lib/python3.9/site-packages/tornado/ioloop.py:687, <TaskWakeupMethWrapper object at 0x7f9558a53af0>()]> cb=[IOLoop.add_future.<locals>.<lambda>() at /Users/zane/miniconda3/lib/python3.9/site-packages/tornado/ioloop.py:687]>
相关代码
@router.websocket('/donations') async def scan_donations(websocket: WebSocket, username): await websocket.accept() client = None gift_queue = [] try: client = TikTokLiveClient(unique_id=username, sign_api_key=api_key) def on_gift(event: GiftEvent): gift_cost = 1 print('[DEBUG] Could not find gift cost for {}. Using 1.'.format( event.gift.extended_gift.name)) try: if event.gift.streakable: if not event.gift.streaking: if gift_cost is not None: gift_total_cost = gift_cost * event.gift.repeat_count # if it cost less than 99 coins, skip it if gift_total_cost < 99: return gift_queue.append(json.dumps({ "type": 'tiktok_gift', "data": dict(name=event.user.nickname) })) else: if gift_cost < 99: return gift_queue.append(json.dumps({ 'type': 'tiktok_gift', 'data': dict(name=event.user.nickname) })) except: print('[DEBUG] Could not parse gift event: {}'.format(event)) client.add_listener('gift', on_gift) await client.start() while True: if gift_queue: gift = gift_queue.pop(0) await websocket.send_text(gift) del gift else: try: await asyncio.wait_for(websocket.receive_text(), timeout=1) except asyncio.TimeoutError: continue except ConnectionClosedOK: pass except ConnectionClosedError as e: print("[ERROR] TikTok: ConnectionClosedError {}".format(e)) pass except FailedConnection as e: print("[ERROR] TikTok: FailedConnection {}".format(e)) pass except Exception as e: print("[ERROR] TikTok: {}".format(e)) pass finally: print("[DEBUG] TikTok: Stopping listener...") if client is not None: print('Stopping TTL') client.stop() await websocket.close()
补充日志信息
Postman中断开连接后,日志输出:
[ERROR] TikTok: 1000 [DEBUG] TikTok: Stopping listener... Stopping TTL
看起来部分资源未正常退出。另有一个逻辑几乎完全相同的scan_chat函数,未出现此类错误。
问题原因与解决方案
核心问题分析
- 普通列表非异步安全:
gift_queue用普通列表实现,而on_gift回调可能在独立线程/协程中执行,直接append/pop会引发线程安全问题,同时无法触发协程调度,导致任务阻塞。 - TikTokLiveClient停止不彻底:
client.stop()可能只是发送停止信号,但未等待后台WebSocket任务(即错误中的_receive_frame_loop)结束,导致任务在销毁时仍处于pending状态。 - WebSocket监听逻辑冗余:while循环中用
wait_for(receive_text(), timeout=1)轮询,会持续创建短暂的pending任务,客户端断开时可能有未清理的任务残留。
修改后的代码
import asyncio import json from fastapi import WebSocket, APIRouter from TikTokLive import TikTokLiveClient from TikTokLive.events import GiftEvent from TikTokLive.exceptions import FailedConnection from fastapi.websockets import ConnectionClosedOK, ConnectionClosedError router = APIRouter() api_key = "your_api_key_here" @router.websocket('/donations') async def scan_donations(websocket: WebSocket, username: str): await websocket.accept() client = None gift_queue = asyncio.Queue() # 用于标记任务是否需要停止 stop_flag = asyncio.Event() async def process_gift_queue(): """独立协程处理队列发送,避免阻塞主线程""" while not stop_flag.is_set(): try: # 等待队列有数据,或停止信号触发 gift = await asyncio.wait_for(gift_queue.get(), timeout=0.5) await websocket.send_text(gift) gift_queue.task_done() except asyncio.TimeoutError: continue def on_gift(event: GiftEvent): gift_cost = 1 print('[DEBUG] Could not find gift cost for {}. Using 1.'.format( event.gift.extended_gift.name)) try: gift_data = None if event.gift.streakable: if not event.gift.streaking: if gift_cost is not None: gift_total_cost = gift_cost * event.gift.repeat_count if gift_total_cost < 99: return gift_data = json.dumps({ "type": 'tiktok_gift', "data": dict(name=event.user.nickname) }) else: if gift_cost < 99: return gift_data = json.dumps({ 'type': 'tiktok_gift', 'data': dict(name=event.user.nickname) }) if gift_data: # 用create_task异步添加到队列,避免同步回调阻塞 asyncio.create_task(gift_queue.put(gift_data)) except Exception as e: print('[DEBUG] Could not parse gift event: {}'.format(e)) try: client = TikTokLiveClient(unique_id=username, sign_api_key=api_key) client.add_listener('gift', on_gift) # 启动队列处理协程 queue_task = asyncio.create_task(process_gift_queue()) # 启动TikTok客户端 await client.start() # 阻塞主线程,直到WebSocket断开 while not stop_flag.is_set(): try: await websocket.receive_text() except ConnectionClosedOK: break except ConnectionClosedError as e: print("[ERROR] TikTok: ConnectionClosedError {}".format(e)) break except FailedConnection as e: print("[ERROR] TikTok: FailedConnection {}".format(e)) except Exception as e: print("[ERROR] TikTok: {}".format(e)) finally: print("[DEBUG] TikTok: Stopping listener...") # 设置停止信号,终止队列处理协程 stop_flag.set() # 停止TikTok客户端并等待后台任务结束 if client is not None: print('Stopping TTL') client.stop() # 短暂等待确保后台任务清理完成 await asyncio.sleep(0.5) # 等待队列处理协程结束 if 'queue_task' in locals(): await queue_task # 关闭WebSocket await websocket.close()
关键修改点
- 替换普通列表为
asyncio.Queue,保证异步环境下的线程安全与协程调度 - 新增
process_gift_queue独立协程处理队列发送,避免主线程阻塞 - 用
asyncio.create_task在on_gift回调中异步添加数据到队列,解决同步回调无法await的问题 - 添加
stop_flag事件统一管理任务停止,确保所有协程能优雅退出 - 停止客户端后等待短暂时间,确保tornado的WebSocket后台任务有足够时间清理
- 简化WebSocket监听逻辑,直接等待receive_text直到断开,减少冗余的timeout轮询
内容的提问来源于stack exchange,提问作者Zane Helton
相关产品推荐
相关产品推荐

