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

如何在aiohttp应用中正确终止等待设备数据的后台任务?

解决aiohttp应用关闭时,监控字符设备的异步任务无法退出的问题

核心问题在于:读取/dev/scaner0这类字符设备时,无数据状态下await dev.read(1)会陷入底层同步IO阻塞,此时asyncio的任务取消信号(CancelledError)无法被触发——因为任务不在事件循环的调度队列里,直到有数据到来才会回到事件循环,才能处理取消请求。

以下是三种可行的解决方案:

方案1:给读取操作加超时(改动最小)

通过asyncio.wait_for给每次read操作设置超时,让任务定期回到事件循环,即使没有数据,也能及时响应取消信号。

修改后的watchscaner函数:

async def watchscaner(app):
    try:
        async with aiofiles.open("/dev/scaner0", "r") as dev:
            buffer = ""
            while True:
                try:
                    # 每隔0.5秒触发一次超时,回到事件循环检查取消状态
                    b = await asyncio.wait_for(dev.read(1), timeout=0.5)
                except asyncio.TimeoutError:
                    # 超时后主动检查任务是否被取消
                    if asyncio.current_task().cancelled():
                        raise asyncio.CancelledError
                    continue
                
                # 原有的数据处理逻辑不变
                if b == "\n":
                    buffer = ""
                elif b == "\x00":
                    if len(buffer) > 0:
                        logging.info(f"Scanned {buffer=}")
                        buffer = ""
                else:
                    buffer += b
    except asyncio.CancelledError:
        logging.warning('Cancelled error')
    finally:
        logging.info('Finally came here')

方案2:同步线程+asyncio队列(适合复杂设备交互)

将阻塞的设备读取操作放到单独的同步线程中,用asyncio.Queue传递数据到异步任务。取消任务时,通过事件通知线程停止读取。

import threading
import asyncio

async def watchscaner(app):
    data_queue = asyncio.Queue()
    stop_event = threading.Event()

    # 同步读取设备的线程函数
    def read_device_loop():
        try:
            with open("/dev/scaner0", "r") as dev:
                while not stop_event.is_set():
                    b = dev.read(1)
                    if b:
                        data_queue.put_nowait(b)
        except Exception as e:
            logging.error(f"Device read failed: {str(e)}")
        finally:
            data_queue.put_nowait(None)  # 发送结束标记

    # 启动后台线程
    read_thread = threading.Thread(target=read_device_loop, daemon=True)
    read_thread.start()

    try:
        buffer = ""
        while True:
            b = await data_queue.get()
            if b is None:
                break  # 线程结束,退出循环
            
            # 原数据处理逻辑
            if b == "\n":
                buffer = ""
            elif b == "\x00":
                if len(buffer) > 0:
                    logging.info(f"Scanned {buffer=}")
                    buffer = ""
                else:
                    buffer += b
    except asyncio.CancelledError:
        logging.warning('Cancelled error')
        stop_event.set()
        read_thread.join(timeout=1)  # 等待线程退出
    finally:
        logging.info('Finally came here')

方案3:非阻塞IO+事件监听(性能最优)

用os.open以非阻塞模式打开设备,通过asyncio的add_reader监听设备的可读事件,避免主动阻塞,让事件循环可以自由调度任务,包括处理取消信号。

import os
import asyncio

async def watchscaner(app):
    # 以非阻塞模式打开设备
    fd = os.open("/dev/scaner0", os.O_RDONLY | os.O_NONBLOCK)
    loop = asyncio.get_running_loop()
    buffer = ""
    stop_flag = False

    # 可读事件回调函数
    def handle_read():
        nonlocal buffer, stop_flag
        try:
            # 非阻塞读取,无数据时会抛出BlockingIOError
            b = os.read(fd, 1).decode()
            if not b:
                stop_flag = True
                return
            
            # 原数据处理逻辑
            if b == "\n":
                buffer = ""
            elif b == "\x00":
                if len(buffer) > 0:
                    logging.info(f"Scanned {buffer=}")
                    buffer = ""
            else:
                buffer += b
        except BlockingIOError:
            # 无数据时忽略,等待下一次可读事件
            pass
        except Exception as e:
            logging.error(f"Read error: {str(e)}")
            stop_flag = True

    try:
        loop.add_reader(fd, handle_read)
        while not stop_flag:
            await asyncio.sleep(0.1)  # 让出事件循环
    except asyncio.CancelledError:
        logging.warning('Cancelled error')
    finally:
        loop.remove_reader(fd)
        os.close(fd)
        logging.info('Finally came here')

方案选择建议

  • 优先选方案1:代码改动最小,适配现有aiofiles的使用习惯,适合大多数场景。
  • 追求性能选方案3:无线程开销,完全基于事件驱动,适合高并发或资源受限的环境。
  • 复杂设备交互选方案2:如果设备读取需要额外的同步逻辑(比如配置、异常重试),线程模式更易实现。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 17:44:52