liusuyi
2026-05-30 f4f4fc53260eb67483dce406a963628273786a61
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
package com.ard.work.sdk.hb.service;
 
import java.util.Map;
import java.util.Optional;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.TimeUnit;
 
/**
 * 相机推流队列管理器(静态版本)
 * 全局唯一的队列管理器,无需创建实例,直接通过类名调用
 * 每个相机(设备ID/通道号)对应独立的阻塞队列,存储待推流的H264数据
 */
public class CameraStreamQueueManager {
 
    /**
     * 静态队列容器:Key=相机唯一标识,Value=该相机的H264数据队列
     * ConcurrentHashMap保证高并发下的线程安全
     */
    private static final ConcurrentHashMap<String, BlockingQueue<byte[]>> CAMERA_QUEUE_MAP = new ConcurrentHashMap<>();
 
    /**
     * 每个相机队列的默认容量
     */
    private static final int DEFAULT_QUEUE_CAPACITY = 10;
 
    // 私有化构造方法,禁止创建实例(核心:确保静态单例)
    private CameraStreamQueueManager() {
        throw new UnsupportedOperationException("该类为静态工具类,禁止实例化!");
    }
 
    // ===================== 静态核心操作 =====================
 
    /**
     * 创建指定相机的独立队列(不存在则创建,已存在则返回现有队列)
     * @param cameraId 相机唯一标识(不可为空)
     * @return 该相机的阻塞队列
     */
    public static BlockingQueue<byte[]> createCameraQueue(String cameraId) {
        if (cameraId == null || cameraId.trim().isEmpty()) {
            throw new IllegalArgumentException("相机唯一标识不能为空!");
        }
        // ConcurrentHashMap的原子操作,保证线程安全
        return CAMERA_QUEUE_MAP.computeIfAbsent(cameraId.trim(), k -> new LinkedBlockingQueue<>(DEFAULT_QUEUE_CAPACITY));
    }
 
    /**
     * 向指定相机的队列中添加H264数据(非阻塞,队列满则返回false)
     * @param cameraId 相机唯一标识
     * @param h264Data H264裸流数据(完整NALU单元)
     * @return true=入队成功,false=队列满/队列不存在
     */
    public static boolean offerDataToQueue(String cameraId, byte[] h264Data) {
        if (h264Data == null || h264Data.length == 0) {
            return false;
        }
        BlockingQueue<byte[]> queue = CAMERA_QUEUE_MAP.get(cameraId);
        return queue != null && queue.offer(h264Data);
    }
 
    /**
     * 向指定相机的队列中添加H264数据(阻塞式,队列满则等待超时)
     * @param cameraId 相机唯一标识
     * @param h264Data H264裸流数据
     * @param timeout 超时时间(单位:毫秒)
     * @return true=入队成功,false=超时/队列不存在
     * @throws InterruptedException 线程中断异常
     */
    public static boolean putDataToQueueWithTimeout(String cameraId, byte[] h264Data, long timeout) throws InterruptedException {
        if (h264Data == null || h264Data.length == 0) {
            return false;
        }
        BlockingQueue<byte[]> queue = CAMERA_QUEUE_MAP.get(cameraId);
        if (queue == null) {
            return false;
        }
        return queue.offer(h264Data, timeout, TimeUnit.MILLISECONDS);
    }
 
    /**
     * 从指定相机的队列中获取H264数据(阻塞式,队列为空则等待)
     * @param cameraId 相机唯一标识
     * @return H264数据,队列为空则阻塞直到有数据/线程中断
     * @throws InterruptedException 线程中断异常
     */
    public static byte[] takeDataFromQueue(String cameraId) throws InterruptedException {
        BlockingQueue<byte[]> queue = CAMERA_QUEUE_MAP.get(cameraId);
        return queue != null ? queue.take() : null;
    }
 
    /**
     * 从指定相机的队列中获取H264数据(非阻塞,队列为空则返回null)
     * @param cameraId 相机唯一标识
     * @return H264数据,队列为空/队列不存在则返回null
     */
    public static byte[] pollDataFromQueue(String cameraId) {
        BlockingQueue<byte[]> queue = CAMERA_QUEUE_MAP.get(cameraId);
        return queue != null ? queue.poll() : null;
    }
 
    /**
     * 销毁指定相机的队列(清空数据并从容器中移除)
     * @param cameraId 相机唯一标识
     * @return true=销毁成功,false=队列不存在
     */
    public static boolean destroyCameraQueue(String cameraId) {
        if (cameraId == null || cameraId.trim().isEmpty()) {
            return false;
        }
        BlockingQueue<byte[]> queue = CAMERA_QUEUE_MAP.remove(cameraId.trim());
        if (queue != null) {
            queue.clear();
            return true;
        }
        return false;
    }
 
    // ===================== 静态监控/辅助方法 =====================
 
    /**
     * 获取指定相机队列的当前数据量
     * @param cameraId 相机唯一标识
     * @return 队列大小,队列不存在则返回0
     */
    public static int getQueueSize(String cameraId) {
        return Optional.ofNullable(CAMERA_QUEUE_MAP.get(cameraId)).map(BlockingQueue::size).orElse(0);
    }
 
    /**
     * 检查指定相机的队列是否存在
     * @param cameraId 相机唯一标识
     * @return true=存在,false=不存在
     */
    public static boolean isQueueExists(String cameraId) {
        return CAMERA_QUEUE_MAP.containsKey(cameraId);
    }
 
    /**
     * 获取所有相机队列的信息(用于监控/日志)
     * @return 队列信息Map:Key=相机ID,Value=队列大小
     */
    public static Map<String, Integer> getAllCameraQueueInfo() {
        ConcurrentHashMap<String, Integer> queueInfoMap = new ConcurrentHashMap<>();
        CAMERA_QUEUE_MAP.forEach((cameraId, queue) -> queueInfoMap.put(cameraId, queue.size()));
        return queueInfoMap;
    }
 
    /**
     * 清空所有相机队列(慎用,一般用于系统关闭/重置)
     */
    public static void clearAllQueues() {
        CAMERA_QUEUE_MAP.forEach((cameraId, queue) -> queue.clear());
    }
}