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

Python 3.10 asyncio服务器递归转发时任务被销毁问题求助

Asyncio P2P节点转发请求时任务被销毁的问题解决

问题背景

我使用asyncio开发P2P节点,节点可将特定请求转发至其他节点。当节点A转发请求到节点B,节点B又转发回节点A时,处理请求的任务有时会被取消,出现Task was destroyed but it was pending!错误,该错误发生在读取asyncio.start_server提供的reader时。

复现代码

import sys
import asyncio

global neighbor
neighbor = None

async def make_request(host, port, n):
    task = asyncio.current_task()
    reader, writer = await asyncio.open_connection(host, port)
    print(f"Making request to {(host, port)}")

    writer.write(n.to_bytes(4, 'big'))
    await writer.drain()

    task.set_name(task.get_name()+"_here")
    #
    ###
    ##### FAILS HERE
    ###
    #
    resp_len = await reader.read(1)
    task.set_name(task.get_name()+"_not_here")
    resp_len = int.from_bytes(resp_len, 'big')
    resp = await reader.read(resp_len)
    print(f"Response from {(host, port)}: {resp}")

    writer.close()
    await writer.wait_closed()
    return resp_len, resp


async def handle_request(reader, writer):
    print(f'Got request from {writer.get_extra_info("peername")}')
    n = int.from_bytes(await reader.read(4), 'big')
    if n == 0:
        writer.write(b'\x04\x01\x02\x03\x04')
        await writer.drain()
    else:
        resp_len, resp = await make_request(*neighbor, n-1)
        writer.write(resp_len.to_bytes(1, 'big') + resp)
        await writer.drain()
    writer.close()
    await writer.wait_closed()


async def run_server(host='127.0.0.1', port=0):
    server = await asyncio.start_server(handle_request, host, port)
    addrs = ', '.join(str(sock.getsockname()) for sock in server.sockets)
    print(f"Serving on {addrs}")

    async with server:
        await server.serve_forever()


if __name__=="__main__":
    if sys.argv[1] == 'node1':
        neighbor = ('127.0.0.1', 7776)
        asyncio.run(run_server(port=8889))
    elif sys.argv[1] == 'node2':
        neighbor = ('127.0.0.1', 8889)
        asyncio.run(run_server(port=7776))
    elif sys.argv[1] == 'client':
        for i in range(10000):
            asyncio.run(make_request('127.0.0.1', 7776, n=2))

错误现象

在三个终端分别运行参数为"node1"、"node2"、"client"的程序后,node2的输出如下:

Serving on ('127.0.0.1', 7776)
Got request from ('127.0.0.1', 46620)
Making request to ('127.0.0.1', 8889)
Got request from ('127.0.0.1', 46630)
Response from ('127.0.0.1', 8889): b'\x01\x02\x03\x04'
Got request from ('127.0.0.1', 46646)
Making request to ('127.0.0.1', 8889)
Got request from ('127.0.0.1', 46652)
Response from ('127.0.0.1', 8889): b'\x01\x02\x03\x04'
Got request from ('127.0.0.1', 46656)
Making request to ('127.0.0.1', 8889)
Got request from ('127.0.0.1', 46670)
Response from ('127.0.0.1', 8889): b'\x01\x02\x03\x04'
Got request from ('127.0.0.1', 46678)
Making request to ('127.0.0.1', 8889)
Got request from ('127.0.0.1', 46686)
Response from ('127.0.0.1', 8889): b'\x01\x02\x03\x04'
Got request from ('127.0.0.1', 46700)
Making request to ('127.0.0.1', 8889)
Got request from ('127.0.0.1', 46706)
Response from ('127.0.0.1', 8889): b'\x01\x02\x03\x04'
Got request from ('127.0.0.1', 46722)
Making request to ('127.0.0.1', 8889)
Task was destroyed but it is pending!
task: <Task pending name='Task-24_here' coro=<handle_request() done, defined at /home/colm/Documents/Exodus/node_test.py:32> wait_for=<Future pending cb=[Task.task_wakeup()]>>
Got request from ('127.0.0.1', 46734)

node1无错误输出,client在多次打印请求和响应后挂起,停留在“Making request ...”阶段,推测是等待最后一个请求完成。该错误仅在节点将请求转发回请求来源节点时出现,可通过将client的make_request中n设为1复现。

问题原因

这个错误的核心是请求链路形成循环时,TCP连接被提前关闭,导致异步任务在等待reader数据时被销毁:

  • 当client发送n=2的请求给node2,node2转发n=1给node1,node1再转发n=0回node2
  • 此时node2的handle_request任务在处理node1的n=0请求时,会快速写完响应并关闭连接;但同时node2还有另一个任务(处理client请求的任务)正在等待node1的响应(也就是刚才那个n=0请求的响应)
  • 当网络调度时机巧合时,node1关闭连接的操作会导致node2的等待任务被中断,而asyncio会检测到这个未完成的pending任务被垃圾回收,从而抛出错误。

另外,asyncio.start_server会自动为每个新连接创建任务,这些任务的生命周期由asyncio管理,但如果连接被提前关闭,任务内的await操作就会处于pending状态,最终触发该错误。

解决方案

1. 捕获IO操作的异常

在读取reader时捕获连接相关异常,确保任务能正常结束,而非处于pending状态:

async def make_request(host, port, n):
    task = asyncio.current_task()
    try:
        reader, writer = await asyncio.open_connection(host, port)
        print(f"Making request to {(host, port)}")

        writer.write(n.to_bytes(4, 'big'))
        await writer.drain()

        task.set_name(task.get_name()+"_here")
        try:
            resp_len = await reader.read(1)
            if not resp_len:
                raise ConnectionResetError("Connection closed while reading response length")
            task.set_name(task.get_name()+"_not_here")
            resp_len = int.from_bytes(resp_len, 'big')
            resp = await reader.read(resp_len)
            if not resp:
                raise ConnectionResetError("Connection closed while reading response data")
            print(f"Response from {(host, port)}: {resp}")
        finally:
            writer.close()
            await writer.wait_closed()
        return resp_len, resp
    except (asyncio.CancelledError, ConnectionResetError, BrokenPipeError) as e:
        print(f"Request to {(host, port)} failed: {e}")
        return 0, b''

2. 确保handle_request任务优雅结束

在处理请求的核心逻辑外添加异常捕获,避免子任务失败导致主任务处于pending状态:

async def handle_request(reader, writer):
    peername = writer.get_extra_info("peername")
    print(f'Got request from {peername}')
    try:
        n_data = await reader.read(4)
        if not n_data:
            print(f"Connection closed by {peername}")
            return
        n = int.from_bytes(n_data, 'big')
        if n == 0:
            writer.write(b'\x04\x01\x02\x03\x04')
            await writer.drain()
        else:
            resp_len, resp = await make_request(*neighbor, n-1)
            writer.write(resp_len.to_bytes(1, 'big') + resp)
            await writer.drain()
    except Exception as e:
        print(f"Error handling request from {peername}: {e}")
    finally:
        writer.close()
        await writer.wait_closed()

3. 避免无意义的循环请求(可选)

在实际P2P场景中,建议给每个请求添加唯一ID,记录节点已处理过的请求ID,避免重复转发形成循环,从根源减少这类错误的触发概率。

验证效果

修改代码后,再次运行三个终端的程序,即使出现连接关闭的情况,任务也会正常捕获异常并结束,不会再出现Task was destroyed but it was pending!的错误,client也不会挂起。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 12:37:58