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

多子进程调用父进程函数 解决数据库过载与文件读取延迟问题

替代文件共享的进程间数据同步方案

针对你遇到的文件IO延迟问题,以下几种进程间通信(IPC)方案可以让子进程直接从父进程获取最新数据,完全绕开磁盘IO的开销,同时避免直接访问数据库:

1. 双向IPC管道请求-响应模式

父进程在内存中维护一份最新数据副本,为每个子进程创建专属的双向管道。子进程需要数据时向父进程发送请求,父进程直接返回内存中的最新数据,全程在内存中完成数据传递,延迟远低于文件读写。

伪代码示例(Python)

父进程侧:

import multiprocessing
import threading

def handle_child_conn(conn, latest_data):
    # 持续处理对应子进程的请求
    while True:
        if conn.poll():
            req = conn.recv()
            if req == "FETCH_LATEST":
                conn.send(latest_data)

def main():
    latest_data = {}
    # 启动MongoDB监听线程,实时更新latest_data
    threading.Thread(target=listen_mongo_updates, args=(latest_data,)).start()

    child_connections = []
    # 为每个子进程创建管道并启动处理线程
    for _ in range(16):
        parent_conn, child_conn = multiprocessing.Pipe()
        child_connections.append(child_conn)
        threading.Thread(target=handle_child_conn, args=(parent_conn, latest_data)).start()
    
    # 启动16个子进程,传入各自的管道连接
    for conn in child_connections:
        multiprocessing.Process(target=child_worker, args=(conn,)).start()

子进程侧:

def child_worker(conn):
    while True:
        # 按需向父进程请求最新数据
        conn.send("FETCH_LATEST")
        latest_data = conn.recv()
        # 执行业务操作
        process_updated_data(latest_data)

2. 共享内存+信号量同步

利用操作系统提供的共享内存机制,让父子进程直接访问同一块内存区域,同时用信号量保证数据更新时的线程安全,避免子进程读取到半更新的脏数据。这种方式下子进程可以直接读取内存数据,不需要额外的请求交互,延迟最低。

伪代码示例(Python)

父进程侧:

from multiprocessing import shared_memory, Semaphore, Process
import json

def listen_mongo_updates(shm, sem):
    # 监听MongoDB变更,更新共享内存
    while True:
        new_data = fetch_mongo_updates()
        data_bytes = json.dumps(new_data).encode()
        # 加锁避免并发写入/读取冲突
        sem.acquire()
        shm.buf[:len(data_bytes)] = data_bytes
        # 填充空字节覆盖旧数据
        shm.buf[len(data_bytes):] = b'\x00'
        sem.release()

def main():
    # 创建共享内存(按需调整大小)
    shm = shared_memory.SharedMemory(create=True, size=4096)
    # 创建二元信号量保证互斥访问
    sem = Semaphore(1)

    # 启动MongoDB监听线程
    import threading
    threading.Thread(target=listen_mongo_updates, args=(shm, sem)).start()

    # 启动16个子进程
    for _ in range(16):
        Process(target=child_worker, args=(shm.name, sem)).start()

子进程侧:

from multiprocessing import shared_memory, Semaphore
import json

def child_worker(shm_name, sem):
    shm = shared_memory.SharedMemory(name=shm_name)
    while True:
        # 加锁读取共享内存
        sem.acquire()
        data_bytes = shm.buf.tobytes().strip(b'\x00')
        sem.release()
        if data_bytes:
            latest_data = json.loads(data_bytes)
            process_updated_data(latest_data)

3. 消息队列广播模式

父进程每次获取到MongoDB更新后,将最新数据写入一个固定长度的消息队列(仅保留最新一条数据),子进程从队列中读取最新数据。这种方式适合需要主动推送更新的场景,子进程无需主动请求,等待队列消息即可。

伪代码示例(Python)

父进程侧:

from multiprocessing import Queue, Process
import threading

def listen_mongo_updates(queue):
    while True:
        new_data = fetch_mongo_updates()
        # 清空队列并写入最新数据,确保子进程拿到的总是最新的
        while not queue.empty():
            queue.get()
        queue.put(new_data)

def main():
    # 创建队列,设置最大长度为1
    update_queue = Queue(maxsize=1)
    # 启动MongoDB监听线程
    threading.Thread(target=listen_mongo_updates, args=(update_queue,)).start()
    # 启动子进程
    for _ in range(16):
        Process(target=child_worker, args=(update_queue,)).start()

子进程侧:

def child_worker(queue):
    while True:
        # 阻塞等待最新数据
        latest_data = queue.get()
        process_updated_data(latest_data)
        # 标记任务完成,让队列可以写入新数据
        queue.task_done()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 06:41:11