You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

同一机器多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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.22 18:22:54