同一机器多JVM进程Chronicle Queues Pub/Sub配置及消息排序服务咨询
基于Chronicle Queues实现多JVM进程的Pub/Sub及消息排序方案
关于现成服务的说明
Chronicle生态中没有直接提供这种带中心排序节点的开箱即用服务,但可以基于Chronicle Queue的核心能力快速搭建符合需求的Pub/Sub架构。
用Chronicle Queue实现需求的具体方案
你的场景可以通过"多发布者→中心排序节点→多订阅者"的三层架构实现,以下是具体步骤:
1. 架构设计
- 多发布者JVM:各自向共享的输入Chronicle Queue写入消息(可选择单队列或多队列,根据业务隔离需求决定)
- 中心排序节点JVM:作为独立进程,监听所有输入队列,拉取消息后按指定规则排序,再将排序后的消息写入输出Chronicle Queue
- 多订阅者JVM:从共享的输出队列消费排序完成的消息
2. 核心代码实现示例
发布者端(每个JVM进程)
// 初始化输入队列,同一机器使用本地文件路径实现进程共享 ChronicleQueue inputQueue = ChronicleQueue.singleBuilder("/opt/chronicle/input-queue").build(); ExcerptAppender appender = inputQueue.acquireAppender(); // 发送带排序键的业务消息 try (DocumentContext dc = appender.writingDocument()) { dc.wire().write("sortKey").int32((int) System.currentTimeMillis()); // 用时间戳作为排序依据 dc.wire().write("payload").text("发布者消息内容"); }
中心排序节点端
// 初始化输入、输出队列 ChronicleQueue inputQueue = ChronicleQueue.singleBuilder("/opt/chronicle/input-queue").build(); ChronicleQueue outputQueue = ChronicleQueue.singleBuilder("/opt/chronicle/output-queue").build(); ExcerptTailer tailer = inputQueue.createTailer("sort-node-tailer"); ExcerptAppender outputAppender = outputQueue.acquireAppender(); // 用TreeMap实现按sortKey自动排序的缓存 TreeMap<Integer, String> sortedCache = new TreeMap<>(); // 轮询拉取消息并排序输出 while (true) { try (DocumentContext dc = tailer.readingDocument()) { if (dc.isPresent()) { int sortKey = dc.wire().read("sortKey").int32(); String payload = dc.wire().read("payload").text(); sortedCache.put(sortKey, payload); } } // 达到批量阈值时写入输出队列(示例:批量大小100) if (sortedCache.size() >= 100) { for (Map.Entry<Integer, String> entry : sortedCache.entrySet()) { try (DocumentContext dc = outputAppender.writingDocument()) { dc.wire().write("sortKey").int32(entry.getKey()); dc.wire().write("payload").text(entry.getValue()); } } sortedCache.clear(); } Thread.sleep(10); // 避免空轮询占用CPU }
订阅者端(每个JVM进程)
// 初始化输出队列 ChronicleQueue outputQueue = ChronicleQueue.singleBuilder("/opt/chronicle/output-queue").build(); // 每个订阅者用唯一ID标记消费位置,避免重复消费 ExcerptTailer tailer = outputQueue.createTailer("subscriber-1"); // 持续消费排序后的消息 while (true) { try (DocumentContext dc = tailer.readingDocument()) { if (dc.isPresent()) { int sortKey = dc.wire().read("sortKey").int32(); String payload = dc.wire().read("payload").text(); // 处理业务逻辑 System.out.printf("订阅者收到排序后消息:sortKey=%d, payload=%s%n", sortKey, payload); } } Thread.sleep(10); }
3. 关键注意事项
- 队列共享机制:Chronicle Queue通过内存映射文件实现进程间共享,同一机器上的JVM只需指向相同本地文件路径即可
- 排序逻辑扩展:可根据业务需求替换排序规则(如业务ID、优先级等),若中心节点采用多线程拉取消息,需保证排序缓存的线程安全
- 性能优化:调整队列的
rollCycle参数适配消息量,采用批量读写减少IO开销;中心节点可使用多线程拉取输入队列提升吞吐量 - 持久化保障:Chronicle Queue的消息默认持久化到磁盘,进程重启后不会丢失未消费的消息
内容的提问来源于stack exchange,提问作者Androiduser14919 Starters
相关产品推荐
相关产品推荐

