如何在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
相关产品推荐
相关产品推荐

