Python中能否异步使用Stream Socket与multiprocessing Pipe?附场景需求
实现方案与Task1示例代码
关键要点
- 将
multiprocessing.Pipe转换为异步IO对象,适配asyncio事件循环 - 使用
asyncio.StreamReader/StreamWriter实现支持receiveuntil()的异步Socket通信 - 在Task1中通过
asyncio.gather同时运行Socket通信和Pipe通信两个异步任务
Task1实现示例
1. 异步Pipe封装类
先封装一个异步Pipe类,让multiprocessing的Pipe可以被asyncio await:
import asyncio import os from multiprocessing import Pipe class AsyncPipe: def __init__(self, pipe_end): self.pipe_end = pipe_end self.read_queue = asyncio.Queue() self.loop = asyncio.get_event_loop() # 注册文件描述符读事件监听 self.loop.add_reader(pipe_end.fileno(), self._read_from_pipe) def _read_from_pipe(self): # 非阻塞读取Pipe数据,放入队列 try: data = os.read(self.pipe_end.fileno(), 1024) if data: self.read_queue.put_nowait(data) except BlockingIOError: pass async def read(self): # 异步读取Pipe数据 return await self.read_queue.get() async def write(self, data): # 异步写入Pipe数据 os.write(self.pipe_end.fileno(), data)
2. Task1核心逻辑
Task1同时处理Socket的receiveuntil()和Pipe的异步通信:
async def task1_core(pipe1_end): # 初始化异步Pipe async_pipe = AsyncPipe(pipe1_end) # 连接到目标Socket服务端(示例地址,根据实际修改) reader, writer = await asyncio.open_connection('127.0.0.1', 8888) print("Task1: Socket连接成功") async def handle_socket(): while True: # 使用receiveuntil()读取Socket数据,直到指定分隔符(示例用b'\r\n') data = await reader.receiveuntil(b'\r\n') print(f"Task1: 从Socket收到数据: {data.decode().strip()}") # 将Socket数据通过Pipe发送给Task2 await async_pipe.write(data) async def handle_pipe(): while True: # 从Pipe读取Task2发来的数据 data = await async_pipe.read() print(f"Task1: 从Pipe收到数据: {data.decode().strip()}") # 将Pipe数据发送到Socket writer.write(data + b'\r\n') await writer.drain() # 同时运行两个异步任务 await asyncio.gather(handle_socket(), handle_pipe()) def task1_process(pipe1_end): # 启动asyncio事件循环 asyncio.run(task1_core(pipe1_end))
3. MyApp启动逻辑
在myapp中创建Pipe并启动所有进程:
from multiprocessing import Process def task2_process(pipe1_end, pipe2_end): # 可参照Task1的异步Pipe封装实现Task2逻辑 pass def task3_process(pipe2_end): # 实现Task3的USB读写+Pipe通信逻辑 pass if __name__ == "__main__": # 创建两个双向Pipe pipe1_parent, pipe1_child = Pipe() pipe2_parent, pipe2_child = Pipe() # 启动各进程 p1 = Process(target=task1_process, args=(pipe1_child,)) p2 = Process(target=task2_process, args=(pipe1_parent, pipe2_child)) p3 = Process(target=task3_process, args=(pipe2_parent,)) p1.start() p2.start() p3.start() p1.join() p2.join() p3.join()
注意事项
receiveuntil()的分隔符需根据实际协议调整,示例中使用b'\r\n'作为换行分隔符- 异步Pipe封装中使用
os.read/os.write而非Pipe自带的recv/send,是为了配合事件循环的非阻塞监听 - CPU密集型操作需放入单独的线程或进程处理,避免阻塞asyncio事件循环
内容的提问来源于stack exchange,提问作者AeroClassics
相关产品推荐
相关产品推荐

