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

Python中非阻塞流读取:Asyncio还是Concurrent Futures?

解决方案:异步/并发实现带定时任务的流读取循环

核心澄清:Asyncio的await不是线程阻塞

你之前的误解在于把协程的暂停和线程的阻塞混为一谈。Asyncio的await只会暂停当前协程,把控制权交还给事件循环,让其他就绪的协程有机会执行,并不会阻塞整个线程——这是Asyncio实现并发的核心逻辑。


方案一:用Asyncio实现(推荐IO密集场景)

通过asyncio.wait同时监听「流读取完成」和「定时时间到」两个事件,完美匹配你的伪代码逻辑。如果原有readStream()是阻塞函数,无需修改,用asyncio.to_thread就能把它适配为异步接口。

示例代码

import asyncio
import time

# 原阻塞式流读取函数(无需修改)
def readStream():
    # 模拟阻塞等待流数据(实际场景可能是socket/管道读取)
    ...

# 处理流数据的函数
def processData(data):
    ...

# 定时执行的 housekeeping 函数
def houseKeepingTasks():
    ...

# 将阻塞读取转为异步接口
async def async_read_stream():
    return await asyncio.to_thread(readStream)

async def main(min_step=1.0):
    last_loop_time = time.time()
    while True:
        now = time.time()
        # 计算距离下一次housekeeping还需等待的时间
        sleep_duration = max(0.0, min_step - (now - last_loop_time))
        
        # 同时等待两个任务:读取流 或 定时睡眠结束
        done, pending = await asyncio.wait(
            [async_read_stream(), asyncio.sleep(sleep_duration)],
            return_when=asyncio.FIRST_COMPLETED
        )
        
        # 处理流数据(如果读取任务先完成)
        for task in done:
            if task.get_name() == async_read_stream.__name__:
                data = task.result()
                processData(data)
        
        # 无论有没有流数据,都执行housekeeping
        houseKeepingTasks()
        
        # 更新循环时间戳
        last_loop_time = time.time()
        
        # 取消未完成的任务(比如流先就绪时,睡眠任务还在等待)
        for task in pending:
            task.cancel()

asyncio.run(main())

逻辑说明

  • asyncio.wait会监控两个任务,哪个先触发就继续执行,不会卡住事件循环
  • 若流先有数据:读取任务完成,睡眠任务被取消,处理数据后执行housekeeping
  • 若流无数据:睡眠任务到时间触发,直接执行housekeeping,进入下一轮循环

方案二:用Concurrent Futures实现(线程/进程模型)

如果更习惯线程模型,也可以用线程池实现,优化你原本的轮询思路,避免time.sleep(0.01)的低效轮询,改用wait等待事件触发。

示例代码

from concurrent.futures import ThreadPoolExecutor, wait, FIRST_COMPLETED
import time
import threading

# 原阻塞式流读取函数
def readStream():
    ...

def processData(data):
    ...

def houseKeepingTasks():
    ...

def main(min_step=1.0):
    last_loop_time = time.time()
    executor = ThreadPoolExecutor(max_workers=2)
    
    while True:
        now = time.time()
        sleep_duration = max(0.0, min_step - (now - last_loop_time))
        
        # 提交流读取任务
        read_future = executor.submit(readStream)
        # 启动定时线程
        sleep_event = threading.Event()
        threading.Thread(target=lambda: (time.sleep(sleep_duration), sleep_event.set())).start()
        
        # 等待读取完成或定时触发(二选一)
        wait([read_future], timeout=sleep_duration)
        
        # 处理流数据(如果读取完成)
        if read_future.done():
            data = read_future.result()
            processData(data)
        
        # 执行housekeeping
        houseKeepingTasks()
        
        last_loop_time = time.time()
        
        # 清理未完成的读取任务(可选,根据业务场景决定是否保留)
        if not read_future.done():
            read_future.cancel()

方案选择建议

  • IO密集场景:优先选Asyncio,协程切换开销远低于线程,资源占用更少
  • CPU密集场景:选Concurrent Futures的进程池(或线程池),Asyncio受GIL限制无法利用多核

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 07:20:29