MJPEG 视频流 Java 转发
·
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.");
}
};
}
}
更多推荐




所有评论(0)