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

ESP32 JPEG视频流转发求助:解决家庭监控多客户端访问问题

解决方案:ESP32 JPEG流转发器实现(NodeJS/Python)

核心原理

ESP32输出的是连续拼接的JPEG帧序列(无HTTP多部分分隔),转发器需要完成三个核心动作:

  1. 与每个ESP32建立唯一长连接,持续读取原始流
  2. 通过JPEG帧的特征标记(起始0xFFD8、结束0xFFD9)拆分完整帧
  3. 将帧封装为标准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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 11:57:33