如何将Pipe引用传入Python多进程异步任务?
问题解决与代码修正
你的需求是让Task1(TCP服务)和Task2(数据处理)运行在不同核心,通过multiprocessing.Pipe通信,核心问题是Pipe的传递和异步监听的实现,以下是修正后的完整代码:
Task1.py
import asyncio from multiprocessing import Pipe class Task1: def __init__(self, pipe, writer): self.comm_pipe = pipe self.writer = writer # 注册pipe可读事件监听 asyncio.get_event_loop().add_reader(self.comm_pipe.fileno(), self._handle_pipe_data) async def run(self, reader): while True: data = await reader.read(100) if data: self.comm_pipe.send(data) else: break # 关闭连接时移除监听并释放资源 asyncio.get_event_loop().remove_reader(self.comm_pipe.fileno()) self.writer.close() await self.writer.wait_closed() def _handle_pipe_data(self): # 异步回调处理pipe返回的数据 if self.comm_pipe.poll(): data_from_pipe = self.comm_pipe.recv() self.writer.write(data_from_pipe) # 异步刷新写入缓冲区 asyncio.create_task(self.writer.drain()) async def task_main(pipe): # 封装客户端处理逻辑,传递pipe给Task1实例 async def handle_client(reader, writer): task1 = Task1(pipe, writer) await task1.run(reader) server = await asyncio.start_server( handle_client, '0.0.0.0', 4365) addr = server.sockets[0].getsockname() print(f'服务启动于 {addr}') async with server: await server.serve_forever() def main(pipe): asyncio.run(task_main(pipe))
修改点:
task_main新增pipe参数,通过内部的handle_client函数将pipe传递给每个新连接的Task1实例- 用
asyncio.get_event_loop().add_reader监听pipe的可读事件,避免同步poll阻塞异步循环,影响TCP连接处理 - 拆分
run方法职责,专门处理TCP数据读取,新增_handle_pipe_data回调处理Pipe返回的数据
Task2.py
import asyncio from multiprocessing import Pipe class Task2: def __init__(self, pipe): self.pipe = pipe # 注册pipe可读事件监听 asyncio.get_event_loop().add_reader(self.pipe.fileno(), self._handle_data) def _handle_data(self): if self.pipe.poll(): data = self.pipe.recv() # 此处可替换为后续的串口硬件交互逻辑 reversed_data = data[::-1] self.pipe.send(reversed_data) async def run(self): # 保持事件循环持续运行 await asyncio.Event().wait() def main(pipe): task2 = Task2(pipe) asyncio.run(task2.run())
修改点:
- 新增
main函数,接收pipe参数并创建Task2实例,对应App中的进程调用入口 - 同样用
add_reader异步监听Pipe的可读事件,避免同步阻塞 run方法通过等待永不触发的Event,保持进程存活以持续处理数据
App.py
import multiprocessing import Task1 import Task2 if __name__ == "__main__": parent_pipe, child_pipe = multiprocessing.Pipe() # 创建两个独立进程,分别绑定Task1和Task2的main函数 p1 = multiprocessing.Process(target=Task1.main, args=(parent_pipe,)) p2 = multiprocessing.Process(target=Task2.main, args=(child_pipe,)) # 可选:强制指定进程运行的核心(比如p1用核心0,p2用核心1) # p1.cpu_affinity([0]) # p2.cpu_affinity([1]) p1.start() p2.start() p1.join() p2.join()
修改点:
- 确保调用的是Task1和Task2的
main函数,作为两个进程的执行入口 - 注释了CPU亲和性设置代码,可取消注释强制让两个进程运行在不同核心
关键说明:
- Pipe传递:通过
multiprocessing.Process的args参数将Pipe的一端传递给子进程,这是多进程间传递Pipe的标准方式 - 异步监听Pipe:使用
asyncio.add_reader将Pipe的文件描述符注册到事件循环,实现异步监听,避免同步poll阻塞TCP服务的正常处理 - 多核心运行:
multiprocessing默认会调度子进程到不同核心运行,也可以通过cpu_affinity强制指定核心
内容的提问来源于stack exchange,提问作者AeroClassics
相关产品推荐
相关产品推荐

