如何取消阻塞的await asyncio.StreamReader.readline()调用?
如何中止或取消阻塞的
await asyncio.StreamReader.readline()调用? 背景
我用asyncio.StreamReader处理基于行的持续通信,代码如下:
reader, writer = await asyncio.open_connection(host, port) while running: line = await reader.readline() process(line)
readline()会一直阻塞直到收到一行数据。我想中止这个阻塞操作,于是把调用包装进Task:
reader, writer = await asyncio.open_connection(host, port) while running: read_task = loop.create_task(reader.readline()) await read_task line = read_task.result() process(line)
打算从其他地方调用:
read_task.cancel()
来停止阻塞的readline(),但这个方法无效。
最小可复现示例
import asyncio import socket import threading import time read_task: asyncio.Task = None async def wait_for_data(): global read_task loop = asyncio.get_running_loop() # 创建一对连接的套接字 rsock, wsock = socket.socketpair() reader, writer = await asyncio.open_connection(sock=rsock) # 创建任务调用阻塞的readline read_task = loop.create_task(reader.readline()) print("等待一行数据...") await read_task data = read_task.result() # ==== 永远到不了这里 ==== print("收到:", data.decode()) writer.close() await writer.wait_closed() # 关闭第二个套接字 wsock.close() def thread_entry(): asyncio.run(wait_for_data()) if __name__ == '__main__': thread = threading.Thread(target=thread_entry) thread.start() time.sleep(2) # 给线程启动时间 print("取消读取任务") read_task.cancel() # 取消任务模拟超时 thread.join() print("完成。")
临时解决方案
我找到一种不够优雅的临时方案:用asyncio.wait_for设置超时:
reader, writer = await asyncio.open_connection(host, port) while running: try: line = await asyncio.wait_for(reader.readline(), timeout=5) except asyncio.TimeoutError: if not running: # 冗余但为了清晰 break continue process(line)
正确解决方案
问题出在StreamReader.readline()底层依赖的asyncio.StreamReader.read()不会响应任务取消——它的实现里没有检查取消请求的挂起点。以下是两种可靠的解决方式:
方法1:关闭底层套接字
当需要取消读取操作时,直接关闭StreamReader关联的套接字,这会让readline()抛出IO异常(比如ConnectionResetError),从而退出阻塞。
修改示例代码如下:
import asyncio import socket import threading import time read_task: asyncio.Task = None rsock: socket.socket = None # 保存套接字引用 async def wait_for_data(): global read_task, rsock loop = asyncio.get_running_loop() rsock, wsock = socket.socketpair() reader, writer = await asyncio.open_connection(sock=rsock) read_task = loop.create_task(reader.readline()) print("等待一行数据...") try: await read_task data = read_task.result() print("收到:", data.decode()) except (asyncio.CancelledError, ConnectionResetError): print("读取操作已取消") finally: writer.close() await writer.wait_closed() wsock.close() def thread_entry(): asyncio.run(wait_for_data()) if __name__ == '__main__': thread = threading.Thread(target=thread_entry) thread.start() time.sleep(2) print("取消读取任务并关闭套接字") if read_task: read_task.cancel() if rsock: rsock.close() # 关键:关闭套接字触发IO异常 thread.join() print("完成。")
方法2:用asyncio.wait监听取消信号和读取任务
创建一个asyncio.Event作为取消信号,用asyncio.wait同时等待读取任务和事件触发。一旦事件被设置,就取消读取任务并退出。
示例代码:
import asyncio import socket import threading import time stop_event = asyncio.Event() async def wait_for_data(): loop = asyncio.get_running_loop() rsock, wsock = socket.socketpair() reader, writer = await asyncio.open_connection(sock=rsock) read_task = loop.create_task(reader.readline()) print("等待一行数据...") # 同时等待读取任务和停止事件 done, pending = await asyncio.wait( [read_task, stop_event.wait()], return_when=asyncio.FIRST_COMPLETED ) if stop_event.is_set(): read_task.cancel() print("读取操作已取消") else: data = read_task.result() print("收到:", data.decode()) writer.close() await writer.wait_closed() wsock.close() def thread_entry(): asyncio.run(wait_for_data()) if __name__ == '__main__': thread = threading.Thread(target=thread_entry) thread.start() time.sleep(2) print("触发停止事件") # 在事件循环线程中设置事件(必须用call_threadsafe) loop = asyncio.get_event_loop_policy().get_event_loop() loop.call_soon_threadsafe(stop_event.set) thread.join() print("完成。")
为什么直接cancel()无效?
asyncio.Task.cancel()会向任务抛出CancelledError,但只有当任务中的代码处于可取消的挂起点时才会响应。StreamReader.readline()的底层等待套接字数据的过程没有检查取消请求,所以即使调用了cancel(),任务也会继续阻塞直到有数据或套接字关闭。
内容的提问来源于stack exchange,提问作者Dragoner
相关产品推荐
相关产品推荐

