master
go 404 lines 11.6 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package cassandra
4
5 import (
6 "errors"
7 "strings"
8
9 "github.com/netdata/netdata/go/plugins/pkg/prometheus"
10 )
11
12 const (
13 suffixCount = "_count"
14 suffixValue = "_value"
15 )
16
17 func (c *Collector) collect() (map[string]int64, error) {
18 pms, err := c.prom.ScrapeSeries()
19 if err != nil {
20 return nil, err
21 }
22
23 if c.validateMetrics {
24 if !isCassandraMetrics(pms) {
25 return nil, errors.New("collected metrics aren't Collector metrics")
26 }
27 c.validateMetrics = false
28 }
29
30 mx := make(map[string]int64)
31
32 c.resetMetrics()
33 c.collectMetrics(pms)
34 c.processMetric(mx)
35
36 return mx, nil
37 }
38
39 func (c *Collector) resetMetrics() {
40 cm := newCassandraMetrics()
41 for key, p := range c.mx.threadPools {
42 cm.threadPools[key] = &threadPoolMetrics{
43 name: p.name,
44 hasCharts: p.hasCharts,
45 }
46 }
47 c.mx = cm
48 }
49
50 func (c *Collector) processMetric(mx map[string]int64) {
51 c.mx.clientReqTotalLatencyReads.write(mx, "client_request_total_latency_reads")
52 c.mx.clientReqTotalLatencyWrites.write(mx, "client_request_total_latency_writes")
53 c.mx.clientReqLatencyReads.write(mx, "client_request_latency_reads")
54 c.mx.clientReqLatencyWrites.write(mx, "client_request_latency_writes")
55 c.mx.clientReqTimeoutsReads.write(mx, "client_request_timeouts_reads")
56 c.mx.clientReqTimeoutsWrites.write(mx, "client_request_timeouts_writes")
57 c.mx.clientReqUnavailablesReads.write(mx, "client_request_unavailables_reads")
58 c.mx.clientReqUnavailablesWrites.write(mx, "client_request_unavailables_writes")
59 c.mx.clientReqFailuresReads.write(mx, "client_request_failures_reads")
60 c.mx.clientReqFailuresWrites.write(mx, "client_request_failures_writes")
61
62 c.mx.clientReqReadLatencyP50.write(mx, "client_request_read_latency_p50")
63 c.mx.clientReqReadLatencyP75.write(mx, "client_request_read_latency_p75")
64 c.mx.clientReqReadLatencyP95.write(mx, "client_request_read_latency_p95")
65 c.mx.clientReqReadLatencyP98.write(mx, "client_request_read_latency_p98")
66 c.mx.clientReqReadLatencyP99.write(mx, "client_request_read_latency_p99")
67 c.mx.clientReqReadLatencyP999.write(mx, "client_request_read_latency_p999")
68 c.mx.clientReqWriteLatencyP50.write(mx, "client_request_write_latency_p50")
69 c.mx.clientReqWriteLatencyP75.write(mx, "client_request_write_latency_p75")
70 c.mx.clientReqWriteLatencyP95.write(mx, "client_request_write_latency_p95")
71 c.mx.clientReqWriteLatencyP98.write(mx, "client_request_write_latency_p98")
72 c.mx.clientReqWriteLatencyP99.write(mx, "client_request_write_latency_p99")
73 c.mx.clientReqWriteLatencyP999.write(mx, "client_request_write_latency_p999")
74
75 c.mx.rowCacheHits.write(mx, "row_cache_hits")
76 c.mx.rowCacheMisses.write(mx, "row_cache_misses")
77 c.mx.rowCacheSize.write(mx, "row_cache_size")
78 if c.mx.rowCacheHits.isSet && c.mx.rowCacheMisses.isSet {
79 if s := c.mx.rowCacheHits.value + c.mx.rowCacheMisses.value; s > 0 {
80 mx["row_cache_hit_ratio"] = int64((c.mx.rowCacheHits.value * 100 / s) * 1000)
81 } else {
82 mx["row_cache_hit_ratio"] = 0
83 }
84 }
85 if c.mx.rowCacheCapacity.isSet && c.mx.rowCacheSize.isSet {
86 if s := c.mx.rowCacheCapacity.value; s > 0 {
87 mx["row_cache_utilization"] = int64((c.mx.rowCacheSize.value * 100 / s) * 1000)
88 } else {
89 mx["row_cache_utilization"] = 0
90 }
91 }
92
93 c.mx.keyCacheHits.write(mx, "key_cache_hits")
94 c.mx.keyCacheMisses.write(mx, "key_cache_misses")
95 c.mx.keyCacheSize.write(mx, "key_cache_size")
96 if c.mx.keyCacheHits.isSet && c.mx.keyCacheMisses.isSet {
97 if s := c.mx.keyCacheHits.value + c.mx.keyCacheMisses.value; s > 0 {
98 mx["key_cache_hit_ratio"] = int64((c.mx.keyCacheHits.value * 100 / s) * 1000)
99 } else {
100 mx["key_cache_hit_ratio"] = 0
101 }
102 }
103 if c.mx.keyCacheCapacity.isSet && c.mx.keyCacheSize.isSet {
104 if s := c.mx.keyCacheCapacity.value; s > 0 {
105 mx["key_cache_utilization"] = int64((c.mx.keyCacheSize.value * 100 / s) * 1000)
106 } else {
107 mx["key_cache_utilization"] = 0
108 }
109 }
110
111 c.mx.droppedMessages.write1k(mx, "dropped_messages")
112
113 c.mx.storageLoad.write(mx, "storage_load")
114 c.mx.storageExceptions.write(mx, "storage_exceptions")
115
116 c.mx.compactionBytesCompacted.write(mx, "compaction_bytes_compacted")
117 c.mx.compactionPendingTasks.write(mx, "compaction_pending_tasks")
118 c.mx.compactionCompletedTasks.write(mx, "compaction_completed_tasks")
119
120 c.mx.jvmMemoryHeapUsed.write(mx, "jvm_memory_heap_used")
121 c.mx.jvmMemoryNonHeapUsed.write(mx, "jvm_memory_nonheap_used")
122 c.mx.jvmGCParNewCount.write(mx, "jvm_gc_parnew_count")
123 c.mx.jvmGCParNewTime.write1k(mx, "jvm_gc_parnew_time")
124 c.mx.jvmGCCMSCount.write(mx, "jvm_gc_cms_count")
125 c.mx.jvmGCCMSTime.write1k(mx, "jvm_gc_cms_time")
126
127 for _, p := range c.mx.threadPools {
128 if !p.hasCharts {
129 p.hasCharts = true
130 c.addThreadPoolCharts(p)
131 }
132
133 px := "thread_pool_" + p.name + "_"
134 p.activeTasks.write(mx, px+"active_tasks")
135 p.pendingTasks.write(mx, px+"pending_tasks")
136 p.blockedTasks.write(mx, px+"blocked_tasks")
137 p.totalBlockedTasks.write(mx, px+"total_blocked_tasks")
138 }
139 }
140
141 func (c *Collector) collectMetrics(pms prometheus.Series) {
142 c.collectClientRequestMetrics(pms)
143 c.collectDroppedMessagesMetrics(pms)
144 c.collectThreadPoolsMetrics(pms)
145 c.collectStorageMetrics(pms)
146 c.collectCacheMetrics(pms)
147 c.collectJVMMetrics(pms)
148 c.collectCompactionMetrics(pms)
149 }
150
151 func (c *Collector) collectClientRequestMetrics(pms prometheus.Series) {
152 const metric = "org_apache_cassandra_metrics_clientrequest"
153
154 var rw struct{ read, write *metricValue }
155 for _, pm := range pms.FindByName(metric + suffixCount) {
156 name := pm.Labels.Get("name")
157 scope := pm.Labels.Get("scope")
158
159 switch name {
160 case "TotalLatency":
161 rw.read, rw.write = &c.mx.clientReqTotalLatencyReads, &c.mx.clientReqTotalLatencyWrites
162 case "Latency":
163 rw.read, rw.write = &c.mx.clientReqLatencyReads, &c.mx.clientReqLatencyWrites
164 case "Timeouts":
165 rw.read, rw.write = &c.mx.clientReqTimeoutsReads, &c.mx.clientReqTimeoutsWrites
166 case "Unavailables":
167 rw.read, rw.write = &c.mx.clientReqUnavailablesReads, &c.mx.clientReqUnavailablesWrites
168 case "Failures":
169 rw.read, rw.write = &c.mx.clientReqFailuresReads, &c.mx.clientReqFailuresWrites
170 default:
171 continue
172 }
173
174 switch scope {
175 case "Read":
176 rw.read.add(pm.Value)
177 case "Write":
178 rw.write.add(pm.Value)
179 }
180 }
181
182 rw = struct{ read, write *metricValue }{}
183
184 for _, pm := range pms.FindByNames(
185 metric+"_50thpercentile",
186 metric+"_75thpercentile",
187 metric+"_95thpercentile",
188 metric+"_98thpercentile",
189 metric+"_99thpercentile",
190 metric+"_999thpercentile",
191 ) {
192 name := pm.Labels.Get("name")
193 scope := pm.Labels.Get("scope")
194
195 if name != "Latency" {
196 continue
197 }
198
199 switch {
200 case strings.HasSuffix(pm.Name(), "_50thpercentile"):
201 rw.read, rw.write = &c.mx.clientReqReadLatencyP50, &c.mx.clientReqWriteLatencyP50
202 case strings.HasSuffix(pm.Name(), "_75thpercentile"):
203 rw.read, rw.write = &c.mx.clientReqReadLatencyP75, &c.mx.clientReqWriteLatencyP75
204 case strings.HasSuffix(pm.Name(), "_95thpercentile"):
205 rw.read, rw.write = &c.mx.clientReqReadLatencyP95, &c.mx.clientReqWriteLatencyP95
206 case strings.HasSuffix(pm.Name(), "_98thpercentile"):
207 rw.read, rw.write = &c.mx.clientReqReadLatencyP98, &c.mx.clientReqWriteLatencyP98
208 case strings.HasSuffix(pm.Name(), "_99thpercentile"):
209 rw.read, rw.write = &c.mx.clientReqReadLatencyP99, &c.mx.clientReqWriteLatencyP99
210 case strings.HasSuffix(pm.Name(), "_999thpercentile"):
211 rw.read, rw.write = &c.mx.clientReqReadLatencyP999, &c.mx.clientReqWriteLatencyP999
212 default:
213 continue
214 }
215
216 switch scope {
217 case "Read":
218 rw.read.add(pm.Value)
219 case "Write":
220 rw.write.add(pm.Value)
221 }
222 }
223 }
224
225 func (c *Collector) collectCacheMetrics(pms prometheus.Series) {
226 const metric = "org_apache_cassandra_metrics_cache"
227
228 var hm struct{ hits, misses *metricValue }
229 for _, pm := range pms.FindByName(metric + suffixCount) {
230 name := pm.Labels.Get("name")
231 scope := pm.Labels.Get("scope")
232
233 switch scope {
234 case "KeyCache":
235 hm.hits, hm.misses = &c.mx.keyCacheHits, &c.mx.keyCacheMisses
236 case "RowCache":
237 hm.hits, hm.misses = &c.mx.rowCacheHits, &c.mx.rowCacheMisses
238 default:
239 continue
240 }
241
242 switch name {
243 case "Hits":
244 hm.hits.add(pm.Value)
245 case "Misses":
246 hm.misses.add(pm.Value)
247 }
248 }
249
250 var cs struct{ cap, size *metricValue }
251 for _, pm := range pms.FindByName(metric + suffixValue) {
252 name := pm.Labels.Get("name")
253 scope := pm.Labels.Get("scope")
254
255 switch scope {
256 case "KeyCache":
257 cs.cap, cs.size = &c.mx.keyCacheCapacity, &c.mx.keyCacheSize
258 case "RowCache":
259 cs.cap, cs.size = &c.mx.rowCacheCapacity, &c.mx.rowCacheSize
260 default:
261 continue
262 }
263
264 switch name {
265 case "Capacity":
266 cs.cap.add(pm.Value)
267 case "Size":
268 cs.size.add(pm.Value)
269 }
270 }
271 }
272
273 func (c *Collector) collectThreadPoolsMetrics(pms prometheus.Series) {
274 const metric = "org_apache_cassandra_metrics_threadpools"
275
276 for _, pm := range pms.FindByName(metric + suffixValue) {
277 name := pm.Labels.Get("name")
278 scope := pm.Labels.Get("scope")
279 pool := c.getThreadPoolMetrics(scope)
280
281 switch name {
282 case "ActiveTasks":
283 pool.activeTasks.add(pm.Value)
284 case "PendingTasks":
285 pool.pendingTasks.add(pm.Value)
286 }
287 }
288 for _, pm := range pms.FindByName(metric + suffixCount) {
289 name := pm.Labels.Get("name")
290 scope := pm.Labels.Get("scope")
291 pool := c.getThreadPoolMetrics(scope)
292
293 switch name {
294 case "CompletedTasks":
295 pool.totalBlockedTasks.add(pm.Value)
296 case "TotalBlockedTasks":
297 pool.totalBlockedTasks.add(pm.Value)
298 case "CurrentlyBlockedTasks":
299 pool.blockedTasks.add(pm.Value)
300 }
301 }
302 }
303
304 func (c *Collector) collectStorageMetrics(pms prometheus.Series) {
305 const metric = "org_apache_cassandra_metrics_storage"
306
307 for _, pm := range pms.FindByName(metric + suffixCount) {
308 name := pm.Labels.Get("name")
309
310 switch name {
311 case "Load":
312 c.mx.storageLoad.add(pm.Value)
313 case "Exceptions":
314 c.mx.storageExceptions.add(pm.Value)
315 }
316 }
317 }
318
319 func (c *Collector) collectDroppedMessagesMetrics(pms prometheus.Series) {
320 const metric = "org_apache_cassandra_metrics_droppedmessage"
321
322 for _, pm := range pms.FindByName(metric + suffixCount) {
323 c.mx.droppedMessages.add(pm.Value)
324 }
325 }
326
327 func (c *Collector) collectJVMMetrics(pms prometheus.Series) {
328 const metricMemUsed = "jvm_memory_bytes_used"
329 const metricGC = "jvm_gc_collection_seconds"
330
331 for _, pm := range pms.FindByName(metricMemUsed) {
332 area := pm.Labels.Get("area")
333
334 switch area {
335 case "heap":
336 c.mx.jvmMemoryHeapUsed.add(pm.Value)
337 case "nonheap":
338 c.mx.jvmMemoryNonHeapUsed.add(pm.Value)
339 }
340 }
341
342 for _, pm := range pms.FindByName(metricGC + suffixCount) {
343 gc := pm.Labels.Get("gc")
344
345 switch gc {
346 case "ParNew":
347 c.mx.jvmGCParNewCount.add(pm.Value)
348 case "ConcurrentMarkSweep":
349 c.mx.jvmGCCMSCount.add(pm.Value)
350 }
351 }
352
353 for _, pm := range pms.FindByName(metricGC + "_sum") {
354 gc := pm.Labels.Get("gc")
355
356 switch gc {
357 case "ParNew":
358 c.mx.jvmGCParNewTime.add(pm.Value)
359 case "ConcurrentMarkSweep":
360 c.mx.jvmGCCMSTime.add(pm.Value)
361 }
362 }
363 }
364
365 func (c *Collector) collectCompactionMetrics(pms prometheus.Series) {
366 const metric = "org_apache_cassandra_metrics_compaction"
367
368 for _, pm := range pms.FindByName(metric + suffixValue) {
369 name := pm.Labels.Get("name")
370
371 switch name {
372 case "CompletedTasks":
373 c.mx.compactionCompletedTasks.add(pm.Value)
374 case "PendingTasks":
375 c.mx.compactionPendingTasks.add(pm.Value)
376 }
377 }
378 for _, pm := range pms.FindByName(metric + suffixCount) {
379 name := pm.Labels.Get("name")
380
381 switch name {
382 case "BytesCompacted":
383 c.mx.compactionBytesCompacted.add(pm.Value)
384 }
385 }
386 }
387
388 func (c *Collector) getThreadPoolMetrics(name string) *threadPoolMetrics {
389 pool, ok := c.mx.threadPools[name]
390 if !ok {
391 pool = &threadPoolMetrics{name: name}
392 c.mx.threadPools[name] = pool
393 }
394 return pool
395 }
396
397 func isCassandraMetrics(pms prometheus.Series) bool {
398 for _, pm := range pms {
399 if strings.HasPrefix(pm.Name(), "org_apache_cassandra_metrics") {
400 return true
401 }
402 }
403 return false
404 }