如何从asyncio事件循环启动新Process且不阻塞原有协程逻辑运行
问题根因
multiprocessing.Process.join()是同步阻塞方法,直接在async协程中调用会卡住整个事件循环,导致shoot任务无法被调度执行- 删除
p.join()后,Executor.run()方法在启动进程和shoot任务后会直接返回,主协程退出会导致整个程序终止 - 现有代码中
multiprocessing.Queue.get()也是同步阻塞调用,直接在shoot协程中调用同样会阻塞事件循环
解决方法
核心思路是把所有同步阻塞操作都放到异步执行器中运行,避免阻塞事件循环,同时让主协程等待所有任务完成再退出。
完整修改后代码
import asyncio from multiprocessing import Process, Queue from cryptofeed import FeedHandler from cryptofeed.defines import L2_BOOK from cryptofeed.exchanges.ftx import FTX class Pricefeed(Process): def __init__(self, queue: Queue): super().__init__() self.daemon = True # 设为守护进程,主程序退出时自动销毁 self.coin_symbol = 'SOL-USD' self.fut_symbol = 'SOL-USD-PERP' self.queue = queue async def _book_update(self, feed, symbol, book, timestamp, receipt_timestamp): self.queue.put(book) def run(self): fh = FeedHandler() fh.add_feed(FTX(symbols=[self.fut_symbol, self.coin_symbol], channels=[L2_BOOK], callbacks={L2_BOOK: self._book_update})) fh.run() class Executor: def __init__(self): self.q = Queue() async def shoot(self): print('in shoot') loop = asyncio.get_running_loop() # 如果需要持续运行就把下面的for循环改成while True for i in range(5): # 把同步的get操作放到线程池中运行,避免阻塞事件循环 msg = await loop.run_in_executor(None, self.q.get) print(msg) await asyncio.sleep(1) # do some stuff async def run(self): shoot_task = asyncio.create_task(self.shoot()) p = Pricefeed(self.q) p.start() # 把同步的join操作放到线程池中异步执行 join_task = asyncio.create_task(asyncio.to_thread(p.join)) # 等待shoot任务和进程任务任意一个完成就退出,可根据需求调整等待逻辑 await asyncio.gather(shoot_task, join_task) async def main(): g = Executor() await g.run() if __name__ == '__main__': asyncio.run(main())
关键修改说明
- 给
Pricefeed进程设置daemon = True,主程序退出时会自动终止该进程,避免出现孤儿进程 - 用
asyncio.to_thread把阻塞的p.join()包装成异步任务,不会阻塞事件循环 - 把
shoot方法中的同步q.get()用loop.run_in_executor包装成异步调用,避免阻塞事件循环 - 用
asyncio.gather同时等待shoot任务和进程join任务完成,保证主协程不会提前退出 - 如果需要
shoot方法持续运行,把for i in range(5)改成while True即可
内容的提问来源于stack exchange,提问作者mchangun
相关产品推荐
相关产品推荐

