simplememorymq基于arrayblockingqueue实现线程安全的内存消息队列,支持发送、接收与优雅关闭;适用于轻量级异步解耦,但无持久化、无重试、单消费者串行处理,不适用于强一致性场景。

用 Java 阻塞队列(BlockingQueue)手写一个简易内存消息队列,核心是封装线程安全的生产-消费模型,不依赖外部中间件,适合轻量级异步解耦或内部事件通知场景。关键不是造轮子,而是理解其边界与设计取舍。
选对阻塞队列实现类
Java 提供了多个线程安全的阻塞队列实现,根据使用场景选择:
-
ArrayBlockingQueue:有界、基于数组、公平/非公平可选,适合内存可控、需防 OOM 的场景; -
LinkedBlockingQueue:默认无界(实际容量为Integer.MAX_VALUE),基于链表,吞吐高,但要注意内存失控风险; -
SynchronousQueue:不存储元素,纯“交接”队列,适合高实时性、零缓冲的直传场景(如任务调度桥接)。
一般入门推荐 ArrayBlockingQueue,显式控制容量,避免内存无限增长。
封装基础消息队列类
定义一个泛型消息队列,支持发送(offer)、接收(poll/take)、关闭(优雅停用):
public class SimpleMemoryMQ<t> {
private final BlockingQueue<t> queue;
private final Thread consumerThread;
private volatile boolean running = true;
public SimpleMemoryMQ(int capacity) {
this.queue = new ArrayBlockingQueue(capacity);
this.consumerThread = new Thread(this::consumeLoop, "SimpleMQ-Consumer");
this.consumerThread.setDaemon(false); // 非守护线程,确保能处理完
}
public void start() {
this.consumerThread.start();
}
public boolean send(T message) {
return queue.offer(message); // 非阻塞发送,失败返回 false
}
// 可选:提供阻塞发送(带超时)
public boolean send(T message, long timeout, TimeUnit unit) throws InterruptedException {
return queue.offer(message, timeout, unit);
}
private void consumeLoop() {
while (running || !queue.isEmpty()) {
try {
T msg = queue.poll(100, TimeUnit.MILLISECONDS); // 短等待避免空转
if (msg != null) {
processMessage(msg);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
}
}
}
// 子类或匿名方式实现具体业务逻辑
protected void processMessage(T message) {
System.out.println("Received: " + message);
}
public void shutdown() {
this.running = false;
this.consumerThread.interrupt();
try {
this.consumerThread.join(3000); // 最多等 3 秒
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
}</t></t>
说明:
- 使用
poll(timeout)而非take(),避免关机时永久阻塞; -
running标志 +queue.isEmpty()组合,确保未消费完的消息也被处理; -
processMessage设计为protected,方便子类定制逻辑(如日志、转发、重试)。
简单使用示例
发几条字符串消息,观察消费效果:
public class MQDemo {
public static void main(String[] args) throws InterruptedException {
SimpleMemoryMQ<string> mq = new SimpleMemoryMQ(10);
mq.start();
// 发送
mq.send("Hello");
mq.send("World");
mq.send("From SimpleMQ");
// 主线程稍作等待
Thread.sleep(500);
// 关闭
mq.shutdown();
}
}</string>
输出类似:
Received: HelloReceived: World
Received: From SimpleMQ
注意事项和简化边界
这个简易版只解决「内存内、单 JVM、无持久化、无 ACK、无重试」的基础通信:
- 消息丢失风险:JVM 崩溃即丢,不适用于金融、订单等强一致性场景;
- 无消费者分组/广播:所有消费逻辑都在同一个线程里串行执行;
- 无监控与统计:没暴露积压量、吞吐、延迟等指标;
- 异常需自行捕获:
processMessage中抛异常会导致当前消息跳过,后续消息继续——如需重试,得自己加 try-catch + 重入逻辑。
它本质是一个带消费线程的包装器,价值在于快速验证逻辑、教学演示或内部模块间低耦合通信,不是 RocketMQ 或 Kafka 的替代品。
Java免费学习笔记:立即使用
解锁 Java 大师之旅:从入门到精通的终极指南











