mongodb_plugin: oplog window and "optimeDate" diff between nodes charts added
Ilya committed
Feb 24, 2017 at 19:30 UTC
878a350b94fda1838240811d126bf39cbf659500
1 file changed
+85
-28
python.d/mongodb.chart.py
+85
-28
@@ -5,8 +5,9 @@
5
from base import SimpleService
6
from copy import deepcopy
7
from datetime import datetime
8
+from sys import exc_info
9
try:
9
- from pymongo import MongoClient
10
+ from pymongo import MongoClient, ASCENDING, DESCENDING
11
from pymongo.errors import PyMongoError
12
PYMONGO = True
13
except ImportError:
@@ -123,6 +124,7 @@ CHARTS = {
124
'lines': [
125
['memory_virtual', 'virtual', 'absolute', 1, 1],
126
['memory_resident', 'resident', 'absolute', 1, 1],
127
+ ['memory_nonmapped', 'nonmapped', 'absolute', 1, 1],
128
['memory_mapped', 'mapped', 'absolute', 1, 1]
129
]},
130
'page_faults': {
@@ -154,11 +156,11 @@ CHARTS = {
156
['errors_user', 'user', 'incremental', 1, 1]
157
]},
158
'wiredtiger_cache': {
157
- 'options': [None, 'Amount of space taken by cached data and by dirty data in the cache',
158
- 'KB', 'resource utilization', 'mongodb.wiredtiger_cache', 'stacked'],
159
+ 'options': [None, 'The percentage of the wiredTiger cache that is in use and cache with dirty bytes',
160
+ 'percent', 'resource utilization', 'mongodb.wiredtiger_cache', 'stacked'],
161
'lines': [
160
- ['wiredTiger_bytes_in_cache', 'cached', 'absolute', 1, 1024],
161
- ['wiredTiger_dirty_in_cache', 'dirty', 'absolute', 1, 1024]
162
+ ['wiredTiger_bytes_in_cache', 'inuse', 'absolute', 1, 1000],
163
+ ['wiredTiger_dirty_in_cache', 'dirty', 'absolute', 1, 1000]
164
]},
165
'wiredtiger_pages_evicted': {
166
'options': [None, 'Pages evicted from the cache',
@@ -231,25 +233,29 @@ class Service(SimpleService):
233
self.error(error)
234
return False
235
234
- self.repl = 'repl' in server_status
236
try:
237
self.databases = self.connection.database_names()
238
except PyMongoError as error:
239
self.databases = list()
240
self.info('Can\'t collect databases: %s' % str(error))
241
241
- self.create_charts_(server_status)
242
+ self.ss = dict()
243
+ for elem in ['dur', 'backgroundFlushing', 'wiredTiger', 'tcmalloc', 'cursor', 'commands', 'repl']:
244
+ self.ss[elem] = in_server_status(elem, server_status)
245
243
- return True
246
+ try:
247
+ self._get_data()
248
+ except (LookupError, SyntaxError, AttributeError):
249
+ self.error('Type: %s, error: %s' % (str(exc_info()[0]), str( exc_info()[1])))
250
+ return False
251
+ else:
252
+ self.create_charts_(server_status)
253
+ return True
254
255
def create_charts_(self, server_status):
256
257
self.order = ORDER[:]
258
self.definitions = deepcopy(CHARTS)
249
- self.ss = dict()
250
-
251
- for elem in ['dur', 'backgroundFlushing', 'wiredTiger', 'tcmalloc', 'cursor', 'commands']:
252
- self.ss[elem] = in_server_status(elem, server_status)
259
260
if not self.ss['dur']:
261
self.order.remove('journaling_transactions')
@@ -288,11 +294,11 @@ class Service(SimpleService):
294
]}
295
self.definitions['dbstats_objects']['lines'].append(['_'.join([dbase, 'objects']), dbase, 'absolute'])
296
291
- if server_status.get('repl'):
292
- def create_heartbeat_lines(hosts):
297
+ if self.ss['repl']:
298
+ def create_lines(hosts, string):
299
lines = list()
300
for host in hosts:
295
- dim_id = '_'.join([host, 'heartbeat_lag'])
301
+ dim_id = '_'.join([host, string])
302
lines.append([dim_id, host, 'absolute', 1, 1000])
303
return lines
304
@@ -307,12 +313,25 @@ class Service(SimpleService):
313
this_host = server_status['repl']['me']
314
other_hosts = [host for host in all_hosts if host != this_host]
315
310
- # Create "heartbeat delay" charts
316
+ if 'local' in self.databases:
317
+ self.order.append('oplog_window')
318
+ self.definitions['oplog_window'] = {
319
+ 'options': [None, 'Interval of time between the oldest and the latest entries in the oplog',
320
+ 'seconds', 'replication', 'mongodb.oplog_window', 'line'],
321
+ 'lines': [['timeDiff', 'window', 'absolute', 1, 1000]]}
322
+ # Create "heartbeat delay" chart
323
self.order.append('heartbeat_delay')
324
self.definitions['heartbeat_delay'] = {
325
'options': [None, 'Latency between this node and replica set members (lastHeartbeatRecv)',
326
'seconds', 'replication', 'mongodb.replication_heartbeat_delay', 'stacked'],
315
- 'lines': create_heartbeat_lines(other_hosts)}
327
+ 'lines': create_lines(other_hosts, 'heartbeat_lag')}
328
+ # Create "optimedate delay" chart
329
+ self.order.append('optimedate_delay')
330
+ self.definitions['optimedate_delay'] = {
331
+ 'options': [None, '"optimeDate"(time when last entry from the oplog was applied)'
332
+ ' diff between all nodes',
333
+ 'seconds', 'replication', 'mongodb.replication_optimedate_delay', 'stacked'],
334
+ 'lines': create_lines(all_hosts, 'optimedate')}
335
# Create "replica set members state" chart
336
for host in all_hosts:
337
chart_name = '_'.join([host, 'state'])
@@ -328,6 +347,7 @@ class Service(SimpleService):
347
raw_data.update(self.get_serverstatus_() or dict())
348
raw_data.update(self.get_dbstats_() or dict())
349
raw_data.update(self.get_replsetgetstatus_() or dict())
350
+ raw_data.update(self.get_getreplicationinfo_() or dict())
351
352
return raw_data or None
353
@@ -355,7 +375,7 @@ class Service(SimpleService):
375
return raw_data
376
377
def get_replsetgetstatus_(self):
358
- if not self.repl:
378
+ if not self.ss['repl']:
379
return None
380
381
raw_data = dict()
@@ -366,6 +386,22 @@ class Service(SimpleService):
386
else:
387
return raw_data
388
389
+ def get_getreplicationinfo_(self):
390
+ if not (self.ss['repl'] and 'local' in self.databases):
391
+ return None
392
+
393
+ raw_data = dict()
394
+ raw_data['getReplicationInfo'] = dict()
395
+ try:
396
+ raw_data['getReplicationInfo']['ASCENDING'] = self.connection.local.oplog.rs.find().sort(
397
+ "$natural", ASCENDING).limit(1)[0]
398
+ raw_data['getReplicationInfo']['DESCENDING'] = self.connection.local.oplog.rs.find().sort(
399
+ "$natural", DESCENDING).limit(1)[0]
400
+ except PyMongoError:
401
+ return None
402
+ else:
403
+ return raw_data
404
+
405
def _get_data(self):
406
"""
407
:return: dict
@@ -379,6 +415,7 @@ class Service(SimpleService):
415
serverStatus = raw_data['serverStatus']
416
dbStats = raw_data.get('dbStats')
417
replSetGetStatus = raw_data.get('replSetGetStatus')
418
+ getReplicationInfo = raw_data.get('getReplicationInfo')
419
utc_now = datetime.utcnow()
420
421
# serverStatus
@@ -386,6 +423,8 @@ class Service(SimpleService):
423
to_netdata.update(update_dict_key(serverStatus['globalLock']['activeClients'], 'activeClients'))
424
to_netdata.update(update_dict_key(serverStatus['connections'], 'connections'))
425
to_netdata.update(update_dict_key(serverStatus['mem'], 'memory'))
426
+ to_netdata['memory_nonmapped'] = (serverStatus['mem']['virtual']
427
+ - serverStatus['mem'].get('mappedWithJournal', serverStatus['mem']['mapped']))
428
to_netdata.update(update_dict_key(serverStatus['globalLock']['currentQueue'], 'currentQueue'))
429
to_netdata.update(update_dict_key(serverStatus['asserts'], 'errors'))
430
to_netdata['page_faults'] = serverStatus['extra_info']['page_faults']
@@ -410,8 +449,12 @@ class Service(SimpleService):
449
'wiredTigerRead'))
450
to_netdata.update(update_dict_key(serverStatus['wiredTiger']['concurrentTransactions']['write'],
451
'wiredTigerWrite'))
413
- to_netdata['wiredTiger_bytes_in_cache'] = wired_tiger['cache']['bytes currently in the cache']
414
- to_netdata['wiredTiger_dirty_in_cache'] = wired_tiger['cache']['tracked dirty bytes in the cache']
452
+ to_netdata['wiredTiger_bytes_in_cache'] = (int(wired_tiger['cache']['bytes currently in the cache']
453
+ * 100 / wired_tiger['cache']['maximum bytes configured']
454
+ * 1000))
455
+ to_netdata['wiredTiger_dirty_in_cache'] = (int(wired_tiger['cache']['tracked dirty bytes in the cache']
456
+ * 100 / wired_tiger['cache']['maximum bytes configured']
457
+ * 1000))
458
to_netdata['wiredTiger_unmodified_pages_evicted'] = wired_tiger['cache']['unmodified pages evicted']
459
to_netdata['wiredTiger_modified_pages_evicted'] = wired_tiger['cache']['modified pages evicted']
460
@@ -433,19 +476,33 @@ class Service(SimpleService):
476
if replSetGetStatus:
477
other_hosts = list()
478
members = replSetGetStatus['members']
479
+ unix_epoch = datetime(1970, 1, 1, 0, 0)
480
+
481
for member in members:
482
if not member.get('self'):
483
other_hosts.append(member)
484
+ # Replica set time diff between current time and time when last entry from the oplog was applied
485
+ if member['optimeDate'] != unix_epoch:
486
+ member_optimedate = member['name'] + '_optimedate'
487
+ to_netdata.update({member_optimedate: int(delta_calculation(delta=utc_now - member['optimeDate'],
488
+ multiplier=1000))})
489
# Replica set members state
490
+ member_state = member['name'] + '_state'
491
for elem in REPLSET_STATES:
492
state = elem[0]
442
- to_netdata.update({'_'.join([member['name'], 'state', state]): 0})
443
- to_netdata.update({'_'.join([member['name'], 'state', str(member['state'])]): member['state']})
493
+ to_netdata.update({'_'.join([member_state, state]): 0})
494
+ to_netdata.update({'_'.join([member_state, str(member['state'])]): member['state']})
495
# Heartbeat lag calculation
496
for other in other_hosts:
446
- if other['lastHeartbeatRecv'] != datetime(1970, 1, 1, 0, 0):
497
+ if other['lastHeartbeatRecv'] != unix_epoch:
498
node = other['name'] + '_heartbeat_lag'
448
- to_netdata[node] = int(lag_calculation(utc_now - other['lastHeartbeatRecv']) * 1000)
499
+ to_netdata[node] = int(delta_calculation(delta=utc_now - other['lastHeartbeatRecv'],
500
+ multiplier=1000))
501
+
502
+ if getReplicationInfo:
503
+ first_event = getReplicationInfo['ASCENDING']['ts'].as_datetime()
504
+ last_event = getReplicationInfo['DESCENDING']['ts'].as_datetime()
505
+ to_netdata['timeDiff'] = int(delta_calculation(delta=last_event - first_event, multiplier=1000))
506
507
return to_netdata
508
@@ -478,8 +535,8 @@ def in_server_status(elem, server_status):
535
return elem in server_status or elem in server_status['metrics']
536
537
481
-def lag_calculation(lag):
482
- if hasattr(lag, 'total_seconds'):
483
- return lag.total_seconds()
538
+def delta_calculation(delta, multiplier=1):
539
+ if hasattr(delta, 'total_seconds'):
540
+ return delta.total_seconds() * multiplier
541
else:
485
- return (lag.microseconds + (lag.seconds + lag.days * 24 * 3600) * 10 ** 6) / 10.0 ** 6
542
+ return (delta.microseconds + (delta.seconds + delta.days * 24 * 3600) * 10 ** 6) / 10.0 ** 6 * multiplier