rabbitmq_plugin: initial version added
Ilya committed
May 13, 2017 at 00:43 UTC
b07a775ba85d41b592217d4d6076bb2cb4204727
1 file changed
+189
python.d/rabbitmq.chart.py
new
+189
@@ -0,0 +1,189 @@
1
+# -*- coding: utf-8 -*-
2
+# Description: rabbitmq netdata python.d module
3
+# Author: l2isbad
4
+
5
+from base import UrlService
6
+from socket import gethostbyname, gaierror
7
+try:
8
+ from queue import Queue
9
+except ImportError:
10
+ from Queue import Queue
11
+from threading import Thread
12
+from collections import namedtuple
13
+from json import loads
14
+
15
+# default module values (can be overridden per job in `config`)
16
+# update_every = 2
17
+update_every = 1
18
+priority = 60000
19
+retries = 60
20
+
21
+METHODS = namedtuple('METHODS', ['get_data_function', 'url', 'stats'])
22
+
23
+NODE_STATS = [('fd_used', None),
24
+ ('fd_total', None),
25
+ ('mem_used', None),
26
+ ('mem_limit', None),
27
+ ('sockets_used', None),
28
+ ('sockets_total', None),
29
+ ('proc_total', None),
30
+ ('proc_used', None)
31
+ ]
32
+OVERVIEW_STATS = [('object_totals.channels', None),
33
+ ('object_totals.consumers', None),
34
+ ('object_totals.connections', None),
35
+ ('object_totals.queues', None),
36
+ ('object_totals.exchanges', None),
37
+ ('queue_totals.messages_ready', None),
38
+ ('queue_totals.messages_unacknowledged', None),
39
+ ('message_stats.ack', None),
40
+ ('message_stats.redeliver', None),
41
+ ('message_stats.deliver', None),
42
+ ('message_stats.publish', None)
43
+ ]
44
+ORDER = ['queued_messages', 'message_rates', 'global_counts',
45
+ 'file_descriptors', 'memory', 'sockets', 'processes']
46
+
47
+CHARTS = {
48
+ 'file_descriptors': {
49
+ 'options': [None, 'File Descriptors', 'descriptors', 'overview',
50
+ 'rabbitmq.file_descriptors', 'stacked'],
51
+ 'lines': [
52
+ ['fd_total', 'total', 'absolute'],
53
+ ['fd_used', 'used', 'absolute']
54
+ ]},
55
+ 'memory': {
56
+ 'options': [None, 'Memory', 'MB', 'overview',
57
+ 'rabbitmq.memory', 'stacked'],
58
+ 'lines': [
59
+ ['mem_limit', 'limit', 'absolute', 1, 1024 << 10],
60
+ ['mem_used', 'used', 'absolute', 1, 1024 << 10]
61
+ ]},
62
+ 'sockets': {
63
+ 'options': [None, 'Sockets', 'sockets', 'overview',
64
+ 'rabbitmq.sockets', 'stacked'],
65
+ 'lines': [
66
+ ['sockets_total', 'total', 'absolute'],
67
+ ['sockets_used', 'used', 'absolute']
68
+ ]},
69
+ 'processes': {
70
+ 'options': [None, 'Erlang Processes', 'processes', 'overview',
71
+ 'rabbitmq.processes', 'stacked'],
72
+ 'lines': [
73
+ ['proc_total', 'total', 'absolute'],
74
+ ['proc_used', 'used', 'absolute']
75
+ ]},
76
+ 'global_counts': {
77
+ 'options': [None, 'Global Counts', 'counts', 'overview',
78
+ 'rabbitmq.global_counts', 'line'],
79
+ 'lines': [
80
+ ['channels', None, 'absolute'],
81
+ ['consumers', None, 'absolute'],
82
+ ['connections', None, 'absolute'],
83
+ ['queues', None, 'absolute'],
84
+ ['exchanges', None, 'absolute']
85
+ ]},
86
+ 'queued_messages': {
87
+ 'options': [None, 'Queued Messages', 'messages', 'overview',
88
+ 'rabbitmq.queued_messages', 'stacked'],
89
+ 'lines': [
90
+ ['messages_ready', 'ready', 'incremental'],
91
+ ['messages_unacknowledged', 'unacknowledged', 'incremental']
92
+ ]},
93
+ 'message_rates': {
94
+ 'options': [None, 'Message Rates', 'messages/s', 'overview',
95
+ 'rabbitmq.message_rates', 'stacked'],
96
+ 'lines': [
97
+ ['ack', None, 'incremental'],
98
+ ['redeliver', None, 'incremental'],
99
+ ['deliver', None, 'incremental'],
100
+ ['publish', None, 'incremental']
101
+ ]}
102
+}
103
+
104
+
105
+class Service(UrlService):
106
+ def __init__(self, configuration=None, name=None):
107
+ UrlService.__init__(self, configuration=configuration, name=name)
108
+ self.order = ORDER
109
+ self.definitions = CHARTS
110
+ self.host = self.configuration.get('host', '127.0.0.1')
111
+ self.port = self.configuration.get('port', 15672)
112
+ self.scheme = self.configuration.get('scheme', 'http')
113
+
114
+ def check(self):
115
+ # We can't start if <host> AND <port> not specified
116
+ if not (self.host and self.port):
117
+ self.error('Host is not defined in the module configuration file')
118
+ return False
119
+
120
+ # Hostname -> ip address
121
+ try:
122
+ self.host = gethostbyname(self.host)
123
+ except gaierror as error:
124
+ self.error(str(error))
125
+ return False
126
+
127
+ # Add handlers (auth, self signed cert accept)
128
+ url = '%s://%s:%s/api' % (self.scheme, self.host, self.port)
129
+ self.opener = self._build_opener(url=url)
130
+ if not self.opener:
131
+ return False
132
+ # Add methods
133
+ api_node = url + '/nodes'
134
+ api_overview = url + '/overview'
135
+ self.methods = [METHODS(get_data_function=self._get_overview_stats, url=api_node, stats=NODE_STATS),
136
+ METHODS(get_data_function=self._get_overview_stats, url=api_overview, stats=OVERVIEW_STATS)]
137
+
138
+ result = self._get_data()
139
+ if not result:
140
+ self.error('_get_data() returned no data')
141
+ return False
142
+ self._data_from_check = result
143
+ return True
144
+
145
+ def _get_data(self):
146
+ threads = list()
147
+ queue = Queue()
148
+ result = dict()
149
+
150
+ for method in self.methods:
151
+ th = Thread(target=method.get_data_function, args=(queue, method.url, method.stats))
152
+ th.start()
153
+ threads.append(th)
154
+
155
+ for thread in threads:
156
+ thread.join()
157
+ result.update(queue.get())
158
+
159
+ return result or None
160
+
161
+ def _get_overview_stats(self, queue, url, stats):
162
+ """
163
+ Format data received from http request
164
+ :return: dict
165
+ """
166
+
167
+ raw_data = self._get_raw_data(url)
168
+
169
+ if not raw_data:
170
+ return queue.put(dict())
171
+ data = loads(raw_data)
172
+ data = data[0] if isinstance(data, list) else data
173
+
174
+ to_netdata = fetch_data_(raw_data=data, metrics_list=stats)
175
+ return queue.put(to_netdata)
176
+
177
+
178
+def fetch_data_(raw_data, metrics_list):
179
+ to_netdata = dict()
180
+ for metric, new_name in metrics_list:
181
+ value = raw_data
182
+ for key in metric.split('.'):
183
+ try:
184
+ value = value[key]
185
+ except (KeyError, TypeError):
186
+ break
187
+ if not isinstance(value, dict):
188
+ to_netdata[new_name or key] = value
189
+ return to_netdata