用阻塞队列结合字节流实现单机生产者-消费者模型,核心是将字节数据作为产品通过队列传递:生产者读字节块入队,消费者取块写出;队列承担线程安全缓冲,字节流专注io,二者通过byte[]桥接,需处理结束标记、中断恢复与资源释放。
用阻塞队列结合字节流实现单机环境下的简单生产者-消费者模型,核心在于把字节数据作为“产品”在生产者与消费者之间传递,而阻塞队列天然承担缓冲区角色——它自动控制容量、线程安全、空满阻塞,无需手动加锁或 wait/notify。
明确字节流与阻塞队列的分工
字节流(如 InputStream / OutputStream)负责原始数据读写;阻塞队列(如 ArrayBlockingQueue
- 生产者从输入流读取一段字节(例如每次读 1024 字节),封装为
byte[]或ByteBuffer,调用queue.put()入队 - 消费者从队列取出该对象,写入输出流(如
outputStream.write(byteArr)),再释放引用(避免内存泄漏) - 队列大小即为内存中最大缓存字节数(例如
new ArrayBlockingQueue(10)表示最多缓存 10 个字节块)
选对阻塞队列类型和元素粒度
缓冲区控制效果直接受队列实现和单次传输单元影响:
- 固定容量优先选 ArrayBlockingQueue:底层数组 + 可重入锁,内存占用可控,适合明确限制总缓存上限(如最多缓存 5MB 数据)
-
按块而非按字节入队:不要把单个
byte放进队列(开销大、无意义),应以合理 chunk 大小(如 4KB、8KB)读取后整体入队 -
推荐用 byte[] 而非 ByteBuffer:避免堆外内存管理复杂性;若需复用缓冲区,可用
ByteBuffer.allocate()并 clear(),但须确保消费者处理完才被重用
处理边界与资源释放
字节流场景下容易忽略流关闭、异常中断和队列清空问题:
- 生产者读到
-1(流末尾)后,应发送一个特殊结束标记(如null或自定义 sentinel 对象),通知消费者停止 - 消费者检测到该标记后,不再
take(),并主动关闭输出流 - 若使用
put()/take(),需捕获InterruptedException并恢复中断状态(Thread.currentThread().interrupt()),防止线程中断被吞 - 避免用
offer()/poll()配合循环忙等——这会绕过阻塞机制,失去缓冲区节流能力
一个轻量可运行的示意结构
不依赖外部框架,仅用 JDK 标准类:
- 构造:
BlockingQueue<byte> queue = new ArrayBlockingQueue(8); // 最多存 8 个 4KB 块 ≈ 32KB 缓冲</byte> - 生产者:用
InputStream.read(byte[], 0, len)读块 → 封装新数组(避免复用导致脏读)→queue.put(chunk) - 消费者:
byte[] data = queue.take(); outputStream.write(data);→ 遇null则 break 并 close - 启动:两个线程分别执行生产者/消费者逻辑,主线程
join()等待完成











