ESP32 JPEG视频流转发求助:解决家庭监控多客户端访问问题
解决方案:ESP32 JPEG流转发器实现(NodeJS/Python)
核心原理
ESP32输出的是连续拼接的JPEG帧序列(无HTTP多部分分隔),转发器需要完成三个核心动作:
- 与每个ESP32建立唯一长连接,持续读取原始流
- 通过JPEG帧的特征标记(起始
0xFFD8、结束0xFFD9)拆分完整帧 - 将帧封装为标准MJPEG格式,分发给所有连接的客户端
NodeJS 实现方案
使用NodeJS的http模块搭建转发服务器,内置流处理和重连机制:
const http = require('http'); // 摄像头配置:ID -> ESP32流地址 const cameras = { cam1: 'http://192.168.1.100/stream', cam2: 'http://192.168.1.101/stream' // 扩展到30台摄像头 }; // 存储每个摄像头的客户端集合与连接状态 const cameraState = new Map(); // 处理ESP32原始流,拆分JPEG帧并推送 function processCameraStream(camId, espResponse) { let buffer = Buffer.alloc(0); const clients = new Set(); cameraState.set(camId, { clients }); espResponse.on('data', (chunk) => { buffer = Buffer.concat([buffer, chunk]); // 循环查找完整JPEG帧 let startIdx = buffer.indexOf(0xFFD8); while (startIdx !== -1) { const endIdx = buffer.indexOf(0xFFD9, startIdx + 2); if (endIdx === -1) break; // 未找到完整帧,等待后续数据 const frame = buffer.slice(startIdx, endIdx + 2); buffer = buffer.slice(endIdx + 2); // 保留剩余数据 // 推送给所有客户端 clients.forEach((res) => { res.write(`--frame\r\nContent-Type: image/jpeg\r\nContent-Length: ${frame.length}\r\n\r\n`); res.write(frame); res.write('\r\n'); }); startIdx = buffer.indexOf(0xFFD8); } }); // 连接断开时自动重连 espResponse.on('end', () => { console.log(`Cam ${camId} stream ended, reconnecting...`); connectToCamera(camId); }); espResponse.on('error', (err) => { console.error(`Cam ${camId} stream error:`, err); connectToCamera(camId); }); } // 建立与ESP32的长连接 function connectToCamera(camId) { const camUrl = cameras[camId]; if (!camUrl) return; const req = http.get(camUrl, (res) => { console.log(`Connected to cam ${camId}`); processCameraStream(camId, res); }); req.on('error', (err) => { console.error(`Failed to connect cam ${camId}:`, err); setTimeout(() => connectToCamera(camId), 5000); // 5秒后重试 }); } // 启动转发服务器 const server = http.createServer((req, res) => { const camId = req.url.slice(1); if (!cameras[camId]) { res.writeHead(404); res.end('Camera not found'); return; } // 返回标准MJPEG响应头 res.writeHead(200, { 'Content-Type': 'multipart/x-mixed-replace; boundary=frame', 'Cache-Control': 'no-cache', 'Connection': 'keep-alive' }); // 将客户端加入对应摄像头的订阅列表 const camData = cameraState.get(camId); if (camData) { camData.clients.add(res); req.on('close', () => camData.clients.delete(res)); } else { // 首次访问时建立ESP32连接 connectToCamera(camId); setTimeout(() => { const newCamData = cameraState.get(camId); if (newCamData) { newCamData.clients.add(res); req.on('close', () => newCamData.clients.delete(res)); } else { res.writeHead(503); res.end('Failed to connect to camera'); } }, 1000); } }); const PORT = 3000; server.listen(PORT, () => { console.log(`Relay server running on port ${PORT}`); // 预连接所有摄像头 Object.keys(cameras).forEach(connectToCamera); });
Python 实现方案
使用aiohttp实现异步转发,更适合高并发客户端场景:
import asyncio from aiohttp import web, ClientSession, ClientTimeout # 摄像头配置 CAMERAS = { "cam1": "http://192.168.1.100/stream", "cam2": "http://192.168.1.101/stream" # 扩展到30台摄像头 } # 存储每个摄像头的订阅者集合 camera_subscribers = {cam: set() for cam in CAMERAS} async def fetch_and_relay(cam_id): """持续读取ESP32流并拆分转发""" url = CAMERAS[cam_id] buffer = b"" while True: try: async with ClientSession(timeout=ClientTimeout(total=None)) as session: async with session.get(url) as resp: if resp.status != 200: print(f"Cam {cam_id} returned {resp.status}, retrying...") await asyncio.sleep(5) continue print(f"Connected to cam {cam_id}") # 分块读取流 async for chunk in resp.content.iter_chunked(1024): buffer += chunk # 查找JPEG帧边界 start_idx = buffer.find(b"\xff\xd8") while start_idx != -1: end_idx = buffer.find(b"\xff\xd9", start_idx + 2) if end_idx == -1: break # 提取完整帧 frame = buffer[start_idx:end_idx+2] buffer = buffer[end_idx+2:] # 推送给所有订阅者 for subscriber in camera_subscribers[cam_id].copy(): try: await subscriber.send_bytes(b"--frame\r\nContent-Type: image/jpeg\r\nContent-Length: %d\r\n\r\n" % len(frame)) await subscriber.send_bytes(frame) await subscriber.send_bytes(b"\r\n") except Exception as e: print(f"Failed to send to subscriber: {e}") camera_subscribers[cam_id].remove(subscriber) start_idx = buffer.find(b"\xff\xd8") except Exception as e: print(f"Cam {cam_id} stream error: {e}, retrying in 5s...") await asyncio.sleep(5) async def handle_client(request): """处理客户端请求""" cam_id = request.match_info.get("cam_id") if cam_id not in CAMERAS: return web.Response(status=404, text="Camera not found") # 设置MJPEG响应头 headers = { "Content-Type": "multipart/x-mixed-replace; boundary=frame", "Cache-Control": "no-cache", "Connection": "keep-alive" } resp = web.StreamResponse(headers=headers) await resp.prepare(request) # 添加到订阅列表 camera_subscribers[cam_id].add(resp) try: await request.wait_for_disconnection() finally: camera_subscribers[cam_id].remove(resp) return resp async def main(): app = web.Application() app.add_routes([web.get("/{cam_id}", handle_client)]) # 启动所有摄像头的转发任务 for cam_id in CAMERAS: asyncio.create_task(fetch_and_relay(cam_id)) runner = web.AppRunner(app) await runner.setup() site = web.TCPSite(runner, "0.0.0.0", 3000) await site.start() print("Relay server running on http://0.0.0.0:3000") await asyncio.Event().wait() if __name__ == "__main__": asyncio.run(main())
关键注意事项
- 帧拆分逻辑:必须依赖JPEG的二进制特征(
0xFFD8起始、0xFFD9结束)拆分帧,ESP32的裸流没有HTTP分隔符,这是读取流的核心。 - MJPEG格式兼容:转发给客户端时必须封装为标准MJPEG多部分格式,否则浏览器/仪表盘无法解析播放。
- 资源管理:当客户端断开连接时,必须从订阅列表中移除,避免内存泄漏。
- 重连机制:ESP32连接不稳定时,转发器需自动重试,保证服务连续性。
内容的提问来源于stack exchange,提问作者Alexey Khachatryan
相关产品推荐
相关产品推荐

