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

并行进程运行时双向通信实现方法咨询

解决Python多进程双向并行通信问题

咱们先拆解下你现有代码里的核心问题,再一步步改成符合需求的双向并行通信逻辑:

现有代码的主要问题

  • 进程重复启停:你在主循环里每次都新建process_1和process_2,这会导致进程反复创建销毁,根本没法实现持续的双向通信,反而变成了串行执行的多次进程启动。
  • 共享内存操作错误:multiprocessing.Array是固定大小的共享内存数组,不能直接用append方法,运行时会直接报错。如果需要动态记录日志,应该用multiprocessing.Manager().list()。
  • 通信逻辑混乱:两个进程的启动时机和信号处理没对齐,比如process2一开始判断q.empty()就赋值1,但如果队列空的时候process1还在等待数据,很容易出现死锁或者逻辑混乱。
  • 队列关闭时机错误:主进程在启动后立刻q.close(),这会直接切断进程间的通信通道,后续数据根本传不起来。

修正后的实现方案

我们要让两个进程持续运行,互相收发数据,直到主进程发送终止信号;用动态共享列表记录数据;明确初始触发信号和终止逻辑。

完整修正代码

from multiprocessing import Process, Queue, Manager

# 终止信号定义
SENTINEL = -1

def process1(q, logbook):
    """接收数据、乘2后回传,直到收到终止信号"""
    while True:
        data = q.get()
        # 收到终止信号时,把信号放回队列让另一个进程也能接收,然后退出
        if data == SENTINEL:
            q.put(SENTINEL)
            break
        # 处理数据并记录日志
        processed_data = data * 2
        logbook.append(f"Process1: 接收{data} → 处理为{processed_data}")
        # 将处理结果回传给process2
        q.put(processed_data)

def process2(q, logbook):
    """接收数据、加1后回传,直到收到终止信号"""
    # 初始启动:先发送第一个数据触发整个交互循环
    initial_data = 1
    logbook.append(f"Process2: 初始启动 → 发送{initial_data}")
    q.put(initial_data)
    
    while True:
        data = q.get()
        if data == SENTINEL:
            break
        # 处理数据并记录日志
        processed_data = data + 1
        logbook.append(f"Process2: 接收{data} → 处理为{processed_data}")
        # 将处理结果回传给process1
        q.put(processed_data)

if __name__ == "__main__":
    # 进程间通信队列
    comm_queue = Queue()
    # 用Manager创建可安全在多进程中动态修改的共享列表
    with Manager() as manager:
        logbook = manager.list()
        # 只启动一次两个进程,让它们持续运行
        p1 = Process(target=process1, args=(comm_queue, logbook))
        p2 = Process(target=process2, args=(comm_queue, logbook))
        
        p1.start()
        p2.start()
        
        # 控制交互次数:让两个进程完成10次双向交互
        interaction_count = 0
        while interaction_count < 10:
            # 通过日志条目数量判断交互进度:每次完整交互对应两条日志(两个进程各一条)
            if len(logbook) >= (interaction_count + 1) * 2:
                interaction_count += 1
        
        # 交互完成后,发送终止信号
        comm_queue.put(SENTINEL)
        
        # 等待两个进程正常退出
        p1.join()
        p2.join()
        
        # 打印完整交互日志
        print("=== 双向交互日志 ===")
        for entry in logbook:
            print(entry)

关键逻辑解释

  1. 进程持续运行:两个进程只启动一次,通过while True循环持续监听队列,直到收到SENTINEL终止信号才退出,真正实现并行通信。
  2. 初始触发:process2启动后先发送初始数据1,主动触发process1的处理逻辑,让两个进程进入互相传递数据的循环。
  3. 共享日志:用Manager().list()代替Array,可以安全地在多进程中动态添加日志条目,不需要担心固定大小的限制。
  4. 终止逻辑:主进程在完成指定交互次数后发送终止信号,process1收到后会把信号放回队列,确保process2也能收到并退出,避免单个进程阻塞。
  5. 交互计数:通过日志列表的长度判断交互进度,比固定计数更准确,能适配进程调度的延迟情况。

运行效果说明

运行后你会看到两个进程交替处理数据:

  • Process2先发送初始数据1
  • Process1把1乘2得到2,回传给Process2
  • Process2把2加1得到3,回传给Process1
  • 以此类推,直到完成10次交互,最后两个进程收到终止信号退出

内容的提问来源于stack exchange,提问作者ceske

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:54:05