java程序产生RocketMQ消息,怎么查询到RocketMQ消息堆积
·
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);
}
}
重要考虑因素
-
监控粒度:上面的示例监控的是整个主题的堆积,你还可以细化到具体 MessageQueue。
-
性能影响:频繁调用 Admin API 会对 Broker 产生压力,建议监控间隔不要太短(10-30秒比较合适)。
-
消费者组选择:如果一个主题有多个消费者组,你需要监控正确的那个。
-
策略灵活性:
- 可以设置多级阈值(如警告阈值、停止阈值)
- 可以动态调整发送速率而不是完全停止
- 可以结合其他指标(如消费速率、延迟时间等)
-
异常处理:在查询堆积信息失败时,要有降级策略(比如继续发送、停止发送或告警)。
这种"智能生产者"模式在需要保证系统稳定性的场景中非常有用,可以有效防止因消费端故障导致的消息无限堆积,避免磁盘被打满等严重问题。
魔乐社区(Modelers.cn) 是一个中立、公益的人工智能社区,提供人工智能工具、模型、数据的托管、展示与应用协同服务,为人工智能开发及爱好者搭建开放的学习交流平台。社区通过理事会方式运作,由全产业链共同建设、共同运营、共同享有,推动国产AI生态繁荣发展。
更多推荐


所有评论(0)