MJPEG (Motion JPEG) 是一种视频压缩格式,它将视频的每一帧都独立编码为 JPEG 图像,然后按顺序播放这些图像形成视频效果。实际上这种效果并不是很好,冗余的数据比较大,但胜在发送和转发包括播放都比较简单,也是许多低端 MCU 可以推送视频流的一种方式。

同样在 Java 中处理也比较简单。

下面 MjpegFrameService用于存储和分发 MJPEG 帧的服务。使用 BlockingQueue 实现简单的生产者-消费者模型, * 注意:这是一个简化的单例实现,适用于单实例场景。生产环境可能需要考虑持久化、容量限制、分布式等。

import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;

import java.util.Map;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.atomic.AtomicLong;

/**
 * 用于存储和分发 MJPEG 帧的服务。
 * 使用 BlockingQueue 实现简单的生产者-消费者模型。
 * 注意:这是一个简化的单例实现,适用于单实例场景。
 * 生产环境可能需要考虑持久化、容量限制、分布式等。
 */
@Service
@Slf4j
public class MjpegFrameService {
    private final AtomicLong droppedFramesCounter = new AtomicLong(0); // 计数器,线程安全

    // 使用 LinkedBlockingQueue 作为帧队列,有界队列可以防止单个流无限增长
    // 容量 10 是示例,实际应根据帧率和延迟容忍度调整
    private final BlockingQueue<byte[]> frameQueue = new LinkedBlockingQueue<>(30);

    /**
     * 存储不同设备的帧队列
     */
    private final Map<String, BlockingQueue<byte[]>> deviceFrameQueues = new ConcurrentHashMap<>();

    /**
     * 每个设备队列的容量
     */
    private static final int QUEUE_CAPACITY = 30;

    /**
     * 接收并存储一个新的 MJPEG 帧。
     * 如果队列已满,则丢弃最旧的帧以腾出空间给新帧。
     *
     * @param macAddress 设备 MAC 地址
     * @param frameData  帧的字节数据
     * @return true 表示操作完成(即使有帧被丢弃),总是尝试存储新帧
     */
    public boolean storeFrame(String macAddress, byte[] frameData) {
        if (frameData == null) {
            log.warn("Attempted to store a null frame for device: {}", macAddress);
            return false;
        }

        // 获取或创建设备对应的队列
        BlockingQueue<byte[]> frameQueue = deviceFrameQueues.computeIfAbsent(macAddress, k -> new LinkedBlockingQueue<>(QUEUE_CAPACITY));

        if (frameQueue.remainingCapacity() == 0) { // 检查队列是否已满
            // 队列已满,移除最旧的帧
            byte[] discardedFrame = frameQueue.poll(); // poll() 非阻塞,如果队列空则返回 null

            if (discardedFrame != null) {
                // 实际移除了一个帧
                long totalDropped = droppedFramesCounter.incrementAndGet(); // 原子性增加计数器
                log.debug("Frame queue full for device {}. Discarded oldest frame (size: {} bytes). Total dropped so far: {}",
                        macAddress, discardedFrame.length, totalDropped);
                // 可以在这里添加更详细的丢帧日志或监控逻辑
            } else
                // 理论上不应该发生,因为 capacity=0 意味着队列不空,但以防万一
                log.warn("Queue reported full but poll() returned null for device: {}", macAddress);
        }

        // 此时队列肯定有空间(因为我们刚刚移除了一个元素)
        boolean offered = frameQueue.offer(frameData); // offer() 非阻塞,此时应该总会成功

        if (offered)
            log.trace("Successfully stored new frame (size: {} bytes) for device {} into queue.",
                    frameData.length, macAddress);
        else
            // 理论上不应该发生,因为我们已经确保了空间。记录严重警告。
            log.error("Failed to offer frame to queue even after making space for device: {}!", macAddress);

        return offered; // 因为我们努力确保存储,通常返回 true
    }


    /**
     * 从队列中取出下一个可用的帧。
     * 这是一个阻塞操作,直到有帧可用。
     *
     * @param macAddress 设备MAC地址
     * @return 帧的字节数据
     * @throws InterruptedException 如果等待时被中断
     */
    public byte[] takeFrame(String macAddress) throws InterruptedException {
        BlockingQueue<byte[]> queue = deviceFrameQueues.get(macAddress);
        if (queue == null)
            throw new IllegalStateException("No frame queue found for device: " + macAddress);

        return queue.take();// take 会阻塞直到有元素
    }

    /**
     * 尝试立即取出一个帧,如果队列为空则返回 null。
     *
     * @param macAddress 设备 MAC 地址
     * @return 帧的字节数据或 null
     */
    public byte[] pollFrame(String macAddress) {
        BlockingQueue<byte[]> queue = deviceFrameQueues.get(macAddress);
        if (queue == null)
            return null;

        return queue.poll();// poll 立即返回,不阻塞
    }

    /**
     * 清空队列
     *
     * @param macAddress 设备 MAC 地址
     */
    public void clear(String macAddress) {
        BlockingQueue<byte[]> queue = deviceFrameQueues.remove(macAddress);

        if (queue != null) {
            int sizeBeforeClear = frameQueue.size();
            frameQueue.clear();

            if (sizeBeforeClear > 0)
                log.info("Frame queue cleared, {} frames were discarded.", sizeBeforeClear);
            else
                log.debug("Frame queue cleared (was already empty).");
        } else
            log.debug("No frame queue found for device {} to clear.", macAddress);
    }

    /**
     * 获取自服务启动以来丢弃的总帧数
     *
     * @return 丢弃的帧数
     */
    public long getDroppedFramesCount() {
        return droppedFramesCounter.get();
    }
}

下面这个控制器方法,包括了直播视频的推流和拉流,均支持。

import cn.wyndme.service.MjpegFrameService;
import com.ajaxjs.framework.database.IgnoreDataBaseConnect;
import com.ajaxjs.framework.mvc.unifiedreturn.PureOutput;
import com.ajaxjs.iam.annotation.AllowOpenAccess;
import com.ajaxjs.util.ObjectHelper;
import jakarta.servlet.http.HttpServletRequest;
import jakarta.servlet.http.HttpServletResponse;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.http.HttpHeaders;
import org.springframework.http.HttpStatus;
import org.springframework.http.ResponseEntity;
import org.springframework.web.bind.annotation.*;
import org.springframework.web.servlet.mvc.method.annotation.StreamingResponseBody;

import java.io.IOException;
import java.io.InputStream;
import java.net.HttpURLConnection;
import java.net.URL;
import java.util.Enumeration;

/**
 * 视频
 */
@RestController
@RequestMapping("/video")
@Slf4j
public class VideoController {
    private String esp32StreamUrl = "http://192.168.110.203/stream";

    private int connectTimeoutMs = 5000;

    private int readTimeoutMs = 10000;

    private static final String MJPEG_CONTENT_TYPE = "multipart/x-mixed-replace; boundary=--frameboundary";

    /**
     * 使用 StreamingResponseBody 和 HttpURLConnection 实现流式转发。
     *
     * @param response The HttpServletResponse object.
     * @return A StreamingResponseBody that handles the streaming logic.
     */
    @GetMapping(produces = MJPEG_CONTENT_TYPE)
    @AllowOpenAccess
    @IgnoreDataBaseConnect
    public StreamingResponseBody proxyMjpegStream(HttpServletResponse response) {
        return outputStream -> {
            HttpURLConnection connection = null; // connection 不是 AutoCloseable

            try {
                URL url = new URL(esp32StreamUrl);
                connection = (HttpURLConnection) url.openConnection();
                connection.setRequestMethod("GET");
                connection.setConnectTimeout(connectTimeoutMs);
                connection.setReadTimeout(readTimeoutMs);

                int responseCode = connection.getResponseCode();
                log.debug("Received status code {} from ESP32 stream.", responseCode);

                if (responseCode != HttpURLConnection.HTTP_OK) {
                    log.error("Failed to fetch stream from ESP32. Status: {}", responseCode);
                    response.setStatus(responseCode);
                    outputStream.write(("Error fetching stream: " + responseCode).getBytes());

                    return; // 提前返回,connection.disconnect() 会在 finally 中调用
                }

                // 获取并转发 Content-Type (优先) 或设置默认值
                String contentTypeFromEsp = connection.getContentType();
                if (contentTypeFromEsp != null && !contentTypeFromEsp.isEmpty()) {
                    response.setHeader(HttpHeaders.CONTENT_TYPE, contentTypeFromEsp);
                    log.debug("Forwarding Content-Type from ESP32: {}", contentTypeFromEsp);
                } else {
                    response.setHeader(HttpHeaders.CONTENT_TYPE, MJPEG_CONTENT_TYPE);
                    log.debug("Setting default MJPEG Content-Type: {}", MJPEG_CONTENT_TYPE);
                }

                // --- 使用 try-with-resources 管理 InputStream ---
                // connection.getInputStream() 返回的 InputStream 通常是 AutoCloseable 的
                try (InputStream inputStream = connection.getInputStream()) { // inputStream 会自动关闭
                    log.debug("Starting to stream MJPEG data...");

                    byte[] buffer = new byte[8192];
                    int bytesRead;

                    while ((bytesRead = inputStream.read(buffer)) != -1) {
                        outputStream.write(buffer, 0, bytesRead);
                        outputStream.flush();
                    }

                    log.debug("Finished streaming MJPEG data.");
                } // inputStream 在此处自动关闭

            } catch (IOException e) {
                String msg = e.getMessage();

                if (msg != null && (msg.contains("Broken pipe") || msg.contains("Connection reset")))
                    log.warn("Client disconnected or connection lost during streaming.");
                else {
                    log.error("Network error while connecting to or reading from ESP32 stream: {}", e.getMessage(), e);
                    response.setStatus(HttpStatus.INTERNAL_SERVER_ERROR.value());

                    try {
                        if (!response.isCommitted()) {
                            outputStream.write(("Proxy error: " + e.getMessage()).getBytes());
                        }
                    } catch (IOException ioException) {
                        log.warn("Could not write error message to output stream.", ioException);
                    }
                }
            } finally {
                // connection 不是 AutoCloseable,必须手动 disconnect
                if (connection != null) {
                    connection.disconnect(); // 即使 inputStream 已关闭,也应 disconnect connection
                    log.trace("Disconnected HttpURLConnection.");
                }
            }
        };
    }

    @Autowired
    private MjpegFrameService mjpegFrameService;

    final static private String DEFAULT_MAC = "default";
    final static private String HEAD_KEY = "Device-Id";

    /**
     * 接收 ESP32 推送的 MJPEG 流端点。
     * ESP32 应该向这个 POST 端点发送原始的 MJPEG 数据流。
     *
     * @param mac     设备 MAC 地址
     * @param request HttpServletRequest containing the stream data.
     * @return ResponseEntity indicating success or failure.
     */
    @PutMapping("/receive-stream")
    @ResponseBody
    @AllowOpenAccess
    @IgnoreDataBaseConnect
    @PureOutput
    public ResponseEntity<String> receiveStream(@RequestParam(required = false) String mac, HttpServletRequest request) {
        if (ObjectHelper.isEmptyText(mac)) // 如果没有提供mac参数,则尝试从请求头中获取
            mac = request.getHeader(HEAD_KEY);

        if (ObjectHelper.isEmptyText(mac)) // 如果仍然没有MAC地址,则使用默认值
            mac = DEFAULT_MAC;

        mac = mac.replaceAll(":", "").toUpperCase();
        log.info("Receiving stream data from ESP32...mac: " + mac);

        log.info("======= 开始记录请求头 =======");
        Enumeration<String> headerNames = request.getHeaderNames();

        if (headerNames != null) {
            while (headerNames.hasMoreElements()) {
                String headerName = headerNames.nextElement();
                String headerValue = request.getHeader(headerName);

                // 3. 记录每个请求头
                log.info("请求头 - {}: {}", headerName, headerValue);
                // 或者你想用 debug 级别:
                // log.debug("请求头 - {}: {}", headerName, headerValue);
            }
        } else
            log.info("未找到任何请求头。");
        log.info("======= 结束记录请求头 =======");

        try (InputStream inputStream = request.getInputStream()) {
            // 注意:这里假设 ESP32 是连续推送数据,没有明确的消息边界。
            // 我们需要一种方式来分割帧。最简单的是按固定大小读取(不理想),
            // 或者 ESP32 在推送时按 MJPEG 标准格式(带有 boundary)推送。
            // 为了简化,我们在这里演示读取整个流(如果是一次性推送)或分块读取。

            // 方案 A: 如果 ESP32 一次推送一个完整的帧 (不太常见)
            /*
            byte[] frameData = inputStream.readAllBytes(); // JDK 9+
            if (frameData.length > 0) {
                 mjpegFrameService.storeFrame(frameData);
                 log.debug("Stored frame of size: {} bytes", frameData.length);
            }
            */

            // 方案 B: 如果 ESP32 持续推送流 (更常见)
            // 我们需要解析 MJPEG 流。这比较复杂。
            // 简化处理:假设 ESP32 按照某种方式(例如,每次推送一个完整的 multipart part)
            // 或者我们自己读取并分割。这里演示一种基础的流式读取和存储方式。
            // 但这要求我们知道帧边界,或者 ESP32 以特定方式推送。

            // --- 简化模型:假设 ESP32 推送的是原始的、连续的 MJPEG 流 ---
            // 我们需要解析它。这通常涉及查找 boundary。
            // 为了演示,我们假设 ESP32 已经按照标准格式推送了整个 multipart/mixed 流。
            // 我们将整个流存储起来。但这不是最佳实践,因为 MJPEG 是持续的。
            // 更好的做法是 ESP32 推送单个帧,或者我们在这里解析流。

            // --- 更实际的简化模型:ESP32 推送单个 JPEG 文件 ---
            // 如果 ESP32 每次拍照后 POST 一张 JPEG 图片,我们可以这样处理:
            byte[] buffer = new byte[50 * 1024]; // 50KB buffer, adjust as needed
            int totalBytesRead = 0;
            int bytesRead;

            while ((bytesRead = inputStream.read(buffer, totalBytesRead, buffer.length - totalBytesRead)) != -1) {
                totalBytesRead += bytesRead;

                if (totalBytesRead == buffer.length) {
                    // Buffer full, need bigger buffer or handle chunking differently
                    log.warn("Buffer might be too small for incoming data chunk.");
                    // For simplicity, we'll just use what we have
                    break; // Or reallocate buffer
                }
            }

            if (totalBytesRead > 0) {
                byte[] actualData = new byte[totalBytesRead];
                System.arraycopy(buffer, 0, actualData, 0, totalBytesRead);
                boolean stored = mjpegFrameService.storeFrame(mac, actualData);

                if (stored)
                    log.debug("Stored received data chunk/frame of size: {} bytes", actualData.length);
                else
                    log.warn("Dropped received data chunk/frame of size: {} bytes due to full queue", actualData.length);

                return ResponseEntity.ok("Data received and stored (size: " + actualData.length + " bytes)");
            } else {
                log.warn("Received empty stream data.");
                return ResponseEntity.badRequest().body("Empty stream data received.");
            }
        } catch (IOException e) {
            log.error("Error reading stream data from ESP32: {}", e.getMessage(), e);
            return ResponseEntity.status(HttpStatus.INTERNAL_SERVER_ERROR).body("Failed to read stream data.");
        } catch (Exception e) { // Catch other potential exceptions
            log.error("Unexpected error processing received stream: {}", e.getMessage(), e);
            return ResponseEntity.status(HttpStatus.INTERNAL_SERVER_ERROR).body("Internal server error.");
        }
    }

    /**
     * 向浏览器转发 MJPEG 流的端点。
     * 浏览器访问此 GET 端点以观看视频。
     *
     * @param mac      设备 MAC 地址
     * @param response HttpServletResponse to write the stream to.
     * @return StreamingResponseBody for continuous streaming.
     */
    @GetMapping(value = "/stream", produces = MJPEG_CONTENT_TYPE)
    @AllowOpenAccess
    @IgnoreDataBaseConnect
    public StreamingResponseBody streamToBrowser(@RequestParam(required = false) String mac, HttpServletResponse response) {
        log.info("Browser client connected to stream endpoint.");
        response.setHeader(HttpHeaders.CONTENT_TYPE, MJPEG_CONTENT_TYPE);// 设置响应头 - 由 produces 属性处理,但显式设置也无妨
        response.setHeader(HttpHeaders.CONTENT_ENCODING, "identity"); // 确保没有不必要的压缩

        if (mac == null || mac.isEmpty())
            mac = DEFAULT_MAC;

        final String deviceMac = mac.toUpperCase();
        return outputStream -> {
            try {
                // 写入初始 boundary (MJPEG 标准的一部分)
                outputStream.write("--frameboundary\r\n".getBytes());
                outputStream.flush();

                while (!Thread.currentThread().isInterrupted()) { // 允许优雅关闭
                    try {
                        // 从服务中获取下一帧数据(阻塞直到有数据)
                        byte[] frameData = mjpegFrameService.takeFrame(deviceMac);
                        // 写入帧的 Content-Type 头部 (JPEG)
                        outputStream.write("Content-Type: image/jpeg\r\n".getBytes());
                        // 写入帧的 Content-Length 头部
                        outputStream.write(("Content-Length: " + frameData.length + "\r\n").getBytes());
                        outputStream.write("\r\n".getBytes()); // 写入空行,分隔头部和数据
                        outputStream.write(frameData); // 写入帧的实际 JPEG 数据
                        outputStream.write("\r\n--frameboundary\r\n".getBytes());// 写入帧结束标记和下一个 boundary
                        outputStream.flush(); // 刷新到客户端
                        log.trace("Forwarded frame of size: {} bytes", frameData.length);
                    } catch (InterruptedException e) {
                        log.info("Streaming thread interrupted, stopping stream.");
                        Thread.currentThread().interrupt(); // Preserve interrupt status
                        break;

                    } catch (IOException e) {
                        String msg = e.getMessage();

                        if (msg != null && (msg.contains("Broken pipe") || msg.contains("Connection reset")))
                            log.info("Browser client disconnected.");
                        else
                            log.error("Error writing to browser output stream: {}", e.getMessage(), e);

                        break; // Exit loop on client disconnect or write error
                    } catch (IllegalStateException e) {
                        log.warn("No frame queue found for device: {}. Error: {}", deviceMac, e.getMessage());
                        // 发送错误信息给客户端
                        outputStream.write("Content-Type: text/plain\r\n".getBytes());
                        outputStream.write(("Content-Length: " + e.getMessage().length() + "\r\n").getBytes());
                        outputStream.write("\r\n".getBytes());
                        outputStream.write(e.getMessage().getBytes());
                        outputStream.write("\r\n--frameboundary\r\n".getBytes());
                        outputStream.flush();
                        break;
                    }
                }
            } finally {
                log.info("Browser streaming session ended.");
            }
        };
    }
}
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐