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

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函数,未出现此类错误。


问题原因与解决方案

核心问题分析

  1. 普通列表非异步安全:gift_queue用普通列表实现,而on_gift回调可能在独立线程/协程中执行,直接append/pop会引发线程安全问题,同时无法触发协程调度,导致任务阻塞。
  2. TikTokLiveClient停止不彻底:client.stop()可能只是发送停止信号,但未等待后台WebSocket任务(即错误中的_receive_frame_loop)结束,导致任务在销毁时仍处于pending状态。
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 11:10:52