Python 3.10 asyncio服务器递归转发时任务被销毁问题求助
问题背景
我使用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

