elasticsearch: collect metrics from _cat/indices (#6965)
* get metrics from _cat/indices
Ilya Mashchenko committed
Oct 3, 2019 at 19:11 UTC
f000b41f1088328ebb5cd07f15cc9ef1276eb641
3 files changed
+209
-94
collectors/python.d.plugin/elasticsearch/README.md
+15
-6
@@ -1,6 +1,6 @@
1
# elasticsearch
2
3
-This module monitors Elasticsearch performance and health metrics.
3
+This module monitors [Elasticsearch](https://www.elastic.co/products/elasticsearch) performance and health metrics.
4
5
It produces:
6
@@ -51,19 +51,28 @@ It produces:
51
- Store statistics
52
- Indices and shards statistics
53
54
+9. **Indices** charts (per index statistics, disabled by default):
55
+
56
+ - Docs count
57
+ - Store size
58
+ - Num of replicas
59
+ - Health status
60
+
61
## configuration
62
63
Sample:
64
65
```yaml
66
local:
60
- host : 'ipaddress' # Elasticsearch server ip address or hostname
61
- port : 'port' # Port on which elasticsearch listens
62
- cluster_health : True/False # Calls to cluster health elasticsearch API. Enabled by default.
63
- cluster_stats : True/False # Calls to cluster stats elasticsearch API. Enabled by default.
67
+ host : 'ipaddress' # Elasticsearch server ip address or hostname.
68
+ port : 'port' # Port on which elasticsearch listens.
69
+ node_status : yes/no # Get metrics from "/_nodes/_local/stats". Enabled by default.
70
+ cluster_health : yes/no # Get metrics from "/_cluster/health". Enabled by default.
71
+ cluster_stats : yes/no # Get metrics from "'/_cluster/stats". Enabled by default.
72
+ indices_stats : yes/no # Get metrics from "/_cat/indices". Disabled by default.
73
```
74
66
-If no configuration is given, module will fail to run.
75
+If no configuration is given, module will try to connect to `http://127.0.0.1:9200`.
76
77
---
78
collectors/python.d.plugin/elasticsearch/elasticsearch.chart.py
+188
-83
@@ -168,6 +168,10 @@ ORDER = [
168
'cluster_stats_store',
169
'cluster_stats_indices',
170
'cluster_stats_shards_total',
171
+ 'index_docs_count',
172
+ 'index_store_size',
173
+ 'index_replica',
174
+ 'index_health',
175
]
176
177
CHARTS = {
@@ -196,7 +200,8 @@ CHARTS = {
200
]
201
},
202
'search_latency': {
199
- 'options': [None, 'Query And Fetch Latency', 'milliseconds', 'search performance', 'elastic.search_latency', 'stacked'],
203
+ 'options': [None, 'Query And Fetch Latency', 'milliseconds', 'search performance', 'elastic.search_latency',
204
+ 'stacked'],
205
'lines': [
206
['query_latency', 'query', 'absolute', 1, 1000],
207
['fetch_latency', 'fetch', 'absolute', 1, 1000]
@@ -397,9 +402,6 @@ CHARTS = {
402
'lines': [
403
['status_green', 'green', 'absolute'],
404
['status_red', 'red', 'absolute'],
400
- ['status_foo1', None, 'absolute'],
401
- ['status_foo2', None, 'absolute'],
402
- ['status_foo3', None, 'absolute'],
405
['status_yellow', 'yellow', 'absolute']
406
]
407
},
@@ -483,10 +485,61 @@ CHARTS = {
485
'lines': [
486
['http_current_open', 'opened', 'absolute', 1, 1]
487
]
486
- }
488
+ },
489
+ 'index_docs_count': {
490
+ 'options': [None, 'Docs Count', 'count', 'indices', 'elastic.index_docs', 'line'],
491
+ 'lines': []
492
+ },
493
+ 'index_store_size': {
494
+ 'options': [None, 'Store Size', 'bytes', 'indices', 'elastic.index_store_size', 'line'],
495
+ 'lines': []
496
+ },
497
+ 'index_replica': {
498
+ 'options': [None, 'Replica', 'count', 'indices', 'elastic.index_replica', 'line'],
499
+ 'lines': []
500
+ },
501
+ 'index_health': {
502
+ 'options': [None, 'Health', 'status', 'indices', 'elastic.index_health', 'line'],
503
+ 'lines': []
504
+ },
505
}
506
507
508
+def convert_index_store_size_to_bytes(size):
509
+ # can be b, kb, mb, gb
510
+ if size.endswith('kb'):
511
+ return round(float(size[:-2]) * 1024)
512
+ elif size.endswith('mb'):
513
+ return round(float(size[:-2]) * 1024 * 1024)
514
+ elif size.endswith('gb'):
515
+ return round(float(size[:-2]) * 1024 * 1024 * 1024)
516
+ elif size.endswith('b'):
517
+ return round(float(size[:-1]))
518
+ return -1
519
+
520
+
521
+def convert_index_health(health):
522
+ if health == 'green':
523
+ return 0
524
+ elif health == 'yellow':
525
+ return 1
526
+ elif health == 'read':
527
+ return 2
528
+ return -1
529
+
530
+
531
+def get_survive_any(method):
532
+ def w(*args):
533
+ try:
534
+ method(*args)
535
+ except Exception as error:
536
+ self, queue, url = args[0], args[1], args[2]
537
+ self.error("error during '{0}' : {1}".format(url, error))
538
+ queue.put(dict())
539
+
540
+ return w
541
+
542
+
543
class Service(UrlService):
544
def __init__(self, configuration=None, name=None):
545
UrlService.__init__(self, configuration=configuration, name=name)
@@ -501,34 +554,41 @@ class Service(UrlService):
554
)
555
self.latency = dict()
556
self.methods = list()
557
+ self.collected_indices = set()
558
559
def check(self):
506
- if not all([self.host,
507
- self.port,
508
- isinstance(self.host, str),
509
- isinstance(self.port, (str, int))]):
560
+ if not self.host:
561
self.error('Host is not defined in the module configuration file')
562
return False
563
513
- # Hostname -> ip address
564
try:
565
self.host = gethostbyname(self.host)
566
except gaierror as error:
517
- self.error(str(error))
567
+ self.error(repr(error))
568
return False
569
520
- # Create URL for every Elasticsearch API
521
- self.methods = [METHODS(get_data=self._get_node_stats,
522
- url=self.url + '/_nodes/_local/stats',
523
- run=self.configuration.get('node_stats', True)),
524
- METHODS(get_data=self._get_cluster_health,
525
- url=self.url + '/_cluster/health',
526
- run=self.configuration.get('cluster_health', True)),
527
- METHODS(get_data=self._get_cluster_stats,
528
- url=self.url + '/_cluster/stats',
529
- run=self.configuration.get('cluster_stats', True))]
530
-
531
- # Remove disabled API calls from 'avail methods'
570
+ self.methods = [
571
+ METHODS(
572
+ get_data=self._get_node_stats,
573
+ url=self.url + '/_nodes/_local/stats',
574
+ run=self.configuration.get('node_stats', True),
575
+ ),
576
+ METHODS(
577
+ get_data=self._get_cluster_health,
578
+ url=self.url + '/_cluster/health',
579
+ run=self.configuration.get('cluster_health', True)
580
+ ),
581
+ METHODS(
582
+ get_data=self._get_cluster_stats,
583
+ url=self.url + '/_cluster/stats',
584
+ run=self.configuration.get('cluster_stats', True),
585
+ ),
586
+ METHODS(
587
+ get_data=self._get_indices,
588
+ url=self.url + '/_cat/indices?format=json',
589
+ run=self.configuration.get('indices_stats', False),
590
+ ),
591
+ ]
592
return UrlService.check(self)
593
594
def _get_data(self):
@@ -539,8 +599,11 @@ class Service(UrlService):
599
for method in self.methods:
600
if not method.run:
601
continue
542
- th = threading.Thread(target=method.get_data,
543
- args=(queue, method.url))
602
+ th = threading.Thread(
603
+ target=method.get_data,
604
+ args=(queue, method.url),
605
+ )
606
+ th.daemon = True
607
th.start()
608
threads.append(th)
609
@@ -550,88 +613,128 @@ class Service(UrlService):
613
614
return result or None
615
553
- def _get_cluster_health(self, queue, url):
554
- """
555
- Format data received from http request
556
- :return: dict
557
- """
558
-
616
+ def add_index_to_charts(self, idx_name):
617
+ for name in ('index_docs_count', 'index_store_size', 'index_replica', 'index_health'):
618
+ chart = self.charts[name]
619
+ dim = ['{0}_{1}'.format(idx_name, name), idx_name]
620
+ chart.add_dimension(dim)
621
+
622
+ @get_survive_any
623
+ def _get_indices(self, queue, url):
624
+ # [
625
+ # {
626
+ # "pri.store.size": "650b",
627
+ # "health": "yellow",
628
+ # "status": "open",
629
+ # "index": "twitter",
630
+ # "pri": "5",
631
+ # "rep": "1",
632
+ # "docs.count": "10",
633
+ # "docs.deleted": "3",
634
+ # "store.size": "650b"
635
+ # }
636
+ # ]
637
raw_data = self._get_raw_data(url)
560
-
638
if not raw_data:
639
return queue.put(dict())
640
564
- data = self.json_reply(raw_data)
565
-
566
- if not data:
641
+ indices = self.json_parse(raw_data)
642
+ if not indices:
643
return queue.put(dict())
644
569
- to_netdata = fetch_data_(raw_data=data,
570
- metrics=HEALTH_STATS)
645
+ charts_initialized = len(self.charts) != 0
646
+ data = dict()
647
+ for idx in indices:
648
+ try:
649
+ name = idx['index']
650
+ is_system_index = name.startswith('.')
651
+ if is_system_index:
652
+ continue
653
+
654
+ v = {
655
+ '{0}_index_docs_count'.format(name): idx['docs.count'],
656
+ '{0}_index_replica'.format(name): idx['rep'],
657
+ '{0}_index_health'.format(name): convert_index_health(idx['health']),
658
+ }
659
+ size = convert_index_store_size_to_bytes(idx['store.size'])
660
+ if size != -1:
661
+ v['{0}_index_store_size'.format(name)] = size
662
+ except KeyError as error:
663
+ self.debug("error on parsing index : {0}".format(repr(error)))
664
+ continue
665
572
- to_netdata.update({'status_green': 0, 'status_red': 0, 'status_yellow': 0,
573
- 'status_foo1': 0, 'status_foo2': 0, 'status_foo3': 0})
574
- current_status = 'status_' + data['status']
575
- to_netdata[current_status] = 1
666
+ data.update(v)
667
+ if name not in self.collected_indices and charts_initialized:
668
+ self.collected_indices.add(name)
669
+ self.add_index_to_charts(name)
670
577
- return queue.put(to_netdata)
671
+ return queue.put(data)
672
579
- def _get_cluster_stats(self, queue, url):
580
- """
581
- Format data received from http request
582
- :return: dict
583
- """
584
-
585
- raw_data = self._get_raw_data(url)
673
+ @get_survive_any
674
+ def _get_cluster_health(self, queue, url):
675
+ raw = self._get_raw_data(url)
676
+ if not raw:
677
+ return queue.put(dict())
678
587
- if not raw_data:
679
+ parsed = self.json_parse(raw)
680
+ if not parsed:
681
return queue.put(dict())
682
590
- data = self.json_reply(raw_data)
683
+ data = fetch_data(raw_data=parsed, metrics=HEALTH_STATS)
684
+ dummy = {
685
+ 'status_green': 0,
686
+ 'status_red': 0,
687
+ 'status_yellow': 0,
688
+ }
689
+ data.update(dummy)
690
+ current_status = 'status_' + parsed['status']
691
+ data[current_status] = 1
692
592
- if not data:
593
- return queue.put(dict())
693
+ return queue.put(data)
694
595
- to_netdata = fetch_data_(raw_data=data,
596
- metrics=CLUSTER_STATS)
695
+ @get_survive_any
696
+ def _get_cluster_stats(self, queue, url):
697
+ raw = self._get_raw_data(url)
698
+ if not raw:
699
+ return queue.put(dict())
700
598
- return queue.put(to_netdata)
701
+ parsed = self.json_parse(raw)
702
+ if not parsed:
703
+ return queue.put(dict())
704
600
- def _get_node_stats(self, queue, url):
601
- """
602
- Format data received from http request
603
- :return: dict
604
- """
705
+ data = fetch_data(raw_data=parsed, metrics=CLUSTER_STATS)
706
606
- raw_data = self._get_raw_data(url)
707
+ return queue.put(data)
708
608
- if not raw_data:
709
+ @get_survive_any
710
+ def _get_node_stats(self, queue, url):
711
+ raw = self._get_raw_data(url)
712
+ if not raw:
713
return queue.put(dict())
714
611
- data = self.json_reply(raw_data)
612
-
613
- if not data:
715
+ parsed = self.json_parse(raw)
716
+ if not parsed:
717
return queue.put(dict())
718
616
- node = list(data['nodes'].keys())[0]
617
- to_netdata = fetch_data_(raw_data=data['nodes'][node],
618
- metrics=NODE_STATS)
719
+ node = list(parsed['nodes'].keys())[0]
720
+ data = fetch_data(raw_data=parsed['nodes'][node], metrics=NODE_STATS)
721
722
# Search, index, flush, fetch performance latency
723
for key in LATENCY:
724
try:
623
- to_netdata[key] = self.find_avg(total=to_netdata[LATENCY[key]['total']],
624
- spent_time=to_netdata[LATENCY[key]['spent_time']],
625
- key=key)
725
+ data[key] = self.find_avg(
726
+ total=data[LATENCY[key]['total']],
727
+ spent_time=data[LATENCY[key]['spent_time']],
728
+ key=key)
729
except KeyError:
730
continue
628
- if 'process_open_file_descriptors' in to_netdata and 'process_max_file_descriptors' in to_netdata:
629
- to_netdata['file_descriptors_used'] = round(float(to_netdata['process_open_file_descriptors'])
630
- / to_netdata['process_max_file_descriptors'] * 1000)
731
+ if 'process_open_file_descriptors' in data and 'process_max_file_descriptors' in data:
732
+ v = float(data['process_open_file_descriptors']) / data['process_max_file_descriptors'] * 1000
733
+ data['file_descriptors_used'] = round(v)
734
632
- return queue.put(to_netdata)
735
+ return queue.put(data)
736
634
- def json_reply(self, reply):
737
+ def json_parse(self, reply):
738
try:
739
return json.loads(reply)
740
except ValueError as err:
@@ -640,20 +743,22 @@ class Service(UrlService):
743
744
def find_avg(self, total, spent_time, key):
745
if key not in self.latency:
643
- self.latency[key] = dict(total=total,
644
- spent_time=spent_time)
746
+ self.latency[key] = dict(total=total, spent_time=spent_time)
747
return 0
748
+
749
if self.latency[key]['total'] != total:
647
- latency = float(spent_time - self.latency[key]['spent_time'])\
648
- / float(total - self.latency[key]['total']) * 1000
750
+ spent_diff = spent_time - self.latency[key]['spent_time']
751
+ total_diff = total - self.latency[key]['total']
752
+ latency = float(spent_diff) / float(total_diff) * 1000
753
self.latency[key]['total'] = total
754
self.latency[key]['spent_time'] = spent_time
755
return latency
756
+
757
self.latency[key]['spent_time'] = spent_time
758
return 0
759
760
656
-def fetch_data_(raw_data, metrics):
761
+def fetch_data(raw_data, metrics):
762
data = dict()
763
for metric in metrics:
764
value = raw_data
@@ -661,7 +766,7 @@ def fetch_data_(raw_data, metrics):
766
try:
767
for m in metrics_list:
768
value = value[m]
664
- except KeyError:
769
+ except (KeyError, TypeError):
770
continue
771
data['_'.join(metrics_list)] = value
772
return data
collectors/python.d.plugin/elasticsearch/elasticsearch.conf
+6
-5
@@ -61,11 +61,12 @@
61
#
62
# Additionally to the above, elasticsearch plugin also supports the following:
63
#
64
-# host: 'ipaddress' # Server ip address or hostname.
65
-# port: 'port' # Port on which elasticsearch listen.
66
-# scheme: 'scheme' # URL scheme. Default is 'http'.
67
-# cluster_health: False/True # Calls to cluster health elasticsearch API. Enabled by default.
68
-# cluster_stats: False/True # Calls to cluster stats elasticsearch API. Enabled by default.
64
+# host : 'ipaddress' # Elasticsearch server ip address or hostname.
65
+# port : 'port' # Port on which elasticsearch listens.
66
+# node_status : yes/no # Get metrics from "/_nodes/_local/stats". Enabled by default.
67
+# cluster_health : yes/no # Get metrics from "/_cluster/health". Enabled by default.
68
+# cluster_stats : yes/no # Get metrics from "'/_cluster/stats". Enabled by default.
69
+# indices_stats : yes/no # Get metrics from "/_cat/indices". Disabled by default.
70
#
71
#
72
# if the URL is password protected, the following are supported: