Java 程序可以通过 RocketMQ 提供的 API 来查询消息堆积情况,并根据堆积程度来决定是否停止发送消息。这是一种非常重要的自我保护和生产端限流策略。


方法一:使用 RocketMQ Admin API(推荐)

这是最常用和官方推荐的方式。你需要引入 rocketmq-admin 依赖来获取集群和主题的统计信息。

1. 添加依赖
<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-admin</artifactId>
    <version>4.9.7</version> <!-- 请使用与你的Broker相同的版本 -->
</dependency>

<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-tools</artifactId>
    <version>4.9.7</version>
</dependency>
2. 代码实现:查询堆积并控制发送
import org.apache.rocketmq.admin.MQAdminExt;
import org.apache.rocketmq.admin.impl.MQAdminExtImpl;
import org.apache.rocketmq.common.protocol.body.ConsumerConnection;
import org.apache.rocketmq.common.protocol.body.ConsumerRunningInfo;
import org.apache.rocketmq.common.protocol.heartbeat.ConsumptionType;
import org.apache.rocketmq.common.protocol.heartbeat.MessageModel;
import org.apache.rocketmq.remoting.RPCHook;
import org.apache.rocketmq.remoting.protocol.RemotingCommand;

import java.util.Properties;
import java.util.Set;
import java.util.concurrent.atomic.AtomicBoolean;

public class SmartProducer {

    private final DefaultMQProducer producer;
    private final MQAdminExt mqAdminExt;
    private final String consumerGroupName; // 你要监控的消费者组名
    private final String topicName;
    private final long maxAccumulationThreshold; // 消息堆积阈值,例如 50000 条

    // 全局开关,控制是否允许发送
    private final AtomicBoolean sendingAllowed = new AtomicBoolean(true);

    public SmartProducer(String producerGroup, String namesrvAddr, String consumerGroupName, String topicName, long maxThreshold) throws Exception {
        this.producer = new DefaultMQProducer(producerGroup);
        this.producer.setNamesrvAddr(namesrvAddr);
        this.consumerGroupName = consumerGroupName;
        this.topicName = topicName;
        this.maxAccumulationThreshold = maxThreshold;

        // 初始化 Admin API
        this.mqAdminExt = new MQAdminExtImpl();
        this.mqAdminExt.start();

        // 启动时检查一次堆积情况
        checkAccumulationAndControl();
        
        // 定期检查堆积情况(例如每10秒一次)
        startAccumulationMonitor();
    }

    /**
     * 检查消息堆积并控制发送开关
     */
    private void checkAccumulationAndControl() {
        try {
            long totalAccumulation = getMessageAccumulationCount();
            System.out.println("当前消息堆积量: " + totalAccumulation);

            if (totalAccumulation >= maxAccumulationThreshold) {
                if (sendingAllowed.compareAndSet(true, false)) {
                    System.err.println("警告:消息堆积已达到 " + totalAccumulation + ",超过阈值 " + maxAccumulationThreshold + ",停止发送消息!");
                }
            } else {
                if (sendingAllowed.compareAndSet(false, true)) {
                    System.out.println("消息堆积已回落至 " + totalAccumulation + ",恢复消息发送。");
                }
            }
        } catch (Exception e) {
            // 查询失败时,出于安全考虑,可以停止发送或记录日志
            System.err.println("查询消息堆积失败: " + e.getMessage());
            // sendingAllowed.set(false); // 根据你的安全策略决定
        }
    }

    /**
     * 获取指定消费者组对指定主题的消息堆积量
     */
    private long getMessageAccumulationCount() throws Exception {
        ClusterInfo clusterInfo = mqAdminExt.examineBrokerClusterInfo();
        Set<String> brokerAddrSet = clusterInfo.getBrokerAddrTable().keySet();
        
        long totalAccumulation = 0L;

        for (String brokerAddr : brokerAddrSet) {
            try {
                // 获取消费进度
                ConsumeStats consumeStats = mqAdminExt.examineConsumeStats(consumerGroupName);
                // 从消费进度中提取该主题的堆积量
                TopicStatsTable topicStatsTable = consumeStats.getOffsetTable();
                for (Map.Entry<MessageQueue, OffsetWrapper> entry : topicStatsTable.entrySet()) {
                    MessageQueue mq = entry.getKey();
                    if (mq.getTopic().equals(topicName)) {
                        OffsetWrapper offsetWrapper = entry.getValue();
                        long brokerOffset = offsetWrapper.getBrokerOffset(); // Broker最大Offset
                        long consumerOffset = offsetWrapper.getConsumerOffset(); // 消费者消费到的Offset
                        long accumulation = brokerOffset - consumerOffset;
                        if (accumulation > 0) {
                            totalAccumulation += accumulation;
                        }
                    }
                }
            } catch (Exception e) {
                // 某个Broker查询失败,继续查询下一个
                System.err.println("查询Broker " + brokerAddr + " 的堆积信息失败: " + e.getMessage());
            }
        }
        return totalAccumulation;
    }

    /**
     * 启动定时监控任务
     */
    private void startAccumulationMonitor() {
        ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
        scheduler.scheduleAtFixedRate(this::checkAccumulationAndControl, 10, 10, TimeUnit.SECONDS); // 每10秒检查一次
    }

    /**
     * 安全的发送消息方法
     */
    public SendResult sendSafely(Message msg) throws Exception {
        if (!sendingAllowed.get()) {
            throw new RuntimeException("消息堆积过多,暂停发送。当前堆积量可能超过 " + maxAccumulationThreshold);
        }
        return producer.send(msg);
    }

    public void start() throws Exception {
        producer.start();
    }

    public void shutdown() {
        producer.shutdown();
        mqAdminExt.shutdown();
    }

    public static void main(String[] args) throws Exception {
        SmartProducer smartProducer = new SmartProducer(
            "SmartProducerGroup",
            "localhost:9876",
            "YourConsumerGroupName", // 你要监控的消费者组
            "YourTopicName",
            50000 // 阈值:50000条
        );
        
        smartProducer.start();

        // 示例发送
        for (int i = 0; i < 1000000; i++) {
            try {
                Message msg = new Message("YourTopicName", "TagA", ("Hello " + i).getBytes());
                SendResult result = smartProducer.sendSafely(msg);
                System.out.println("发送成功: " + result.getMsgId());
                Thread.sleep(1000); // 模拟生产速度
            } catch (RuntimeException e) {
                System.err.println("发送被拒绝: " + e.getMessage());
                Thread.sleep(10000); // 被拒绝后等待10秒再试
            }
        }

        smartProducer.shutdown();
    }
}

方法二:通过 HTTP 接口查询(替代方案)

如果不想引入 rocketmq-admin 依赖,RocketMQ 的 Broker 和 Console 也提供了 HTTP 接口可以查询堆积情况。

你可以使用 Spring 的 RestTemplate 或 OkHttp 来调用这些接口:

// 示例:通过 HTTP 查询堆积量
public long getAccumulationViaHttp(String consumerGroup, String topic) {
    try {
        // 这里需要你的 RocketMQ Console 地址或 Broker 的 HTTP 端点
        String url = "http://your-console-server:8080/consumer/consumptionProgress.query?consumerGroup=" + consumerGroup;
        
        RestTemplate restTemplate = new RestTemplate();
        ResponseEntity<String> response = restTemplate.getForEntity(url, String.class);
        
        // 解析返回的JSON,找到对应topic的diffTotal字段(即堆积量)
        // 这里需要根据实际的JSON结构进行解析
        // 假设返回格式中有 "diffTotal" 字段
        ObjectMapper mapper = new ObjectMapper();
        JsonNode root = mapper.readTree(response.getBody());
        long accumulation = root.path("data").path("diffTotal").asLong();
        
        return accumulation;
    } catch (Exception e) {
        throw new RuntimeException("通过HTTP查询堆积量失败", e);
    }
}

重要考虑因素

  1. 监控粒度:上面的示例监控的是整个主题的堆积,你还可以细化到具体 MessageQueue。

  2. 性能影响:频繁调用 Admin API 会对 Broker 产生压力,建议监控间隔不要太短(10-30秒比较合适)。

  3. 消费者组选择:如果一个主题有多个消费者组,你需要监控正确的那个。

  4. 策略灵活性

    • 可以设置多级阈值(如警告阈值、停止阈值)
    • 可以动态调整发送速率而不是完全停止
    • 可以结合其他指标(如消费速率、延迟时间等)
  5. 异常处理:在查询堆积信息失败时,要有降级策略(比如继续发送、停止发送或告警)。

这种"智能生产者"模式在需要保证系统稳定性的场景中非常有用,可以有效防止因消费端故障导致的消息无限堆积,避免磁盘被打满等严重问题。

Logo

魔乐社区(Modelers.cn) 是一个中立、公益的人工智能社区,提供人工智能工具、模型、数据的托管、展示与应用协同服务,为人工智能开发及爱好者搭建开放的学习交流平台。社区通过理事会方式运作,由全产业链共同建设、共同运营、共同享有,推动国产AI生态繁荣发展。

更多推荐