基于Django Channels+OpenCV:如何从consumers.py获取批量帧实现视频分类?
解决方案:基于Django Channels+OpenCV实现批量视频帧获取与分类
问题背景
我正在开发一个基于OpenCV和Django的计算机视觉项目,涉及目标检测与实时帧流传输,当前使用Django Channels实现WebSocket通信。现在遇到新需求:需执行视频分类,这需要向模型发送批量帧。我在前端添加了按钮,期望获取从起始帧到点击按钮时对应帧的时间戳,但不知如何实现批量帧的获取,寻求技术指导。
一、前端逻辑优化
当前前端按钮事件绑定在onmessage回调内,会导致重复绑定且未记录起始时间戳,修改如下:
<script type="text/javascript"> document.addEventListener("DOMContentLoaded", function () { let videoPath = "{{ video_path }}"; let liveImage = document.getElementById("liveImage"); let startTimestamp = null; let currentTimestamp = ""; // 建立WebSocket连接 let wsProtocol = window.location.protocol === "https:" ? "wss" : "ws"; let wsURL = `${wsProtocol}://${window.location.host}/ws/video/${videoPath}/`; let socket = new WebSocket(wsURL); // 绑定按钮点击事件(仅绑定一次) document.getElementById('sendTimestampButton').addEventListener('click', function () { if (!startTimestamp) { // 第一次点击记录起始时间戳 startTimestamp = currentTimestamp; this.textContent = "结束截取"; console.log("已记录起始时间戳:", startTimestamp); } else { // 第二次点击发送时间戳区间请求 this.textContent = "开始截取"; socket.send(JSON.stringify({ 'action': 'get_batch_frames', 'start_ts': startTimestamp, 'end_ts': currentTimestamp })); startTimestamp = null; // 重置状态 } }); socket.onmessage = function (event) { let data = JSON.parse(event.data); // 更新当前帧的时间戳 currentTimestamp = data.timestamp; liveImage.src = "data:image/jpeg;base64," + data.frame_data; }; socket.onopen = function (event) { console.log('WebSocket连接已建立') }; socket.onclose = function (event) { console.log("WebSocket连接已关闭。"); }; socket.onerror = function (error) { console.error("WebSocket错误:", error); }; }); </script>
二、后端Consumer改造
新增帧缓存机制,存储时间戳与帧的映射,并处理前端的批量帧请求:
import cv2 import json import asyncio import datetime import numpy as np from .processing import serialize_frame from channels.generic.websocket import AsyncWebsocketConsumer class VideoConsumer(AsyncWebsocketConsumer): def __init__(self, *args, **kwargs): super().__init__(*args, **kwargs) self.frame_cache = {} # 以时间戳字符串为key存储帧数据 self.cap = None self.running = True # 控制帧发送循环的开关 async def connect(self): await self.accept() video_file_path = self.scope['url_route']['kwargs']['video_path'] # 异步启动帧发送任务 asyncio.create_task(self.send_processed_frames(video_file_path)) async def send_frame(self, frame_data, timestamp_str): await self.send(text_data=json.dumps({ 'frame_data': frame_data, 'timestamp': timestamp_str })) async def send_processed_frames(self, video_file_path): try: self.cap = cv2.VideoCapture(video_file_path) fps = self.cap.get(cv2.CAP_PROP_FPS) frame_interval = 1 / fps # 用视频实际帧率控制发送间隔 while self.running: ret, frame = self.cap.read() if not ret: break # 视频播放完毕退出循环 # 获取当前帧的毫秒级时间戳并格式化 timestamp_ms = self.cap.get(cv2.CAP_PROP_POS_MSEC) timestamp = datetime.timedelta(milliseconds=timestamp_ms) timestamp_str = str(timestamp).split('.', 2)[0] # 缓存帧(视频较长时建议限制缓存大小,比如只保留最近30秒的帧) self.frame_cache[timestamp_str] = frame.copy() frame_data = serialize_frame(frame) await self.send_frame(frame_data, timestamp_str) await asyncio.sleep(frame_interval) self.cap.release() except Exception as e: await self.send(text_data=f"Error: {str(e)}") async def receive(self, text_data): message = json.loads(text_data) if message.get('action') == 'get_batch_frames': start_ts = message['start_ts'] end_ts = message['end_ts'] await self.process_batch_frames(start_ts, end_ts) async def process_batch_frames(self, start_ts, end_ts): # 将时间戳字符串转换为秒数,方便区间判断 def ts_to_seconds(ts_str): h, m, s = map(int, ts_str.split(':')) return h * 3600 + m * 60 + s start_sec = ts_to_seconds(start_ts) end_sec = ts_to_seconds(end_ts) # 收集区间内的帧 batch_frames = [] for ts_str, frame in self.frame_cache.items(): current_sec = ts_to_seconds(ts_str) if start_sec <= current_sec <= end_sec: batch_frames.append(frame) if not batch_frames: await self.send(text_data=json.dumps({'status': 'error', 'message': '未找到区间内的帧'})) return # 替换为你的视频分类模型推理逻辑 # prediction = your_video_classification_model.predict(batch_frames) prediction = "示例分类结果:户外街道场景" # 返回结果给前端 await self.send(text_data=json.dumps({ 'status': 'success', 'start_ts': start_ts, 'end_ts': end_ts, 'frame_count': len(batch_frames), 'prediction': prediction })) async def disconnect(self, close_code): self.running = False if self.cap: self.cap.release()
三、关键注意事项
- 缓存大小控制:若视频较长,直接缓存所有帧会占用大量内存,建议设置缓存上限(如仅保留最近30秒的帧),或改用磁盘临时存储。
- 异步推理:如果模型推理耗时较长,建议将推理任务放到Celery等异步队列中,避免阻塞WebSocket连接。
- 时间戳精度:如需更精确的帧匹配,可保留毫秒级时间戳(不分割字符串的小数部分)。
内容的提问来源于stack exchange,提问作者Shah Zeb
相关产品推荐
相关产品推荐

