Java NIO不直接提供分布式消息队列功能,但可构建单机内存型简易MQ服务端,基于Selector管理多连接、ConcurrentHashMap维护主题与订阅者关系,支持SUB/PUB文本协议及广播消息。

Java NIO 本身不直接提供分布式消息队列功能,但它可以作为高性能网络通信层的基础,支撑自研轻量级消息服务端(如支持发布/订阅、内存队列、简单持久化与多客户端连接)。下面是一个基于 java.nio(非 Netty)实现的**单机、内存型、支持多客户端连接的简易消息队列服务端原型**,具备基本的 publish/subscribe 功能,适用于学习和小规模场景验证。
核心设计思路
用 Selector 管理多个客户端连接,每个连接对应一个 SocketChannel;所有消息在内存中用 ConcurrentHashMap<String, CopyOnWriteArrayList<Consumer>> 模拟主题(topic)与订阅者关系;生产者发送消息时广播给所有该 topic 的活跃消费者(通过 channel 写回);不依赖外部存储,不处理集群同步或消息确认,仅体现 NIO 主干逻辑。
关键组件与步骤
1. 初始化 Selector 和 ServerSocketChannel
- 设置
ServerSocketChannel为非阻塞,绑定端口,注册OP_ACCEPT - 创建
Selector,启动主循环调用select()
2. 处理新连接(OP_ACCEPT)
立即学习“Java免费学习笔记(深入)”;
- 调用
serverChannel.accept()获取SocketChannel - 设为非阻塞,注册
OP_READ,并关联一个简单的会话对象(如ClientSession) - 可选:用
SelectionKey.attach()存储 session,避免额外映射表
3. 处理读请求(OP_READ)
Java项目代码review工具。分析Git变更+完整调用链路上下文,推断业务需求,进行多维度评分和分类汇总,生成完整PRD文档。包含细粒度Java代码审查清单(Null安全、异常处理、Streams、并发、equals/hashCode、资源管理、API设计、性能、MyBatis/ORM、事务边界、SQL/DD...
- 分配
ByteBuffer(建议用allocateDirect()减少 GC) - 调用
channel.read(buf),检查返回值:-1 表示断连,需取消 key 并关闭 channel - 解析协议:本例用简单文本格式,如
PUB topic1 hello world或SUB topic2 - 将命令分发到处理器(如
handlePublish()/handleSubscribe())
4. 处理写响应(OP_WRITE,按需注册)
- 写操作通常不长期注册
OP_WRITE(易 busy-loop),只在写缓冲区满、write()返回 0 时临时注册 - 成功写完后,取消
OP_WRITE,恢复监听OP_READ - 广播消息时,对每个订阅者的
SocketChannel尝试写入,失败则断开该 client
简化协议与内存模型示例
协议定义(每行一条命令,UTF-8 编码):
-
SUB <topic>—— 客户端订阅主题 -
UNSUB <topic>—— 取消订阅 -
PUB <topic> <message>—— 发布消息(空格分隔,message 不含换行) -
LIST—— 列出当前所有 topic
内存结构:
-
Map<String, Set<SocketChannel>> topicSubscribers:主题 → 订阅者 channel 集合 -
Map<SocketChannel, Set<String>> clientSubscriptions:客户端 → 其订阅的主题集合(用于 UNSUB 和断连清理)
注意事项与局限
这不是生产级 MQ,缺少:
- 消息持久化(崩溃即丢)
- 消费者确认与重投机制
- 集群节点间数据同步(ZooKeeper / Raft 等)
- 流量控制、背压、连接认证、TLS 加密
- 消息序列化(本例仅字符串)
若需真正分布式能力,建议基于 Kafka、RocketMQ 或使用 Netty + Redis/ZooKeeper 构建;NIO 更适合作为底层通信模块嵌入其中,而非从零造轮子。
不复杂但容易忽略:记得在每次 read 后调用 buf.flip() 和 buf.clear(),写之前 buf.flip(),写完后 buf.compact()(若未写完)——这是 NIO ByteBuffers 正确使用的前提。

















