This is an automated email from the ASF dual-hosted git repository.
lizhimins pushed a commit to branch rocketmq-studio
in repository https://gitbox.apache.org/repos/asf/rocketmq-dashboard.git
The following commit(s) were added to refs/heads/rocketmq-studio by this push:
new 9d46e0dd feat: support HA replication lag alert rules in alert module
(#658)
9d46e0dd is described below
commit 9d46e0dd0fd799a691dfff68efa8104a872ecf40
Author: zhaohai <[email protected]>
AuthorDate: Fri Jul 31 15:50:51 2026 +0800
feat: support HA replication lag alert rules in alert module (#658)
Rework of the original branch-compensate proposal: instead of adding a
parallel per-metric CRUD stack, extend the existing alert module so
replication lag rules reuse /api/alert-rules:
- Add brokerName/clusterName scope and severity fields to AlertRuleVO
- Render broker/cluster label selectors in exported Prometheus expressions
- Use rule severity (critical/warning/info, default warning) in export
- Route replication/slave metrics to the broker rule group
---
.../rocketmq/studio/ops/alert/AlertRuleVO.java | 6 ++
.../rocketmq/studio/ops/alert/AlertService.java | 35 ++++++++++-
.../studio/ops/alert/AlertServiceTest.java | 67 ++++++++++++++++++++++
3 files changed, 106 insertions(+), 2 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRuleVO.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRuleVO.java
index ec467237..2b8557a2 100644
--- a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRuleVO.java
+++ b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRuleVO.java
@@ -39,4 +39,10 @@ public class AlertRuleVO {
private boolean enabled;
private String lastTriggered;
private String description;
+ /** Target broker name pattern (e.g. "broker-a", or "*" for all) */
+ private String brokerName;
+ /** Target cluster name pattern (e.g. "DefaultCluster", or "*" for all) */
+ private String clusterName;
+ /** Alert severity: critical, warning, info */
+ private String severity;
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertService.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertService.java
index 5fb04726..59df2845 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertService.java
@@ -168,7 +168,7 @@ public class AlertService {
alertName(rule),
expression(rule),
duration(rule),
- "warning",
+ severity(rule),
team,
summary(rule),
description(rule));
@@ -194,7 +194,35 @@ public class AlertService {
private String expression(AlertRuleVO rule) {
String metric = hasText(rule.getMetric()) ? rule.getMetric() :
"rocketmq_consumer_lag_messages";
String operator = hasText(rule.getOperator()) ? rule.getOperator() :
">";
- return metric + " " + operator + " " +
formatThreshold(rule.getThreshold());
+ return metric + labelSelector(rule) + " " + operator + " " +
formatThreshold(rule.getThreshold());
+ }
+
+ private String labelSelector(AlertRuleVO rule) {
+ StringBuilder selector = new StringBuilder();
+ appendLabel(selector, "cluster", rule.getClusterName());
+ appendLabel(selector, "broker", rule.getBrokerName());
+ return selector.isEmpty() ? "" : "{" + selector + "}";
+ }
+
+ private void appendLabel(StringBuilder selector, String label, String
value) {
+ if (!hasText(value) || "*".equals(value.trim())) {
+ return;
+ }
+ if (!selector.isEmpty()) {
+ selector.append(',');
+ }
+ selector.append(label).append("=\"").append(value.trim()).append('"');
+ }
+
+ private String severity(AlertRuleVO rule) {
+ String severity = rule.getSeverity();
+ if (hasText(severity)) {
+ String normalized = severity.trim().toLowerCase();
+ if ("critical".equals(normalized) || "warning".equals(normalized)
|| "info".equals(normalized)) {
+ return normalized;
+ }
+ }
+ return "warning";
}
private String formatThreshold(double threshold) {
@@ -212,6 +240,9 @@ public class AlertService {
if (!hasText(metric)) {
return "broker";
}
+ if (metric.contains("replication") || metric.contains("fall_behind")
|| metric.contains("slave")) {
+ return "broker";
+ }
if (metric.contains("consumer") || metric.contains("lag")) {
return "consumer";
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertServiceTest.java
index 318d34f2..4c110c51 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertServiceTest.java
@@ -108,6 +108,73 @@ class AlertServiceTest {
.contains("description: \"Lag too high\"");
}
+ @Test
+ void
exportPrometheusRulesYamlShouldRenderReplicationLagRuleWithScopeAndSeverity() {
+ AlertRuleVO rule = AlertRuleVO.builder()
+ .name("Replication Lag High")
+ .metric("rocketmq_broker_replication_lag_bytes")
+ .operator(">")
+ .threshold(104857600)
+ .duration("5m")
+ .brokerName("broker-a")
+ .clusterName("DefaultCluster")
+ .severity("critical")
+ .description("Slave falls behind master")
+ .build();
+ when(alertRepository.findAllRules()).thenReturn(List.of(rule));
+
+ String result = alertService.exportPrometheusRulesYaml();
+
+ assertThat(result)
+ .contains("rocketmq-broker.rules")
+ .contains("expr:
rocketmq_broker_replication_lag_bytes{cluster=\"DefaultCluster\",broker=\"broker-a\"}
> 104857600")
+ .contains("for: 5m")
+ .contains("severity: critical");
+ }
+
+ @Test
+ void
exportPrometheusRulesYamlShouldIgnoreWildcardScopeAndInvalidSeverity() {
+ AlertRuleVO rule = AlertRuleVO.builder()
+ .name("Replication Lag Any Broker")
+ .metric("rocketmq_broker_replication_lag_bytes")
+ .operator(">")
+ .threshold(1024)
+ .brokerName("*")
+ .clusterName(" ")
+ .severity("fatal")
+ .build();
+ when(alertRepository.findAllRules()).thenReturn(List.of(rule));
+
+ String result = alertService.exportPrometheusRulesYaml();
+
+ assertThat(result)
+ .contains("expr: rocketmq_broker_replication_lag_bytes > 1024")
+ .contains("severity: warning");
+ }
+
+ @Test
+ void createRuleShouldPreserveReplicationScopeFields() {
+ AlertRuleVO input = AlertRuleVO.builder()
+ .name("Replication Lag High")
+ .metric("rocketmq_broker_replication_lag_bytes")
+ .operator(">")
+ .threshold(104857600)
+ .thresholdUnit("B")
+ .brokerName("broker-a")
+ .clusterName("DefaultCluster")
+ .severity("critical")
+ .build();
+
when(alertRepository.saveRule(any(AlertRuleVO.class))).thenAnswer(invocation ->
invocation.getArgument(0));
+
+ AlertRuleVO result = alertService.createRule(input);
+
+ assertThat(result.getId()).isNotNull().isNotEmpty();
+ assertThat(result.getBrokerName()).isEqualTo("broker-a");
+ assertThat(result.getClusterName()).isEqualTo("DefaultCluster");
+ assertThat(result.getSeverity()).isEqualTo("critical");
+ verify(alertRepository).saveRule(result);
+ }
+
@Test
void createRuleShouldAssignId() {
AlertRuleVO input = AlertRuleVO.builder().name("New
Rule").metric("tps")