Python多进程间实现低延迟事件通知的最优通信方案是什么?
最低延迟进程间事件通知通信方案寻求
我正在寻找两个进程间通知事件发生的最低延迟(以latency为衡量标准)通信方式。
具体场景如下:共享内存中存储了numpy数组,一个生产者进程负责向数组写入更新,另一个消费者进程负责读取数组内容。
由于需要突破GIL限制,必须使用多进程方案。生产者属于CPU/IO密集型进程,负责监听数据流并完成相关数据处理。
消费者负载极轻,大部分时间处于空闲状态,需要在生产者更新数组后以最快速度将其唤醒。
额外要求:以最小延迟触发消费者的优先级高于完整传递所有消息。例如生产者连续无延迟发送三条消息,消费者仅收到第一条、丢失后两条的情况是可以接受的。
我已经尝试了multiprocessing提供的Pipe、Queue、Event三种原语,测试显示三者的延迟表现接近,其中Pipe的稳定性最优。
multiprocessing.Pipe
import multiprocessing as mp import numpy as np import random import time ITER_COUNT = 1000 def get_mcs_diff(ts): return round((time.time() - ts) * 1e6, 0) def main(v, input_pipe): for _ in range(ITER_COUNT): v.value = time.time() input_pipe.send(None) time.sleep((0.1 + random.random()) / 100) if __name__ == "__main__": v = mp.Value('d', time.time()) (ip, op) = mp.Pipe() p = mp.Process(target=main, args=(v, ip,)) measurements = [] p.start() i = 0 while i < ITER_COUNT: op.recv() measurements.append(get_mcs_diff(v.value)) i += 1 print(np.percentile(measurements, [50, 90, 95, 99], axis=0)) p.join()
输出如下(微秒级50、90、95、99分位值):
[138. 206.1 238. 383.21]
multiprocessing.Queue
import multiprocessing as mp import numpy as np import random import time ITER_COUNT = 1000 def get_mcs_diff(ts): return round((time.time() - ts) * 1e6, 0) def main(v, q): for _ in range(ITER_COUNT): v.value = time.time() q.put(None) time.sleep((0.1 + random.random()) / 100) if __name__ == "__main__": v = mp.Value('d', time.time()) q = mp.Queue() p = mp.Process(target=main, args=(v, q,)) measurments = [] p.start() i = 0 while i < ITER_COUNT: q.get() measurments.append(get_mcs_diff(v.value)) i += 1 print(measurments) print(np.percentile(measurments, [50, 90, 95, 99], axis=0)) p.join()
输出如下(微秒级50、90、95、99分位值):
[187. 266. 299.05 444.06]
multiprocessing.Event
import multiprocessing as mp import numpy as np import random import time ITER_COUNT = 1000 def get_mcs_diff(ts): return round((time.time() - ts) * 1e6, 0) def main(v, e): for _ in range(ITER_COUNT): v.value = time.time() e.set() time.sleep((0.1 + random.random()) / 100) if __name__ == "__main__": v = mp.Value('d', time.time()) e = mp.Event() p = mp.Process(target=main, args=(v, e,)) measurments = [] p.start() i = 0 while i < ITER_COUNT: e.wait() measurments.append(get_mcs_diff(v.value)) i += 1 e.clear() print(np.percentile(measurments, [50, 90, 95, 99], axis=0)) p.join()
输出如下(微秒级50、90、95、99分位值):
[142. 222.1 256.05 1754.77]
while True忙等待方案
import multiprocessing as mp import numpy as np import random import time ITER_COUNT = 1000 def get_mcs_diff(ts): return round((time.time() - ts) * 1e6, 0) def main(v): time.sleep(1) for _ in range(ITER_COUNT): v.value = time.time() # print(v.value) time.sleep((0.1 + random.random()) / 100) if __name__ == "__main__": v = mp.Value('d', time.time()) p = mp.Process(target=main, args=(v,)) measurments = [] p.start() i = 0 v_prev = 0 while i < ITER_COUNT: # print(v_prev - v.value) if v_prev < v.value: measurments.append(get_mcs_diff(v.value)) v_prev = float(v.value) i += 1 print(np.percentile(measurments, [50, 90, 95, 99], axis=0)) p.join()
输出如下(微秒级50、90、95、99分位值):
[ 33. 65. 81. 128.05]
截至目前忙等待方案是性能最好的,但由于其占用CPU资源的明显缺陷,不希望采用该方案。
内容的提问来源于stack exchange,提问作者Alex T
相关产品推荐
相关产品推荐

