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

Python3.5中永久运行的异步函数间如何实现值交换?

嘿,我懂你现在的困扰——找了好多资料还是没搞懂怎么让两个一直跑的异步函数互相传变量对吧?别担心,用asyncio.Queue就能完美解决这个问题,它就是专门为异步协程间的安全通信设计的。

示例:用asyncio.Queue实现永久运行异步函数间的变量传递

这里有一个完整的可运行示例,一个协程持续生成数据,另一个持续接收并处理:

import asyncio

async def data_generator(queue: asyncio.Queue):
    """永久运行的生成器,不断产生数据放入队列"""
    counter = 0
    while True:
        # 模拟生成数据(比如从传感器、外部API获取)
        data = f"Generated data {counter}"
        print(f"Putting data into queue: {data}")
        await queue.put(data)
        counter += 1
        await asyncio.sleep(1)  # 模拟生成数据的耗时操作

async def data_processor(queue: asyncio.Queue):
    """永久运行的处理器,不断从队列取数据并处理"""
    while True:
        data = await queue.get()
        print(f"Processing received data: {data}")
        # 这里可以添加你的自定义处理逻辑(计算、存储、转发等)
        await asyncio.sleep(0.5)  # 模拟数据处理的耗时
        queue.task_done()  # 标记队列任务完成(规范写法,非强制但推荐)

async def main():
    # 创建异步队列,可设置最大容量(这里设为5,防止数据积压)
    data_queue = asyncio.Queue(maxsize=5)
    
    # 将两个协程包装成异步任务,让它们在事件循环中并发运行
    generator_task = asyncio.create_task(data_generator(data_queue))
    processor_task = asyncio.create_task(data_processor(data_queue))
    
    # 等待两个任务完成(因为是永久循环,这里会一直阻塞,直到手动终止)
    await asyncio.gather(generator_task, processor_task)

if __name__ == "__main__":
    try:
        asyncio.run(main())
    except KeyboardInterrupt:
        print("\n程序已手动终止")

关键知识点拆解

  • asyncio.Queue:异步环境下的线程安全队列,put()和get()都是异步方法,会自动处理协程间的同步,完全不用担心竞态条件。
  • 永久运行的协程:用while True循环实现,配合await asyncio.sleep()避免阻塞整个事件循环。
  • 任务管理:通过asyncio.create_task()把协程包装成独立任务,让它们能在事件循环中并行执行。
  • 双向通信扩展:如果需要两个函数互相传递数据,只需要再创建一个队列,让两个协程同时读写两个队列即可,比如:
# 双向通信示例
async def func_a(queue_to_b: asyncio.Queue, queue_from_b: asyncio.Queue):
    counter = 0
    while True:
        # 发送数据给func_b
        await queue_to_b.put(f"Message from A: {counter}")
        # 接收func_b的响应
        response = await queue_from_b.get()
        print(f"A got response: {response}")
        counter += 1
        await asyncio.sleep(1)

async def func_b(queue_from_a: asyncio.Queue, queue_to_a: asyncio.Queue):
    while True:
        # 接收func_a的数据
        msg = await queue_from_a.get()
        print(f"B received: {msg}")
        # 发送响应给func_a
        await queue_to_a.put(f"Reply to {msg}")
        await asyncio.sleep(0.8)

async def main():
    queue_a_to_b = asyncio.Queue()
    queue_b_to_a = asyncio.Queue()
    
    task_a = asyncio.create_task(func_a(queue_a_to_b, queue_b_to_a))
    task_b = asyncio.create_task(func_b(queue_a_to_b, queue_b_to_a))
    
    await asyncio.gather(task_a, task_b)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 08:05:09