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
相关产品推荐
相关产品推荐

