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

Python多进程场景下实现进程间变量实时共享的技术咨询

问题解答

核心问题说明

首先纠正你对Value/Array的误解:这两类共享内存对象本身不需要等待进程结束就能读取,你当前拿不到实时数据的核心原因有两个:

  • 你的进程B代码逻辑是等待进程A完全结束后,才会调用do_somethig_with_values处理数据
  • 共享内存没有做读写同步,直接读取很容易拿到进程A写了一半的脏数据,反而会引发逻辑错误

实现方案选择

你这个单生产者(进程A生成数据)、单消费者(进程B处理数据)的场景,用multiprocessing内置的Queue或者Pipe都可以实现实时数据传递,不需要自己手动做同步逻辑,比共享内存方案更省心:

方案1:使用multiprocessing.Queue(更推荐,上手简单)

Queue是线程/进程安全的FIFO队列,进程A每生成一组新数据就写入队列,进程B循环从队列读取数据即可,读取操作默认阻塞,有新数据就会立刻返回,不会空耗CPU。
示例修改代码如下:

进程A修改后代码

def process_A(processid, textid, texttitle, queue):
    iters = 10
    p_var1, p_var2 = 0, 0
    i = 0
    while i < iters:
        # 注意:你原来的new_data_incoming()写在循环外,每次拿的都是同一份数据,实际使用要移到循环内
        h1 = new_data_incoming()
        data  = h1.text.split(" ")
        a_var1 = float(data[1].replace(',',''))
        a_var2 = float(data[1].replace(',',''))
        if a_var1 != p_var1:
            # 组装数据
            values = [0]*6
            values[0] = a_var1
            values[1] = a_var2
            values[2] = float(data[3].replace(',',''))
            values[3] = float(data[4].replace(',',''))
            values[4] = float(data[5][1:]) if data[5][0] == "+" else -float(data[5][1:])
            values[5] = float(data[6][1:-1]) if data[6][0] == "+" else -float(data[6][1:-1])
            # 写入队列
            queue.put(values)
            p_var1 = a_var1
            p_var2 = a_var2
            i = i + 1
    # 所有数据生成完成,写入结束标记
    queue.put(None)

进程B修改后代码

def process_B():
    from multiprocessing import Queue, Process
    ids = getids()
    key2 = "stocastic_process"
    id2  = ids[key2]
    # 初始化队列,maxsize可限制队列最大长度,避免内存溢出
    data_queue = Queue(maxsize=10) 
    process_A_launched = Process(target=process_A, args=(2, id2, key2, data_queue))
    process_A_launched.start()
    print("Process with id %i has started: stocastic %s" % (2, key2))
    
    # 循环读取数据
    while True:
        values = data_queue.get()
        # 收到结束标记,退出循环
        if values is None:
            break
        # 直接处理新收到的实时数据
        do_somethig_with_values(values)
    
    # 等待进程A退出
    process_A_launched.join()
    print("process %i has finished!" % (2))
    print("All process have finished!")

方案2:使用multiprocessing.Pipe(更轻量,适合一对一通信)

如果追求更高的通信性能,可以用Pipe实现,Pipe返回一对连接对象,分别交给两个进程使用,初始化时设置duplex=False即可实现单向通信,性能比Queue更高:

# 初始化Pipe,duplex=False表示只能从conn_recv读,conn_send写
conn_recv, conn_send = multiprocessing.Pipe(duplex=False)

进程A生成数据后调用conn_send.send(values)发送,进程B调用values = conn_recv.recv()接收即可,其余逻辑和Queue方案一致。

原有共享内存方案的修改方式

如果一定要用Array实现,需要额外加multiprocessing.Lock做读写同步:进程A写完数据后释放锁,进程B拿到锁后再读取数据。但这种方案需要自己维护状态标记,很容易出现死锁、脏读问题,不推荐使用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 20:15:02