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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 11:55:40