@cryptotaxi247 / netdata-1 / commits / 8721684ce

rabbitmq plugin fixes and optimization

Ilya committed Jul 19, 2017 at 17:43 UTC 8721684cefb72dffe932308d82f7ed90ae0fdb93
1 file changed +59 -60
python.d/rabbitmq.chart.py
+59 -60
@@ -2,40 +2,41 @@
2 # Description: rabbitmq netdata python.d module
3 # Author: l2isbad
4
5 -from base import UrlService
5 +from collections import namedtuple
6 +from json import loads
7 from socket import gethostbyname, gaierror
8 +from threading import Thread
9 try:
10 from queue import Queue
11 except ImportError:
12 from Queue import Queue
11 -from threading import Thread
12 -from collections import namedtuple
13 -from json import loads
13 +
14 +from base import UrlService
15
16 # default module values (can be overridden per job in `config`)
17 update_every = 1
18 priority = 60000
19 retries = 60
20
20 -METHODS = namedtuple('METHODS', ['get_data_function', 'url', 'stats'])
21 +METHODS = namedtuple('METHODS', ['get_data', 'url', 'stats'])
22
22 -NODE_STATS = [('fd_used', None),
23 - ('mem_used', None),
24 - ('sockets_used', None),
25 - ('proc_used', None),
26 - ('disk_free', None)
23 +NODE_STATS = ['fd_used',
24 + 'mem_used',
25 + 'sockets_used',
26 + 'proc_used',
27 + 'disk_free'
28 ]
28 -OVERVIEW_STATS = [('object_totals.channels', None),
29 - ('object_totals.consumers', None),
30 - ('object_totals.connections', None),
31 - ('object_totals.queues', None),
32 - ('object_totals.exchanges', None),
33 - ('queue_totals.messages_ready', None),
34 - ('queue_totals.messages_unacknowledged', None),
35 - ('message_stats.ack', None),
36 - ('message_stats.redeliver', None),
37 - ('message_stats.deliver', None),
38 - ('message_stats.publish', None)
29 +OVERVIEW_STATS = ['object_totals.channels',
30 + 'object_totals.consumers',
31 + 'object_totals.connections',
32 + 'object_totals.queues',
33 + 'object_totals.exchanges',
34 + 'queue_totals.messages_ready',
35 + 'queue_totals.messages_unacknowledged',
36 + 'message_stats.ack',
37 + 'message_stats.redeliver',
38 + 'message_stats.deliver',
39 + 'message_stats.publish'
40 ]
41 ORDER = ['queued_messages', 'message_rates', 'global_counts',
42 'file_descriptors', 'socket_descriptors', 'erlang_processes', 'memory', 'disk_space']
@@ -75,27 +76,27 @@ CHARTS = {
76 'options': [None, 'Global Counts', 'counts', 'overview',
77 'rabbitmq.global_counts', 'line'],
78 'lines': [
78 - ['channels', None, 'absolute'],
79 - ['consumers', None, 'absolute'],
80 - ['connections', None, 'absolute'],
81 - ['queues', None, 'absolute'],
82 - ['exchanges', None, 'absolute']
79 + ['object_totals_channels', 'channels', 'absolute'],
80 + ['object_totals_consumers', 'consumers', 'absolute'],
81 + ['object_totals_connections', 'connections', 'absolute'],
82 + ['object_totals_queues', 'queues', 'absolute'],
83 + ['object_totals_exchanges', 'exchanges', 'absolute']
84 ]},
85 'queued_messages': {
86 'options': [None, 'Queued Messages', 'messages', 'overview',
87 'rabbitmq.queued_messages', 'stacked'],
88 'lines': [
88 - ['messages_ready', 'ready', 'absolute'],
89 - ['messages_unacknowledged', 'unacknowledged', 'absolute']
89 + ['queue_totals_messages_ready', 'ready', 'absolute'],
90 + ['queue_totals_messages_unacknowledged', 'unacknowledged', 'absolute']
91 ]},
92 'message_rates': {
93 'options': [None, 'Message Rates', 'messages/s', 'overview',
94 'rabbitmq.message_rates', 'stacked'],
95 'lines': [
95 - ['ack', None, 'incremental'],
96 - ['redeliver', None, 'incremental'],
97 - ['deliver', None, 'incremental'],
98 - ['publish', None, 'incremental']
96 + ['message_stats_ack', 'ack', 'incremental'],
97 + ['message_stats_redeliver', 'redeliver', 'incremental'],
98 + ['message_stats_deliver', 'deliver', 'incremental'],
99 + ['message_stats_publish', 'publish', 'incremental']
100 ]}
101 }
102
@@ -123,22 +124,19 @@ class Service(UrlService):
124 return False
125
126 # Add handlers (auth, self signed cert accept)
126 - url = '%s://%s:%s/api' % (self.scheme, self.host, self.port)
127 - self.opener = self._build_opener(url=url)
128 - if not self.opener:
129 - return False
127 + self.url = '{scheme}://{host}:{port}/api'.format(scheme=self.scheme,
128 + host=self.host,
129 + port=self.port)
130 # Add methods
131 - api_node = url + '/nodes'
132 - api_overview = url + '/overview'
133 - self.methods = [METHODS(get_data_function=self._get_overview_stats, url=api_node, stats=NODE_STATS),
134 - METHODS(get_data_function=self._get_overview_stats, url=api_overview, stats=OVERVIEW_STATS)]
135 -
136 - result = self._get_data()
137 - if not result:
138 - self.error('_get_data() returned no data')
139 - return False
140 - self._data_from_check = result
141 - return True
131 + api_node = self.url + '/nodes'
132 + api_overview = self.url + '/overview'
133 + self.methods = [METHODS(get_data=self._get_overview_stats,
134 + url=api_node,
135 + stats=NODE_STATS),
136 + METHODS(get_data=self._get_overview_stats,
137 + url=api_overview,
138 + stats=OVERVIEW_STATS)]
139 + return UrlService.check(self)
140
141 def _get_data(self):
142 threads = list()
@@ -146,7 +144,8 @@ class Service(UrlService):
144 result = dict()
145
146 for method in self.methods:
149 - th = Thread(target=method.get_data_function, args=(queue, method.url, method.stats))
147 + th = Thread(target=method.get_data,
148 + args=(queue, method.url, method.stats))
149 th.start()
150 threads.append(th)
151
@@ -169,19 +168,19 @@ class Service(UrlService):
168 data = loads(raw_data)
169 data = data[0] if isinstance(data, list) else data
170
172 - to_netdata = fetch_data_(raw_data=data, metrics_list=stats)
171 + to_netdata = fetch_data(raw_data=data, metrics=stats)
172 return queue.put(to_netdata)
173
174
176 -def fetch_data_(raw_data, metrics_list):
177 - to_netdata = dict()
178 - for metric, new_name in metrics_list:
175 +def fetch_data(raw_data, metrics):
176 + data = dict()
177 + for metric in metrics:
178 value = raw_data
180 - for key in metric.split('.'):
181 - try:
182 - value = value[key]
183 - except (KeyError, TypeError):
184 - break
185 - if not isinstance(value, dict):
186 - to_netdata[new_name or key] = value
187 - return to_netdata
179 + metrics_list = metric.split('.')
180 + try:
181 + for m in metrics_list:
182 + value = value[m]
183 + except KeyError:
184 + continue
185 + data['_'.join(metrics_list)] = value
186 + return data