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

如何将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亲和性设置代码,可取消注释强制让两个进程运行在不同核心

关键说明:

  1. Pipe传递:通过multiprocessing.Process的args参数将Pipe的一端传递给子进程,这是多进程间传递Pipe的标准方式
  2. 异步监听Pipe:使用asyncio.add_reader将Pipe的文件描述符注册到事件循环,实现异步监听,避免同步poll阻塞TCP服务的正常处理
  3. 多核心运行:multiprocessing默认会调度子进程到不同核心运行,也可以通过cpu_affinity强制指定核心

内容的提问来源于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 17:05:54