Flask视频流跨进程队列传输图像的高延迟问题排查
跨进程机器学习推流延迟排查与解决方案
延迟原因分析
- 队列帧堆积:
multiprocessing.Queue默认无长度限制,若机器学习处理速度跟不上摄像头采集速度,队列会堆积大量旧帧。即使日志显示“生成后立即取出”,前端实际拿到的仍是较早生成的帧,表现为高延迟。 - 图像序列化开销大:直接传递PIL图像、numpy数组等对象时,
Queue会用pickle序列化,大尺寸图像的序列化/反序列化耗时高,且跨进程数据拷贝量大,累积导致延迟。 - 浏览器缓存:Flask推流未设置正确的缓存控制头,浏览器缓存了旧帧,造成“延迟”的错觉。
- 文件存储方案缺陷:普通磁盘IO速度慢,或互斥锁粒度太大,导致进程等待时间过长,无法实时获取最新帧。
正确的跨进程推流方案
方案1:优化multiprocessing.Queue使用
核心优化点:
- 限制队列长度,仅保留最新帧,避免旧帧堆积。
- 提前将图像压缩为JPEG字节流,减少数据传输量与序列化开销。
代码示例:
run_ml.py
import cv2 import multiprocessing def ml_process(frame_queue): cap = cv2.VideoCapture(0) # 设置队列最大长度为1,仅保留最新帧 frame_queue.maxsize = 1 while True: ret, raw_frame = cap.read() if not ret: break # 执行机器学习处理与标注(替换为你的业务代码) # processed_frame = your_ml_inference_and_draw(raw_frame) processed_frame = raw_frame # 示例:直接使用原帧 # 压缩为JPEG字节流,降低数据量 _, encoded_img = cv2.imencode('.jpg', processed_frame, [int(cv2.IMWRITE_JPEG_QUALITY), 80]) img_bytes = encoded_img.tobytes() # 尝试放入队列,若满则丢弃旧帧再放入新帧 try: frame_queue.put(img_bytes, block=False) except multiprocessing.Queue.Full: frame_queue.get(block=False) frame_queue.put(img_bytes, block=False) cap.release()
video_stream_flask.py
from flask import Flask, Response import multiprocessing from run_ml import ml_process app = Flask(__name__) frame_queue = multiprocessing.Queue() def frame_generator(): while True: img_bytes = frame_queue.get() # 直接返回压缩后的字节流,避免重复编码 yield (b'--frame\r\n' b'Content-Type: image/jpeg\r\n' b'Cache-Control: no-cache, no-store, must-revalidate\r\n' b'Pragma: no-cache\r\n' b'Expires: 0\r\n' b'\r\n' + img_bytes + b'\r\n') @app.route('/video_feed') def video_feed(): return Response(frame_generator(), mimetype='multipart/x-mixed-replace; boundary=frame') if __name__ == '__main__': # 启动ML处理子进程 ml_proc = multiprocessing.Process(target=ml_process, args=(frame_queue,)) ml_proc.start() # 关闭Flask debug模式(debug会启动多进程,导致队列访问混乱) app.run(host='0.0.0.0', port=5000, debug=False) ml_proc.join()
方案2:使用共享内存(更低延迟)
共享内存避免了跨进程的数据拷贝,适合对延迟要求极高的场景。
代码示例:
run_ml.py
import cv2 import multiprocessing from multiprocessing import shared_memory import numpy as np def ml_process(shm_name, shm_size): # 连接共享内存 shm = shared_memory.SharedMemory(name=shm_name) # 创建共享内存对应的numpy数组 shared_buffer = np.ndarray(shape=(shm_size,), dtype=np.uint8, buffer=shm.buf) cap = cv2.VideoCapture(0) while True: ret, raw_frame = cap.read() if not ret: break # 执行机器学习处理与标注 processed_frame = raw_frame # 压缩为JPEG _, encoded_img = cv2.imencode('.jpg', processed_frame, [int(cv2.IMWRITE_JPEG_QUALITY), 80]) img_len = len(encoded_img) if img_len > shm_size - 4: # 预留4字节存储图像长度 print("警告:压缩后图像超过共享内存容量") continue # 先写入图像长度(4字节uint32),再写入图像数据 shared_buffer[:4] = np.array([img_len], dtype=np.uint32).tobytes() shared_buffer[4:4+img_len] = encoded_img.tobytes() cap.release() shm.close()
video_stream_flask.py
from flask import Flask, Response import multiprocessing from multiprocessing import shared_memory import numpy as np from run_ml import ml_process app = Flask(__name__) def frame_generator(shm_name, shm_size): shm = shared_memory.SharedMemory(name=shm_name) shared_buffer = np.ndarray(shape=(shm_size,), dtype=np.uint8, buffer=shm.buf) while True: # 读取图像长度 img_len = np.frombuffer(shared_buffer[:4], dtype=np.uint32)[0] if img_len == 0: continue # 读取图像数据 img_bytes = shared_buffer[4:4+img_len].tobytes() yield (b'--frame\r\n' b'Content-Type: image/jpeg\r\n' b'Cache-Control: no-cache\r\n' b'\r\n' + img_bytes + b'\r\n') shm.close() if __name__ == '__main__': # 创建共享内存(1MB足够存储大部分压缩后的JPEG图像) shm_size = 1024 * 1024 shm = shared_memory.SharedMemory(create=True, size=shm_size) # 启动ML子进程 ml_proc = multiprocessing.Process(target=ml_process, args=(shm.name, shm_size)) ml_proc.start() # 注册推流路由 @app.route('/video_feed') def video_feed(): return Response(frame_generator(shm.name, shm_size), mimetype='multipart/x-mixed-replace; boundary=frame') app.run(host='0.0.0.0', port=5000, debug=False) ml_proc.join() # 清理共享内存 shm.unlink()
额外注意事项
- 关闭Flask的
debug模式:debug模式会自动启动多进程,导致队列/共享内存被多个进程访问,引发帧堆积或数据错乱。 - 优化ML处理速度:若ML模型推理耗时过长,即使IPC优化到位,仍会导致延迟。可通过降低摄像头分辨率、帧率,或使用TensorRT、ONNX Runtime等工具加速模型。
- 验证网络延迟:确保前端与Flask服务在同一局域网,排除网络传输导致的延迟。
内容的提问来源于stack exchange,提问作者Hakim_C
相关产品推荐
相关产品推荐

