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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 21:13:28