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

Flask视频流跨进程队列传输图像的高延迟问题排查

跨进程机器学习推流延迟排查与解决方案

延迟原因分析

  1. 队列帧堆积:multiprocessing.Queue默认无长度限制,若机器学习处理速度跟不上摄像头采集速度,队列会堆积大量旧帧。即使日志显示“生成后立即取出”,前端实际拿到的仍是较早生成的帧,表现为高延迟。
  2. 图像序列化开销大:直接传递PIL图像、numpy数组等对象时,Queue会用pickle序列化,大尺寸图像的序列化/反序列化耗时高,且跨进程数据拷贝量大,累积导致延迟。
  3. 浏览器缓存:Flask推流未设置正确的缓存控制头,浏览器缓存了旧帧,造成“延迟”的错觉。
  4. 文件存储方案缺陷:普通磁盘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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 20:05:57