如何将Python脚本改造为异步运行并支持传参创建新线程
异步改造核心原理与实现示例
关键原理拆解
- 持续运行机制:依托
asyncio事件循环维持脚本运行,只要事件循环不退出,脚本就会一直处于活跃状态。FastAPI底层基于asyncio,启动后由uvicorn自动维护事件循环;纯异步脚本则通过无限等待任务(如await asyncio.sleep(...))保持循环不终止。 - 异步+多线程配合:对于阻塞/CPU密集型的测试任务,用
asyncio.to_thread()或ThreadPoolExecutor将任务放到线程池执行,避免阻塞异步事件循环。异步事件循环负责调度IO密集型任务,线程池处理阻塞操作,兼顾并发效率与任务并行。 - 传参创建任务:提供异步接口(如FastAPI路由)或交互入口,接收参数后将任务提交到线程池/事件循环,实现动态创建线程任务。
结合FastAPI的实现示例
这个示例通过FastAPI提供HTTP接口,接收参数后创建线程任务,脚本持续运行等待请求:
import asyncio from fastapi import FastAPI import time from concurrent.futures import ThreadPoolExecutor app = FastAPI() # 配置线程池最大并发数 thread_pool = ThreadPoolExecutor(max_workers=5) # 模拟你的测试任务(同步阻塞函数) def test_task(task_param: str): print(f"线程任务启动 | 参数:{task_param}") time.sleep(3) # 模拟耗时操作 print(f"线程任务完成 | 参数:{task_param}") return f"任务结果:{task_param}" # 异步接口:接收参数并触发线程任务 @app.post("/trigger-task") async def trigger_task(task_param: str): # 用asyncio.to_thread把同步任务丢到线程池,异步等待结果 task_result = await asyncio.to_thread(test_task, task_param) return {"status": "success", "result": task_result} # 保持脚本持续运行(如果用uvicorn启动,此函数可省略,uvicorn会维护事件循环) async def keep_running(): while True: await asyncio.sleep(3600) if __name__ == "__main__": asyncio.run(keep_running())
运行说明
启动脚本后,用uvicorn your_script_name:app --reload启动FastAPI服务,然后通过POST请求http://localhost:8000/trigger-task(携带task_param参数)即可创建新线程任务。
纯异步脚本实现示例
如果不需要Web接口,用控制台交互传参创建线程任务:
import asyncio import time from concurrent.futures import ThreadPoolExecutor thread_pool = ThreadPoolExecutor(max_workers=5) def test_task(task_param: str): print(f"线程任务启动 | 参数:{task_param}") time.sleep(3) print(f"线程任务完成 | 参数:{task_param}") return f"任务结果:{task_param}" async def create_thread_task(task_param: str): # 提交同步任务到线程池,异步等待结果 result = await asyncio.to_thread(test_task, task_param) print(result) async def main(): # 初始化一个任务 await create_thread_task("初始测试参数") # 持续运行,等待用户输入参数创建新任务 while True: # 把阻塞的input操作放到线程,避免卡主事件循环 user_input = await asyncio.to_thread(input, "输入任务参数(输入q退出):") if user_input.lower() == "q": break # 用create_task创建后台任务,不阻塞主循环 asyncio.create_task(create_thread_task(user_input)) # 关闭线程池 thread_pool.shutdown() if __name__ == "__main__": asyncio.run(main())
内容的提问来源于stack exchange,提问作者Flouu
相关产品推荐
相关产品推荐

