如何提升Jupyter中两个Python程序间的任务通知与数据传输效率
针对你用轮询文件修改时间传递数据的低效问题,以下是几种更优的改进方案,均无需频繁轮询:
1. 命名管道(FIFO)—— 跨进程阻塞式通信
命名管道是文件系统中的特殊文件,支持无轮询的跨进程通信:程序2会阻塞等待管道中的数据,直到程序1写入后才继续执行,完全消除轮询开销。
程序1代码:
import pickle import os FIFO_PATH = "data_fifo" # 仅第一次运行创建管道,重复创建会报错,可保留判断 if not os.path.exists(FIFO_PATH): os.mkfifo(FIFO_PATH) def doSomeCalculations(): # 替换为你的实际计算逻辑 return {"result": 123, "status": "completed"} data_dic = doSomeCalculations() # 写入管道 with open(FIFO_PATH, 'wb') as f: pickle.dump(data_dic, f)
程序2代码:
import pickle FIFO_PATH = "data_fifo" # 阻塞等待管道数据 with open(FIFO_PATH, 'rb') as f: data = pickle.load(f) print("Running now") print("Received data:", data)
2. 多进程队列(适用于同一Jupyter Kernel下的进程)
如果两个程序在同一个Jupyter Kernel中运行,可直接用内置Queue实现阻塞式数据传递,完全不需要文件操作:
程序1(计算进程):
import multiprocessing import time def doSomeCalculations(): # 模拟计算耗时 time.sleep(2) return {"result": 456, "status": "completed"} def producer(queue): data_dic = doSomeCalculations() queue.put(data_dic) def consumer(queue): # 阻塞等待队列数据 data = queue.get() print("Running now") print("Received data:", data) if __name__ == "__main__": queue = multiprocessing.Queue() # 启动消费者(程序2逻辑) consumer_process = multiprocessing.Process(target=consumer, args=(queue,)) consumer_process.start() # 执行生产者逻辑(程序1) producer(queue) consumer_process.join()
3. 共享内存(零拷贝,极致性能)
追求最高性能时,可直接用共享内存传递数据,无需磁盘IO:
程序1代码:
import pickle import multiprocessing.shared_memory def doSomeCalculations(): return {"result": 789, "status": "completed"} data_dic = doSomeCalculations() data_bytes = pickle.dumps(data_dic) # 创建共享内存 shm = multiprocessing.shared_memory.SharedMemory(create=True, size=len(data_bytes)) shm.buf[:len(data_bytes)] = data_bytes # 将共享内存名称写入临时文件,供程序2读取 with open("shm_name.txt", 'w') as f: f.write(shm.name)
程序2代码:
import pickle import multiprocessing.shared_memory import os # 读取共享内存名称(仅一次判断,无轮询) while not os.path.exists("shm_name.txt"): pass with open("shm_name.txt", 'r') as f: shm_name = f.read().strip() os.remove("shm_name.txt") # 连接共享内存并读取数据 shm = multiprocessing.shared_memory.SharedMemory(name=shm_name) data_bytes = bytes(shm.buf) data = pickle.loads(data_bytes) # 释放共享内存 shm.close() shm.unlink() print("Running now") print("Received data:", data)
4. 信号通知+文件存储(兼容原有逻辑,消除轮询)
若必须保留文件存储,可通过信号让程序1通知程序2数据就绪,避免轮询:
程序1代码:
import pickle import os import signal def doSomeCalculations(): return {"result": 101112, "status": "completed"} data_dic = doSomeCalculations() with open("data.pkl", 'wb') as f: pickle.dump(data_dic, f) # 读取程序2的PID(程序2启动时会写入文件) with open("pid.txt", 'r') as f: program2_pid = int(f.read()) # 发送自定义信号通知程序2 os.kill(program2_pid, signal.SIGUSR1)
程序2代码:
import pickle import signal import os data = None def handle_signal(signum, frame): global data with open("data.pkl", 'rb') as f: data = pickle.load(f) # 注册信号处理函数 signal.signal(signal.SIGUSR1, handle_signal) # 将自身PID写入文件供程序1读取 with open("pid.txt", 'w') as f: f.write(str(os.getpid())) # 阻塞等待信号 while data is None: signal.pause() print("Running now") print("Received data:", data) os.remove("pid.txt")
内容的提问来源于stack exchange,提问作者user376285
相关产品推荐
相关产品推荐

