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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 05:06:03