elasticsearch: handle json parse error in threads (#4186)
Ilya Mashchenko committed
Sep 13, 2018 at 17:07 UTC
16b2a57ad0ab7ac40a7e1ecdebcc02c881086366
1 file changed
+64
-18
python.d/elasticsearch.chart.py
+64
-18
@@ -3,10 +3,12 @@
3
# Author: l2isbad
4
# SPDX-License-Identifier: GPL-3.0+
5
6
+import json
7
+import threading
8
+
9
from collections import namedtuple
7
-from json import loads
10
from socket import gethostbyname, gaierror
9
-from threading import Thread
11
+
12
try:
13
from queue import Queue
14
except ImportError:
@@ -16,8 +18,6 @@ from bases.FrameworkServices.UrlService import UrlService
18
19
# default module values (can be overridden per job in `config`)
20
update_every = 5
19
-priority = 60000
20
-retries = 60
21
22
METHODS = namedtuple('METHODS', ['get_data', 'url', 'run'])
23
@@ -125,15 +125,43 @@ LATENCY = {
125
}
126
127
# charts order (can be overridden if you want less charts, or different order)
128
-ORDER = ['search_performance_total', 'search_performance_current', 'search_performance_time',
129
- 'search_latency', 'index_performance_total', 'index_performance_current', 'index_performance_time',
130
- 'index_latency', 'index_translog_operations', 'index_translog_size', 'index_segments_count', 'index_segments_memory_writer',
131
- 'index_segments_memory', 'jvm_mem_heap', 'jvm_mem_heap_bytes', 'jvm_buffer_pool_count',
132
- 'jvm_direct_buffers_memory', 'jvm_mapped_buffers_memory', 'jvm_gc_count', 'jvm_gc_time', 'host_metrics_file_descriptors',
133
- 'host_metrics_http', 'host_metrics_transport', 'thread_pool_queued', 'thread_pool_rejected',
134
- 'fielddata_cache', 'fielddata_evictions_tripped', 'cluster_health_status', 'cluster_health_nodes',
135
- 'cluster_health_shards', 'cluster_stats_nodes', 'cluster_stats_query_cache', 'cluster_stats_docs',
136
- 'cluster_stats_store', 'cluster_stats_indices_shards']
128
+ORDER = [
129
+ 'search_performance_total',
130
+ 'search_performance_current',
131
+ 'search_performance_time',
132
+ 'search_latency',
133
+ 'index_performance_total',
134
+ 'index_performance_current',
135
+ 'index_performance_time',
136
+ 'index_latency',
137
+ 'index_translog_operations',
138
+ 'index_translog_size',
139
+ 'index_segments_count',
140
+ 'index_segments_memory_writer',
141
+ 'index_segments_memory',
142
+ 'jvm_mem_heap',
143
+ 'jvm_mem_heap_bytes',
144
+ 'jvm_buffer_pool_count',
145
+ 'jvm_direct_buffers_memory',
146
+ 'jvm_mapped_buffers_memory',
147
+ 'jvm_gc_count',
148
+ 'jvm_gc_time',
149
+ 'host_metrics_file_descriptors',
150
+ 'host_metrics_http',
151
+ 'host_metrics_transport',
152
+ 'thread_pool_queued',
153
+ 'thread_pool_rejected',
154
+ 'fielddata_cache',
155
+ 'fielddata_evictions_tripped',
156
+ 'cluster_health_status',
157
+ 'cluster_health_nodes',
158
+ 'cluster_health_shards',
159
+ 'cluster_stats_nodes',
160
+ 'cluster_stats_query_cache',
161
+ 'cluster_stats_docs',
162
+ 'cluster_stats_store',
163
+ 'cluster_stats_indices_shards',
164
+]
165
166
CHARTS = {
167
'search_performance_total': {
@@ -449,8 +477,8 @@ class Service(UrlService):
477
for method in self.methods:
478
if not method.run:
479
continue
452
- th = Thread(target=method.get_data,
453
- args=(queue, method.url))
480
+ th = threading.Thread(target=method.get_data,
481
+ args=(queue, method.url))
482
th.start()
483
threads.append(th)
484
@@ -471,7 +499,11 @@ class Service(UrlService):
499
if not raw_data:
500
return queue.put(dict())
501
474
- data = loads(raw_data)
502
+ data = self.json_reply(raw_data)
503
+
504
+ if not data:
505
+ return queue.put(dict())
506
+
507
to_netdata = fetch_data_(raw_data=data,
508
metrics=HEALTH_STATS)
509
@@ -493,7 +525,11 @@ class Service(UrlService):
525
if not raw_data:
526
return queue.put(dict())
527
496
- data = loads(raw_data)
528
+ data = self.json_reply(raw_data)
529
+
530
+ if not data:
531
+ return queue.put(dict())
532
+
533
to_netdata = fetch_data_(raw_data=data,
534
metrics=CLUSTER_STATS)
535
@@ -510,7 +546,10 @@ class Service(UrlService):
546
if not raw_data:
547
return queue.put(dict())
548
513
- data = loads(raw_data)
549
+ data = self.json_reply(raw_data)
550
+
551
+ if not data:
552
+ return queue.put(dict())
553
554
node = list(data['nodes'].keys())[0]
555
to_netdata = fetch_data_(raw_data=data['nodes'][node],
@@ -530,6 +569,13 @@ class Service(UrlService):
569
570
return queue.put(to_netdata)
571
572
+ def json_reply(self, reply):
573
+ try:
574
+ return json.loads(reply)
575
+ except ValueError as err:
576
+ self.error(err)
577
+ return None
578
+
579
def find_avg(self, total, spent_time, key):
580
if key not in self.latency:
581
self.latency[key] = dict(total=total,