多子进程调用父进程函数 解决数据库过载与文件读取延迟问题
替代文件共享的进程间数据同步方案
针对你遇到的文件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
相关产品推荐
相关产品推荐

