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

如何取消阻塞的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 20:05:20