Python3.11同步异步函数并行实现需求及代码优化问询
同步采集+无阻塞异步发送的实现方案
现有代码的问题
你当前的代码存在几个关键问题,无法满足需求:
- 异步函数中使用
time.sleep()是同步阻塞操作,会卡住整个事件循环,导致后续的同步采集任务被迫暂停,完全达不到“不阻塞同步函数”的要求 task1被定义为async,但内部是纯同步操作,没有发挥异步IO的优势,反而拖慢了整个循环- 未处理异步任务的异常,长期运行后单个发送任务崩溃可能影响整个程序的稳定性
改进后的Asyncio方案(推荐,资源利用率更高)
将同步采集操作放到线程池执行,避免阻塞事件循环;用asyncio.sleep替代time.sleep;同时做内存和异常控制:
import time import asyncio from concurrent.futures import ThreadPoolExecutor # 纯同步的本地数据采集函数 def sync_collect_data(): print("开始从本地采集数据...") time.sleep(3) # 模拟读文件、查本地DB这类同步IO耗时 print("数据采集完成") return "采集到的本地数据" # 异步发送数据到服务器 async def async_send_to_server(data): try: print(f"启动发送任务: {data}") # 模拟5-30秒的发送耗时 await asyncio.sleep(asyncio.get_event_loop().random() * 25 + 5) print(f"发送完成: {data}") except Exception as e: print(f"发送失败: {str(e)}") async def main(): # 创建线程池,限制最大并发数,防止内存耗尽 executor = ThreadPoolExecutor(max_workers=2) loop = asyncio.get_running_loop() while True: # 在线程池执行同步采集,不阻塞事件循环 data = await loop.run_in_executor(executor, sync_collect_data) # 提交异步发送任务,后台执行,不等待完成 asyncio.create_task(async_send_to_server(data)) # 如果需要给采集任务加间隔,取消下面的注释 # await asyncio.sleep(1) if __name__ == "__main__": asyncio.run(main())
方案优势:
- 同步采集在独立线程运行,完全不影响事件循环,保证采集任务持续执行
- 异步发送任务后台运行,不会阻塞下一次采集操作
- 线程池限制并发数,避免大量线程占用过多内存
- 异常捕获确保单个发送任务失败不会导致整个程序崩溃
非Asyncio方案(Threading+Queue,更直观)
如果不习惯异步编程,用线程+队列也能实现需求,逻辑更简单:
import time import threading import queue import random # 数据队列,限制最大长度,避免数据堆积耗尽内存 data_queue = queue.Queue(maxsize=10) # 同步采集线程 def collect_data_thread(): while True: print("开始从本地采集数据...") time.sleep(3) # 模拟同步IO耗时 data = "采集到的本地数据" print("采集完成,放入等待队列") # 队列满时自动阻塞,防止采集过快 data_queue.put(data) # 可选:添加采集间隔 # time.sleep(1) # 发送线程(可启动多个) def send_to_server_thread(): while True: data = data_queue.get() # 队列空时阻塞,等待新数据 try: print(f"启动发送任务: {data}") # 模拟5-30秒发送耗时 time.sleep(random.uniform(5, 30)) print(f"发送完成: {data}") except Exception as e: print(f"发送失败: {str(e)}") finally: data_queue.task_done() if __name__ == "__main__": # 启动采集线程(守护线程,随主线程退出) collect_thread = threading.Thread(target=collect_data_thread, daemon=True) collect_thread.start() # 启动2个发送线程,可根据需求调整数量 for _ in range(2): send_thread = threading.Thread(target=send_to_server_thread, daemon=True) send_thread.start() # 主线程等待采集线程运行 collect_thread.join()
方案优势:
- 队列实现采集和发送解耦,采集线程持续运行,发送线程按需取数据
- 队列长度限制避免内存堆积
- 多发送线程提高效率,完全不阻塞采集操作
- 逻辑简单,容易调试和维护
两种方案都能满足24小时不间断运行、不阻塞同步任务、控制内存的需求,你可以根据自己的技术栈选择。
内容的提问来源于stack exchange,提问作者g-pane
相关产品推荐
相关产品推荐

