Skip to content

Commit 6993714

Browse files
hai-benhaiben-amzn
authored andcommitted
Added output streaming support and new metrics.
1 parent 17481d4 commit 6993714

7 files changed

Lines changed: 340 additions & 85 deletions

File tree

activemq-prometheus/README.md

Lines changed: 15 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -40,35 +40,39 @@ Two endpoints because brokers with many destinations might produce large respons
4040

4141
| Metric | Type | Description |
4242
|--------|------|-------------|
43-
| `connections_count` | gauge | Current number of connections |
43+
| `connections` | gauge | Current number of connections |
4444
| `connections_total` | counter | Total connections since last start |
4545
| `messages_enqueued_total` | counter | Total messages enqueued since last start |
4646
| `messages_dequeued_total` | counter | Total messages dequeued since last start |
47-
| `consumers_count` | gauge | Current number of consumers |
48-
| `producers_count` | gauge | Current number of producers |
49-
| `message_count` | gauge | Current number of messages across all destinations |
47+
| `consumers` | gauge | Current number of consumers |
48+
| `producers` | gauge | Current number of producers |
49+
| `messages` | gauge | Current number of messages across all destinations |
5050
| `memory_percent_usage` | gauge | Percent of memory limit used |
5151
| `memory_limit_bytes` | gauge | Memory limit in bytes |
5252
| `store_percent_usage` | gauge | Percent of store limit used |
5353
| `store_limit_bytes` | gauge | Store limit in bytes |
5454
| `temp_percent_usage` | gauge | Percent of temp limit used |
5555
| `temp_limit_bytes` | gauge | Temp limit in bytes |
5656
| `uptime_milliseconds` | gauge | Broker uptime in milliseconds |
57+
| `queues` | gauge | Number of queues on the broker |
58+
| `topics` | gauge | Number of topics on the broker |
59+
| `job_scheduler_store_percent_usage` | gauge | Percent of job scheduler store limit used |
60+
| `job_scheduler_store_limit_bytes` | gauge | Job scheduler store limit in bytes |
5761

5862
### Destination metrics (`activemq_queue_*` / `activemq_topic_*`)
5963

6064
Returned only when `?per_object=true` is set.
6165

6266
| Metric | Type | Description |
6367
|--------|------|-------------|
64-
| `message_count` | gauge | Number of messages in destination |
65-
| `enqueue_count_total` | counter | Total messages enqueued since last start |
66-
| `dequeue_count_total` | counter | Total messages dequeued since last start |
67-
| `dispatch_count_total` | counter | Total messages dispatched since last start |
68+
| `messages` | gauge | Number of messages in destination |
69+
| `enqueued_total` | counter | Total messages enqueued since last start |
70+
| `dequeued_total` | counter | Total messages dequeued since last start |
71+
| `dispatched_total` | counter | Total messages dispatched since last start |
6872
| `message_inflight_count` | gauge | Messages dispatched but not acknowledged |
69-
| `expired_count_total` | counter | Total messages expired since last start |
70-
| `consumer_count` | gauge | Number of consumers |
71-
| `producer_count` | gauge | Number of producers |
73+
| `expired_total` | counter | Total messages expired since last start |
74+
| `consumers` | gauge | Number of consumers |
75+
| `producers` | gauge | Number of producers |
7276
| `memory_percent_usage` | gauge | Percent of destination memory limit used |
7377
| `memory_limit_bytes` | gauge | Memory limit for destination in bytes |
7478
| `memory_usage_bytes` | gauge | Memory used by destination in bytes |

activemq-prometheus/pom.xml

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,11 @@
4141
<scope>provided</scope>
4242
</dependency>
4343

44+
<dependency>
45+
<groupId>org.slf4j</groupId>
46+
<artifactId>slf4j-api</artifactId>
47+
</dependency>
48+
4449
<!-- =============================== -->
4550
<!-- Testing Dependencies -->
4651
<!-- =============================== -->

activemq-prometheus/src/main/java/org/apache/activemq/prometheus/PrometheusMetricsServlet.java

Lines changed: 99 additions & 51 deletions
Original file line numberDiff line numberDiff line change
@@ -18,8 +18,10 @@
1818

1919
import java.io.IOException;
2020
import java.io.PrintWriter;
21-
import java.io.StringWriter;
2221
import java.lang.management.ManagementFactory;
22+
import java.util.Collections;
23+
import java.util.LinkedHashMap;
24+
import java.util.Map;
2325
import java.util.Set;
2426

2527
import jakarta.servlet.http.HttpServlet;
@@ -28,9 +30,13 @@
2830
import javax.management.MBeanServer;
2931
import javax.management.ObjectName;
3032

33+
import org.slf4j.Logger;
34+
import org.slf4j.LoggerFactory;
35+
3136
public class PrometheusMetricsServlet extends HttpServlet {
3237

3338
private static final long serialVersionUID = 1L;
39+
private static final Logger LOG = LoggerFactory.getLogger(PrometheusMetricsServlet.class);
3440
private static final String CONTENT_TYPE = "text/plain; version=0.0.4; charset=utf-8";
3541

3642
// Metrics can be easily extended by adding them here
@@ -48,7 +54,11 @@ public class PrometheusMetricsServlet extends HttpServlet {
4854
new MetricDefinition("store_limit_bytes", "Store limit in bytes", "StoreLimit", MetricType.GAUGE),
4955
new MetricDefinition("temp_percent_usage", "Percent (0-100) of temp limit used", "TempPercentUsage", MetricType.GAUGE),
5056
new MetricDefinition("temp_limit_bytes", "Temp limit in bytes", "TempLimit", MetricType.GAUGE),
51-
new MetricDefinition("uptime_milliseconds", "Broker uptime in milliseconds", "UptimeMillis", MetricType.GAUGE)
57+
new MetricDefinition("uptime_milliseconds", "Broker uptime in milliseconds", "UptimeMillis", MetricType.GAUGE),
58+
new MetricDefinition("queues", "Number of queues on the broker", "TotalQueuesCount", MetricType.GAUGE),
59+
new MetricDefinition("topics", "Number of topics on the broker", "TotalTopicsCount", MetricType.GAUGE),
60+
new MetricDefinition("job_scheduler_store_percent_usage", "Percent (0-100) of job scheduler store limit used", "JobSchedulerStorePercentUsage", MetricType.GAUGE),
61+
new MetricDefinition("job_scheduler_store_limit_bytes", "Job scheduler store limit in bytes", "JobSchedulerStoreLimit", MetricType.GAUGE)
5262
};
5363

5464
private static final MetricDefinition[] DESTINATION_METRICS = {
@@ -68,96 +78,134 @@ public class PrometheusMetricsServlet extends HttpServlet {
6878
};
6979

7080
@Override
71-
protected void doGet(HttpServletRequest request, HttpServletResponse response) throws IOException {
72-
boolean perObject = request != null && "true".equalsIgnoreCase(request.getParameter("per_object"));
73-
74-
StringWriter output = new StringWriter();
75-
PrintWriter writer = new PrintWriter(output);
76-
81+
protected void doGet(final HttpServletRequest request, final HttpServletResponse response) throws IOException {
82+
final boolean perObject = request != null && "true".equalsIgnoreCase(request.getParameter("per_object"));
83+
final MBeanServer mBeanServer = ManagementFactory.getPlatformMBeanServer();
84+
85+
// If the mBeanServer is unavailable nothing useful can be produced, so return a 500
86+
final Set<ObjectName> brokers;
87+
final Set<ObjectName> queues;
88+
final Set<ObjectName> topics;
7789
try {
78-
writeMetrics(ManagementFactory.getPlatformMBeanServer(), writer, perObject);
79-
} catch (Exception exception) {
90+
brokers = mBeanServer.queryNames(new ObjectName("org.apache.activemq:type=Broker,brokerName=*"), null);
91+
if (perObject) {
92+
// Scraping destinations on brokers with many queues or topics can be expensive.
93+
queues = mBeanServer.queryNames(new ObjectName(
94+
"org.apache.activemq:type=Broker,brokerName=*,destinationType=Queue,destinationName=*"), null);
95+
topics = mBeanServer.queryNames(new ObjectName(
96+
"org.apache.activemq:type=Broker,brokerName=*,destinationType=Topic,destinationName=*"), null);
97+
} else {
98+
queues = Collections.emptySet();
99+
topics = Collections.emptySet();
100+
}
101+
} catch (final Exception exception) {
102+
LOG.warn("Prometheus scrape failed while querying broker MBeans", exception);
80103
response.sendError(HttpServletResponse.SC_INTERNAL_SERVER_ERROR, "Metrics collection failed");
81104
return;
82105
}
83106

107+
// Stream metrics to avoid keeping the response body in memory
108+
// Failed items are skipped or default to 0 (when not possible) to allow the scrape to continue
84109
response.setContentType(CONTENT_TYPE);
85110
response.setStatus(HttpServletResponse.SC_OK);
86-
response.getWriter().write(output.toString());
87-
}
88-
89-
void writeMetrics(MBeanServer mBeanServer, PrintWriter writer, boolean perObject) throws Exception {
90-
writeBrokerMetrics(mBeanServer, writer);
91-
// Scraping this by default on brokers with lots of queues or topics might be expensive
111+
final PrintWriter writer = response.getWriter();
112+
writeBrokerMetrics(mBeanServer, writer, brokers);
92113
if (perObject) {
93-
writeDestinationMetrics(mBeanServer, writer, "Queue");
94-
writeDestinationMetrics(mBeanServer, writer, "Topic");
114+
writeDestinationMetrics(mBeanServer, writer, "Queue", queues);
115+
writeDestinationMetrics(mBeanServer, writer, "Topic", topics);
95116
}
96117
writer.flush();
97118
}
98119

99-
private void writeBrokerMetrics(MBeanServer mBeanServer, PrintWriter writer) throws Exception {
100-
ObjectName pattern = new ObjectName("org.apache.activemq:type=Broker,brokerName=*");
101-
Set<ObjectName> brokers = mBeanServer.queryNames(pattern, null);
120+
private void writeBrokerMetrics(final MBeanServer mBeanServer, final PrintWriter writer, final Set<ObjectName> brokers) {
121+
// In case there is a network of brokers, only use the local one
122+
// Unreadable names are skipped
123+
final Map<ObjectName, String> identified = new LinkedHashMap<>();
124+
for (final ObjectName broker : brokers) {
125+
final String brokerName = resolveStringAttribute(mBeanServer, broker, "BrokerName");
126+
if (brokerName != null) {
127+
identified.put(broker, "broker=\"" + sanitizeLabel(brokerName) + "\"");
128+
}
129+
}
102130

103-
for (MetricDefinition metric : BROKER_METRICS) {
104-
String metricName = "activemq_broker_" + metric.name;
131+
for (final MetricDefinition metric : BROKER_METRICS) {
132+
final String metricName = "activemq_broker_" + metric.name;
105133
writeMetadata(writer, metricName, metric);
106-
for (ObjectName broker : brokers) {
107-
String brokerName = sanitizeLabel((String) mBeanServer.getAttribute(broker, "BrokerName"));
108-
String labels = "broker=\"" + brokerName + "\"";
109-
writeSample(writer, metricName, labels, getNumber(mBeanServer, broker, metric.attribute));
134+
for (final Map.Entry<ObjectName, String> entry : identified.entrySet()) {
135+
writeSample(writer, metricName, entry.getValue(), getNumber(mBeanServer, entry.getKey(), metric.attribute));
110136
}
111137
}
112138
}
113139

114-
private void writeDestinationMetrics(MBeanServer mBeanServer, PrintWriter writer, String type) throws Exception {
115-
String queryPattern = "org.apache.activemq:type=Broker,brokerName=*,destinationType=" + type + ",destinationName=*";
116-
Set<ObjectName> destinations = mBeanServer.queryNames(new ObjectName(queryPattern), null);
117-
String typeLower = type.toLowerCase();
140+
private void writeDestinationMetrics(final MBeanServer mBeanServer, final PrintWriter writer, final String type,
141+
final Set<ObjectName> destinations) {
142+
final String typeLower = type.toLowerCase();
118143

119-
for (MetricDefinition metric : DESTINATION_METRICS) {
120-
String metricName = "activemq_" + typeLower + "_" + metric.name;
144+
for (final MetricDefinition metric : DESTINATION_METRICS) {
145+
final String metricName = "activemq_" + typeLower + "_" + metric.name;
121146
writeMetadata(writer, metricName, metric.withFormattedHelp(typeLower));
122-
for (ObjectName destination : destinations) {
123-
String brokerName = sanitizeLabel(destination.getKeyProperty("brokerName"));
124-
String destinationName = sanitizeLabel(destination.getKeyProperty("destinationName"));
125-
String labels = String.format("broker=\"%s\",destination=\"%s\"", brokerName, destinationName);
147+
for (final ObjectName destination : destinations) {
148+
final String brokerName = sanitizeLabel(destination.getKeyProperty("brokerName"));
149+
final String destinationName = sanitizeLabel(destination.getKeyProperty("destinationName"));
150+
final String labels = String.format("broker=\"%s\",destination=\"%s\"", brokerName, destinationName);
126151
writeSample(writer, metricName, labels, getNumber(mBeanServer, destination, metric.attribute));
127152
}
128153
}
129154
}
130155

131-
private double getNumber(MBeanServer mBeanServer, ObjectName name, String attribute) {
156+
private String resolveStringAttribute(final MBeanServer mBeanServer, final ObjectName name, final String attribute) {
157+
// Partial results are better than no results if something goes wrong
158+
try {
159+
final Object value = mBeanServer.getAttribute(name, attribute);
160+
if (value instanceof String) {
161+
return (String) value;
162+
}
163+
} catch (final Exception exception) {
164+
LOG.debug("Skipping object {}: identity attribute {} unavailable", name, attribute, exception);
165+
}
166+
return null;
167+
}
168+
169+
private double getNumber(final MBeanServer mBeanServer, final ObjectName name, final String attribute) {
132170
try {
133-
Object value = mBeanServer.getAttribute(name, attribute);
171+
final Object value = mBeanServer.getAttribute(name, attribute);
134172
if (value instanceof Number) {
135173
return ((Number) value).doubleValue();
136174
}
137-
} catch (Exception ignored) {
138-
// Some attributes are not available on every ActiveMQ deployment (eg: bridge metrics)
175+
} catch (final Exception exception) {
176+
// Some attributes are not available on every ActiveMQ deployment (eg: bridge metrics).
177+
LOG.debug("Reporting 0 for {} on {}: attribute unavailable", attribute, name, exception);
139178
}
140179
return 0;
141180
}
142181

143-
private void writeMetadata(PrintWriter writer, String metric, MetricDefinition def) {
182+
private void writeMetadata(final PrintWriter writer, final String metric, final MetricDefinition def) {
144183
writer.println("# HELP " + metric + " " + def.help);
145184
writer.println("# TYPE " + metric + " " + def.type.prometheusName());
146185
}
147186

148-
private void writeSample(PrintWriter writer, String metric, String labels, double value) {
149-
if (value == (long) value) {
150-
writer.println(metric + "{" + labels + "} " + (long) value);
187+
private void writeSample(final PrintWriter writer, final String metric, final String labels, final double value) {
188+
final String rendered;
189+
if (Double.isNaN(value)) {
190+
rendered = "NaN";
191+
} else if (value == Double.POSITIVE_INFINITY) {
192+
rendered = "+Inf";
193+
} else if (value == Double.NEGATIVE_INFINITY) {
194+
rendered = "-Inf";
195+
} else if (value == (long) value) {
196+
rendered = Long.toString((long) value);
151197
} else {
152-
writer.println(metric + "{" + labels + "} " + value);
198+
rendered = Double.toString(value);
153199
}
200+
writer.println(metric + "{" + labels + "} " + rendered);
154201
}
155202

156-
static String sanitizeLabel(String value) {
203+
static String sanitizeLabel(final String value) {
157204
if (value == null) {
158205
return "unknown";
159206
}
160-
return value.replace("\\", "\\\\").replace("\"", "\\\"").replace("\n", "\\n");
207+
// See: https://prometheus.io/docs/instrumenting/exposition_formats/
208+
return value.replace("\\", "\\\\").replace("\"", "\\\"").replace("\n", "\\n").replace("\r", "\\r");
161209
}
162210

163211
private static final class MetricDefinition {
@@ -166,14 +214,14 @@ private static final class MetricDefinition {
166214
private final String attribute;
167215
private final MetricType type;
168216

169-
private MetricDefinition(String name, String help, String attribute, MetricType type) {
217+
private MetricDefinition(final String name, final String help, final String attribute, final MetricType type) {
170218
this.name = name;
171219
this.help = help;
172220
this.attribute = attribute;
173221
this.type = type;
174222
}
175223

176-
private MetricDefinition withFormattedHelp(String arg) {
224+
private MetricDefinition withFormattedHelp(final String arg) {
177225
return new MetricDefinition(name, String.format(help, arg), attribute, type);
178226
}
179227
}
@@ -184,7 +232,7 @@ enum MetricType {
184232

185233
private final String prometheusName;
186234

187-
MetricType(String prometheusName) {
235+
MetricType(final String prometheusName) {
188236
this.prometheusName = prometheusName;
189237
}
190238

0 commit comments

Comments
 (0)