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

如何将阻塞函数转为可在无限循环中运行的非阻塞函数?

问题描述

我有一个运行在同步无限循环中的程序,需要调用外部库的阻塞函数且无法修改该函数。我希望将该函数改造为非阻塞形式,让无限循环不被阻塞,持续打印变量n,仅在get_data()执行完成后打印data。

示例代码

基础示例

import time

SECONDS_TO_WAIT = 3

# 模拟外部库的阻塞函数
def get_data():
    time.sleep(SECONDS_TO_WAIT)
    return "Data received."

def main():
    n = 0

    while True:
        data = get_data()

        print(n)
        print(data)

        n += 1


if __name__ == "__main__":
    main()

贴近实际场景的示例

import time

class LibraryWrapper():
    SECONDS_TO_WAIT = 3

    def __init__(self):
        # 库初始化
        self.data = None

    # 模拟外部库的阻塞函数
    def get_data(self):
        time.sleep(self.SECONDS_TO_WAIT)
        return "Data received."


class Program():
    def __init__(self, library_wrapper):
        self.library_wrapper = library_wrapper

    def run(self, user_input):
        # 执行一些操作
        print(user_input)

        if library_wrapper.data:
            print(library_wrapper.data)


def main():
    library_wrapper = LibraryWrapper()

    program = Program(library_wrapper)

    user_input = 0

    while True:
        # 模拟非阻塞用户输入读取
        user_input += 1

        library_wrapper.data = library_wrapper.get_data()

        program.run(user_input)

if __name__ == "__main__":
    main()

尝试的asyncio代码(仍阻塞)

import asyncio

SECONDS_TO_WAIT = 3

async def get_data():
    await asyncio.sleep(3)
    return "Data received."

async def main():
    n = 0

    while True:
        task = asyncio.create_task(get_data())

        print(n)
        print(await task)

        n += 1

asyncio.run(main())

# 输出:
# 0
# 3秒后...
# Data received.
# 1
# 3秒后...
# Data received.
# 2
# 3秒后...
# Data received.
# 3

我认为asyncio库可能是最佳解决方案,已了解协程、async/await、任务和Future的基础概念,但尝试的代码仍会阻塞循环。也查阅了相关问题,但仍不清楚如何封装阻塞函数为非阻塞形式。

解决方案

核心思路是用asyncio.run_in_executor(或Python3.9+的asyncio.to_thread())将阻塞函数放到线程池中执行,避免阻塞asyncio事件循环。这样事件循环可以继续处理其他逻辑(比如持续打印n),同时等待阻塞函数的结果。

针对基础示例的改造

import asyncio
import time

SECONDS_TO_WAIT = 3

# 不可修改的外部阻塞函数
def get_data():
    time.sleep(SECONDS_TO_WAIT)
    return "Data received."

async def main():
    n = 0
    # 启动第一次数据获取任务,不立即等待结果
    data_task = asyncio.create_task(asyncio.to_thread(get_data))

    while True:
        print(n)
        n += 1
        
        # 检查任务是否完成
        if data_task.done():
            try:
                data = data_task.result()
                print(data)
                # 完成后启动下一次数据获取
                data_task = asyncio.create_task(asyncio.to_thread(get_data))
            except Exception as e:
                print(f"获取数据出错: {e}")
                # 出错也重新启动任务
                data_task = asyncio.create_task(asyncio.to_thread(get_data))
        
        # 让出CPU时间,避免无限循环占用100%资源
        await asyncio.sleep(0.1)

asyncio.run(main())

针对实际场景示例的改造

import asyncio
import time

class LibraryWrapper():
    SECONDS_TO_WAIT = 3

    def __init__(self):
        self.data = None
        self.data_task = None

    # 不可修改的外部阻塞函数
    def get_data(self):
        time.sleep(self.SECONDS_TO_WAIT)
        return "Data received."

    # 后台持续非阻塞获取数据
    async def start_continuous_fetch(self):
        while True:
            # 用线程执行阻塞函数
            self.data = await asyncio.to_thread(self.get_data)
            # 可以在这里添加获取间隔,或者直接连续获取

class Program():
    def __init__(self, library_wrapper):
        self.library_wrapper = library_wrapper

    def run(self, user_input):
        print(user_input)
        if self.library_wrapper.data:
            print(self.library_wrapper.data)
            # 打印后清空数据,避免重复输出
            self.library_wrapper.data = None

async def main():
    library_wrapper = LibraryWrapper()
    # 启动后台数据获取任务
    asyncio.create_task(library_wrapper.start_continuous_fetch())

    program = Program(library_wrapper)
    user_input = 0

    while True:
        user_input += 1
        program.run(user_input)
        # 让出CPU时间
        await asyncio.sleep(0.1)

asyncio.run(main())

关键说明

  • asyncio.to_thread()是Python3.9+的便捷方法,自动将阻塞函数放入线程池执行,替代手动创建ThreadPoolExecutor的繁琐操作。
  • 不要直接await阻塞任务,而是通过done()方法检查任务状态,这样事件循环可以继续执行其他逻辑。
  • 添加await asyncio.sleep(0.1)是为了避免无限循环占用100%CPU,实际场景中可根据需求调整间隔。
  • 任务完成后需重新启动下一次获取,保证数据的持续更新。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 18:15:14