@cryptotaxi247 / netdata-1 / commits / d0de3356f

postgres fix: detect servers version and use the right query (#4944)

* python.d postgres: use query depending on server version * python.d postgres repslot file query fix * python.d postgres minor * python.d postgres format fix * python.d postgres refactor queries * postgres query repslot files v11 fix * @anayrat patch

Ilya Mashchenko committed Dec 12, 2018 at 23:56 UTC d0de3356f3d83d273b2a6a19dfd5920b54a6f391
1 file changed +463 -200
collectors/python.d.plugin/postgres/postgres.chart.py
+463 -200
@@ -16,12 +16,28 @@ except ImportError:
16
17 from bases.FrameworkServices.SimpleService import SimpleService
18
19 -# default module values
20 -update_every = 1
21 -priority = 60000
19 +
20 +WAL = 'WAL'
21 +ARCHIVE = 'ARCHIVE'
22 +BACKENDS = 'BACKENDS'
23 +TABLE_STATS = 'TABLE_STATS'
24 +INDEX_STATS = 'INDEX_STATS'
25 +DATABASE = 'DATABASE'
26 +BGWRITER = 'BGWRITER'
27 +LOCKS = 'LOCKS'
28 +DATABASES = 'DATABASES'
29 +STANDBY = 'STANDBY'
30 +REPLICATION_SLOT = 'REPLICATION_SLOT'
31 +STANDBY_DELTA = 'STANDBY_DELTA'
32 +REPSLOT_FILES = 'REPSLOT_FILES'
33 +IF_SUPERUSER = 'IF_SUPERUSER'
34 +SERVER_VERSION = 'SERVER_VERSION'
35 +AUTOVACUUM = 'AUTOVACUUM'
36 +DIFF_LSN = 'DIFF_LSN'
37 +WAL_WRITES = 'WAL_WRITES'
38
39 METRICS = {
24 - 'DATABASE': [
40 + DATABASE: [
41 'connections',
42 'xact_commit',
43 'xact_rollback',
@@ -37,32 +53,32 @@ METRICS = {
53 'temp_bytes',
54 'size'
55 ],
40 - 'BACKENDS': [
56 + BACKENDS: [
57 'backends_active',
58 'backends_idle'
59 ],
44 - 'INDEX_STATS': [
60 + INDEX_STATS: [
61 'index_count',
62 'index_size'
63 ],
48 - 'TABLE_STATS': [
64 + TABLE_STATS: [
65 'table_size',
66 'table_count'
67 ],
52 - 'WAL': [
68 + WAL: [
69 'written_wal',
70 'recycled_wal',
71 'total_wal'
72 ],
57 - 'WAL_WRITES': [
73 + WAL_WRITES: [
74 'wal_writes'
75 ],
60 - 'ARCHIVE': [
76 + ARCHIVE: [
77 'ready_count',
78 'done_count',
79 'file_count'
80 ],
65 - 'BGWRITER': [
81 + BGWRITER: [
82 'checkpoint_scheduled',
83 'checkpoint_requested',
84 'buffers_checkpoint',
@@ -72,7 +88,7 @@ METRICS = {
88 'buffers_alloc',
89 'buffers_backend_fsync'
90 ],
75 - 'LOCKS': [
91 + LOCKS: [
92 'ExclusiveLock',
93 'RowShareLock',
94 'SIReadLock',
@@ -83,27 +99,34 @@ METRICS = {
99 'ShareLock',
100 'RowExclusiveLock'
101 ],
86 - 'AUTOVACUUM': [
102 + AUTOVACUUM: [
103 'analyze',
104 'vacuum_analyze',
105 'vacuum',
106 'vacuum_freeze',
107 'brin_summarize'
108 ],
93 - 'STANDBY_DELTA': [
109 + STANDBY_DELTA: [
110 'sent_delta',
111 'write_delta',
112 'flush_delta',
113 'replay_delta'
114 ],
99 - 'REPSLOT_FILES': [
115 + REPSLOT_FILES: [
116 'replslot_wal_keep',
117 'replslot_files'
118 ]
119 }
120
105 -QUERIES = {
106 - 'WAL': """
121 +NO_VERSION = 0
122 +DEFAULT = 'DEFAULT'
123 +V96 = 'V96'
124 +V10 = 'V10'
125 +V11 = 'V11'
126 +
127 +
128 +QUERY_WAL = {
129 + DEFAULT: """
130 SELECT
131 count(*) as total_wal,
132 count(*) FILTER (WHERE type = 'recycled') AS recycled_wal,
@@ -111,34 +134,76 @@ SELECT
134 FROM
135 (SELECT
136 wal.name,
114 - pg_{0}file_name(
137 + pg_walfile_name(
138 CASE pg_is_in_recovery()
139 WHEN true THEN NULL
117 - ELSE pg_current_{0}_{1}()
140 + ELSE pg_current_wal_lsn()
141 END ),
142 CASE
120 - WHEN wal.name > pg_{0}file_name(
143 + WHEN wal.name > pg_walfile_name(
144 CASE pg_is_in_recovery()
145 WHEN true THEN NULL
123 - ELSE pg_current_{0}_{1}()
146 + ELSE pg_current_wal_lsn()
147 END ) THEN 'recycled'
148 ELSE 'written'
149 END AS type
127 - FROM pg_catalog.pg_ls_dir('pg_{0}') AS wal(name)
150 + FROM pg_catalog.pg_ls_dir('pg_wal') AS wal(name)
151 WHERE name ~ '^[0-9A-F]{{24}}$'
152 ORDER BY
130 - (pg_stat_file('pg_{0}/'||name)).modification,
153 + (pg_stat_file('pg_wal/'||name)).modification,
154 wal.name DESC) sub;
155 """,
133 - 'ARCHIVE': """
156 + V96: """
157 +SELECT
158 + count(*) as total_wal,
159 + count(*) FILTER (WHERE type = 'recycled') AS recycled_wal,
160 + count(*) FILTER (WHERE type = 'written') AS written_wal
161 +FROM
162 + (SELECT
163 + wal.name,
164 + pg_xlogfile_name(
165 + CASE pg_is_in_recovery()
166 + WHEN true THEN NULL
167 + ELSE pg_current_xlog_location()
168 + END ),
169 + CASE
170 + WHEN wal.name > pg_xlogfile_name(
171 + CASE pg_is_in_recovery()
172 + WHEN true THEN NULL
173 + ELSE pg_current_xlog_location()
174 + END ) THEN 'recycled'
175 + ELSE 'written'
176 + END AS type
177 + FROM pg_catalog.pg_ls_dir('pg_xlog') AS wal(name)
178 + WHERE name ~ '^[0-9A-F]{{24}}$'
179 + ORDER BY
180 + (pg_stat_file('pg_xlog/'||name)).modification,
181 + wal.name DESC) sub;
182 +""",
183 +}
184 +
185 +QUERY_ARCHIVE = {
186 + DEFAULT: """
187 SELECT
188 CAST(COUNT(*) AS INT) AS file_count,
189 CAST(COALESCE(SUM(CAST(archive_file ~ $r$\.ready$$r$ as INT)),0) AS INT) AS ready_count,
190 CAST(COALESCE(SUM(CAST(archive_file ~ $r$\.done$$r$ AS INT)),0) AS INT) AS done_count
191 FROM
139 - pg_catalog.pg_ls_dir('pg_{0}/archive_status') AS archive_files (archive_file);
192 + pg_catalog.pg_ls_dir('pg_wal/archive_status') AS archive_files (archive_file);
193 """,
141 - 'BACKENDS': """
194 + V96: """
195 +SELECT
196 + CAST(COUNT(*) AS INT) AS file_count,
197 + CAST(COALESCE(SUM(CAST(archive_file ~ $r$\.ready$$r$ as INT)),0) AS INT) AS ready_count,
198 + CAST(COALESCE(SUM(CAST(archive_file ~ $r$\.done$$r$ AS INT)),0) AS INT) AS done_count
199 +FROM
200 + pg_catalog.pg_ls_dir('pg_xlog/archive_status') AS archive_files (archive_file);
201 +
202 +""",
203 +}
204 +
205 +QUERY_BACKEND = {
206 + DEFAULT: """
207 SELECT
208 count(*) - (SELECT count(*)
209 FROM pg_stat_activity
@@ -150,21 +215,30 @@ SELECT
215 AS backends_idle
216 FROM pg_stat_activity;
217 """,
153 - 'TABLE_STATS': """
218 +}
219 +
220 +QUERY_TABLE_STATS = {
221 + DEFAULT: """
222 SELECT
223 ((sum(relpages) * 8) * 1024) AS table_size,
224 count(1) AS table_count
225 FROM pg_class
226 WHERE relkind IN ('r', 't');
227 """,
160 - 'INDEX_STATS': """
228 +}
229 +
230 +QUERY_INDEX_STATS = {
231 + DEFAULT: """
232 SELECT
233 ((sum(relpages) * 8) * 1024) AS index_size,
234 count(1) AS index_count
235 FROM pg_class
236 WHERE relkind = 'i';
237 """,
167 - 'DATABASE': """
238 +}
239 +
240 +QUERY_DATABASE = {
241 + DEFAULT: """
242 SELECT
243 datname AS database_name,
244 numbackends AS connections,
@@ -184,7 +258,10 @@ SELECT
258 FROM pg_stat_database
259 WHERE datname IN %(databases)s ;
260 """,
187 - 'BGWRITER': """
261 +}
262 +
263 +QUERY_BGWRITER = {
264 + DEFAULT: """
265 SELECT
266 checkpoints_timed AS checkpoint_scheduled,
267 checkpoints_req AS checkpoint_requested,
@@ -196,7 +273,10 @@ SELECT
273 buffers_backend_fsync
274 FROM pg_stat_bgwriter;
275 """,
199 - 'LOCKS': """
276 +}
277 +
278 +QUERY_LOCKS = {
279 + DEFAULT: """
280 SELECT
281 pg_database.datname as database_name,
282 mode,
@@ -207,7 +287,10 @@ INNER JOIN pg_database
287 GROUP BY datname, mode
288 ORDER BY datname, mode;
289 """,
210 - 'FIND_DATABASES': """
290 +}
291 +
292 +QUERY_DATABASES = {
293 + DEFAULT: """
294 SELECT
295 datname
296 FROM pg_stat_database
@@ -216,48 +299,129 @@ WHERE
299 (SELECT current_user), datname, 'connect')
300 AND NOT datname ~* '^template\d ';
301 """,
219 - 'FIND_STANDBY': """
302 +}
303 +
304 +QUERY_STANDBY = {
305 + DEFAULT: """
306 SELECT
307 application_name
308 FROM pg_stat_replication
309 WHERE application_name IS NOT NULL
310 GROUP BY application_name;
311 """,
226 - 'FIND_REPLICATION_SLOT': """
312 +}
313 +
314 +QUERY_REPLICATION_SLOT = {
315 + DEFAULT: """
316 SELECT slot_name
317 FROM pg_replication_slots;
318 +"""
319 +}
320 +
321 +QUERY_STANDBY_DELTA = {
322 + DEFAULT: """
323 +SELECT
324 + application_name,
325 + pg_wal_lsn_diff(
326 + CASE pg_is_in_recovery()
327 + WHEN true THEN pg_last_wal_receive_lsn()
328 + ELSE pg_current_wal_lsn()
329 + END,
330 + sent_lsn) AS sent_delta,
331 + pg_wal_lsn_diff(
332 + CASE pg_is_in_recovery()
333 + WHEN true THEN pg_last_wal_receive_lsn()
334 + ELSE pg_current_wal_lsn()
335 + END,
336 + write_lsn) AS write_delta,
337 + pg_wal_lsn_diff(
338 + CASE pg_is_in_recovery()
339 + WHEN true THEN pg_last_wal_receive_lsn()
340 + ELSE pg_current_wal_lsn()
341 + END,
342 + flush_lsn) AS flush_delta,
343 + pg_wal_lsn_diff(
344 + CASE pg_is_in_recovery()
345 + WHEN true THEN pg_last_wal_receive_lsn()
346 + ELSE pg_current_wal_lsn()
347 + END,
348 + replay_lsn) AS replay_delta
349 +FROM pg_stat_replication
350 +WHERE application_name IS NOT NULL;
351 """,
230 - 'STANDBY_DELTA': """
352 + V96: """
353 SELECT
354 application_name,
233 - pg_{0}_{1}_diff(
355 + pg_xlog_location_diff(
356 CASE pg_is_in_recovery()
235 - WHEN true THEN pg_last_{0}_receive_{1}()
236 - ELSE pg_current_{0}_{1}()
357 + WHEN true THEN pg_last_xlog_receive_location()
358 + ELSE pg_current_xlog_location()
359 END,
238 - sent_{1}) AS sent_delta,
239 - pg_{0}_{1}_diff(
360 + sent_location) AS sent_delta,
361 + pg_xlog_location_diff(
362 CASE pg_is_in_recovery()
241 - WHEN true THEN pg_last_{0}_receive_{1}()
242 - ELSE pg_current_{0}_{1}()
363 + WHEN true THEN pg_last_xlog_receive_location()
364 + ELSE pg_current_xlog_location()
365 END,
244 - write_{1}) AS write_delta,
245 - pg_{0}_{1}_diff(
366 + write_location) AS write_delta,
367 + pg_xlog_location_diff(
368 CASE pg_is_in_recovery()
247 - WHEN true THEN pg_last_{0}_receive_{1}()
248 - ELSE pg_current_{0}_{1}()
369 + WHEN true THEN pg_last_xlog_receive_location()
370 + ELSE pg_current_xlog_location()
371 END,
250 - flush_{1}) AS flush_delta,
251 - pg_{0}_{1}_diff(
372 + flush_location) AS flush_delta,
373 + pg_xlog_location_diff(
374 CASE pg_is_in_recovery()
253 - WHEN true THEN pg_last_{0}_receive_{1}()
254 - ELSE pg_current_{0}_{1}()
375 + WHEN true THEN pg_last_xlog_receive_location()
376 + ELSE pg_current_xlog_location()
377 END,
256 - replay_{1}) AS replay_delta
378 + replay_location) AS replay_delta
379 FROM pg_stat_replication
380 WHERE application_name IS NOT NULL;
381 """,
260 - 'REPSLOT_FILES': """
382 +}
383 +
384 +QUERY_REPSLOT_FILES = {
385 + DEFAULT: """
386 +WITH wal_size AS (
387 + SELECT
388 + setting::int AS val
389 + FROM pg_settings
390 + WHERE name = 'wal_segment_size'
391 + )
392 +SELECT
393 + slot_name,
394 + slot_type,
395 + replslot_wal_keep,
396 + count(slot_file) AS replslot_files
397 +FROM
398 + (SELECT
399 + slot.slot_name,
400 + CASE
401 + WHEN slot_file <> 'state' THEN 1
402 + END AS slot_file ,
403 + slot_type,
404 + COALESCE (
405 + floor(
406 + (pg_wal_lsn_diff(pg_current_wal_lsn (),slot.restart_lsn)
407 + - (pg_walfile_name_offset (restart_lsn)).file_offset) / (s.val)
408 + ),0) AS replslot_wal_keep
409 + FROM pg_replication_slots slot
410 + LEFT JOIN (
411 + SELECT
412 + slot2.slot_name,
413 + pg_ls_dir('pg_replslot/' || slot2.slot_name) AS slot_file
414 + FROM pg_replication_slots slot2
415 + ) files (slot_name, slot_file)
416 + ON slot.slot_name = files.slot_name
417 + CROSS JOIN wal_size s
418 + ) AS d
419 +GROUP BY
420 + slot_name,
421 + slot_type,
422 + replslot_wal_keep;
423 +""",
424 + V10: """
425 WITH wal_size AS (
426 SELECT
427 current_setting('wal_block_size')::INT * setting::INT AS val
@@ -296,13 +460,22 @@ GROUP BY
460 slot_type,
461 replslot_wal_keep;
462 """,
299 - 'IF_SUPERUSER': """
463 +}
464 +
465 +QUERY_SUPERUSER = {
466 + DEFAULT: """
467 SELECT current_setting('is_superuser') = 'on' AS is_superuser;
468 """,
302 - 'DETECT_SERVER_VERSION': """
469 +}
470 +
471 +QUERY_SHOW_VERSION = {
472 + DEFAULT: """
473 SHOW server_version_num;
474 """,
305 - 'AUTOVACUUM': """
475 +}
476 +
477 +QUERY_AUTOVACUUM = {
478 + DEFAULT: """
479 SELECT
480 count(*) FILTER (WHERE query LIKE 'autovacuum: ANALYZE%%') AS analyze,
481 count(*) FILTER (WHERE query LIKE 'autovacuum: VACUUM ANALYZE%%') AS vacuum_analyze,
@@ -314,23 +487,78 @@ SELECT
487 FROM pg_stat_activity
488 WHERE query NOT LIKE '%%pg_stat_activity%%';
489 """,
317 - 'DIFF_LSN': """
490 +}
491 +
492 +QUERY_DIFF_LSN = {
493 + DEFAULT: """
494 SELECT
319 - pg_{0}_{1}_diff(
495 + pg_wal_lsn_diff(
496 CASE pg_is_in_recovery()
321 - WHEN true THEN pg_last_{0}_receive_{1}()
322 - ELSE pg_current_{0}_{1}()
497 + WHEN true THEN pg_last_wal_receive_lsn()
498 + ELSE pg_current_wal_lsn()
499 END,
500 '0/0') as wal_writes ;
325 -"""
501 +""",
502 + V96: """
503 +SELECT
504 + pg_xlog_location_diff(
505 + CASE pg_is_in_recovery()
506 + WHEN true THEN pg_last_xlog_receive_location()
507 + ELSE pg_current_xlog_location()
508 + END,
509 + '0/0') as wal_writes ;
510 +""",
511 }
512
513
329 -QUERY_STATS = {
330 - QUERIES['DATABASE']: METRICS['DATABASE'],
331 - QUERIES['BACKENDS']: METRICS['BACKENDS'],
332 - QUERIES['LOCKS']: METRICS['LOCKS']
333 -}
514 +def query_factory(name, version=NO_VERSION):
515 + if name == BACKENDS:
516 + return QUERY_BACKEND[DEFAULT]
517 + elif name == TABLE_STATS:
518 + return QUERY_TABLE_STATS[DEFAULT]
519 + elif name == INDEX_STATS:
520 + return QUERY_INDEX_STATS[DEFAULT]
521 + elif name == DATABASE:
522 + return QUERY_DATABASE[DEFAULT]
523 + elif name == BGWRITER:
524 + return QUERY_BGWRITER[DEFAULT]
525 + elif name == LOCKS:
526 + return QUERY_LOCKS[DEFAULT]
527 + elif name == DATABASES:
528 + return QUERY_DATABASES[DEFAULT]
529 + elif name == STANDBY:
530 + return QUERY_STANDBY[DEFAULT]
531 + elif name == REPLICATION_SLOT:
532 + return QUERY_REPLICATION_SLOT[DEFAULT]
533 + elif name == IF_SUPERUSER:
534 + return QUERY_SUPERUSER[DEFAULT]
535 + elif name == SERVER_VERSION:
536 + return QUERY_SHOW_VERSION[DEFAULT]
537 + elif name == AUTOVACUUM:
538 + return QUERY_AUTOVACUUM[DEFAULT]
539 + elif name == WAL:
540 + if version < 100000:
541 + return QUERY_WAL[V96]
542 + return QUERY_WAL[DEFAULT]
543 + elif name == ARCHIVE:
544 + if version < 100000:
545 + return QUERY_ARCHIVE[V96]
546 + return QUERY_ARCHIVE[DEFAULT]
547 + elif name == STANDBY_DELTA:
548 + if version < 100000:
549 + return QUERY_STANDBY_DELTA[V96]
550 + return QUERY_STANDBY_DELTA[DEFAULT]
551 + elif name == REPSLOT_FILES:
552 + if version < 110000:
553 + return QUERY_REPSLOT_FILES[V10]
554 + return QUERY_REPSLOT_FILES[DEFAULT]
555 + elif name == DIFF_LSN:
556 + if version < 100000:
557 + return QUERY_DIFF_LSN[V96]
558 + return QUERY_DIFF_LSN[DEFAULT]
559 +
560 + raise ValueError('unknown query')
561 +
562
563 ORDER = [
564 'db_stat_temp_files',
@@ -553,151 +781,112 @@ CHARTS = {
781 class Service(SimpleService):
782 def __init__(self, configuration=None, name=None):
783 SimpleService.__init__(self, configuration=configuration, name=name)
556 - self.order = ORDER[:]
784 + self.order = list(ORDER)
785 self.definitions = deepcopy(CHARTS)
558 - self.table_stats = configuration.pop('table_stats', False)
559 - self.index_stats = configuration.pop('index_stats', False)
560 - self.database_poll = configuration.pop('database_poll', None)
786 +
787 + self.do_table_stats = configuration.pop('table_stats', False)
788 + self.do_index_stats = configuration.pop('index_stats', False)
789 + self.databases_to_poll = configuration.pop('database_poll', None)
790 self.configuration = configuration
562 - self.connection = False
791 +
792 + self.conn = None
793 self.server_version = None
564 - self.data = dict()
565 - self.locks_zeroed = dict()
794 + self.is_superuser = False
795 + self.alive = False
796 +
797 self.databases = list()
798 self.secondaries = list()
799 self.replication_slots = list()
569 - self.queries = QUERY_STATS.copy()
570 -
571 - def _connect(self):
572 - params = dict(user='postgres',
573 - database=None,
574 - password=None,
575 - host=None,
576 - port=5432)
577 - params.update(self.configuration)
578 -
579 - if not self.connection:
580 - try:
581 - self.connection = psycopg2.connect(**params)
582 - self.connection.set_isolation_level(extensions.ISOLATION_LEVEL_AUTOCOMMIT)
583 - self.connection.set_session(readonly=True)
584 - except OperationalError as error:
585 - return False, str(error)
586 - return True, True
800 +
801 + self.queries = dict()
802 +
803 + self.data = dict()
804 +
805 + def reconnect(self):
806 + return self.connect()
807 +
808 + def connect(self):
809 + if self.conn:
810 + self.conn.close()
811 + self.conn = None
812 +
813 + try:
814 + params = dict(user='postgres',
815 + database=None,
816 + password=None,
817 + host=None,
818 + port=5432)
819 + params.update(self.configuration)
820 +
821 + self.conn = psycopg2.connect(**params)
822 + self.conn.set_isolation_level(extensions.ISOLATION_LEVEL_AUTOCOMMIT)
823 + self.conn.set_session(readonly=True)
824 + except OperationalError as error:
825 + self.error(error)
826 + self.alive = False
827 + else:
828 + self.alive = True
829 +
830 + return self.alive
831
832 def check(self):
833 if not PSYCOPG2:
590 - self.error('\'python-psycopg2\' module is needed to use postgres.chart.py')
834 + self.error("'python-psycopg2' package is needed to use postgres module")
835 return False
592 - result, error = self._connect()
593 - if not result:
594 - conf = dict((k, (lambda k, v: v if k != 'password' else '*****')(k, v))
595 - for k, v in self.configuration.items())
596 - self.error('Failed to connect to %s. Error: %s' % (str(conf), error))
836 +
837 + if not self.connect():
838 + self.error('failed to connect to {0}'.format(hide_password(self.configuration)))
839 return False
840 +
841 try:
599 - cursor = self.connection.cursor()
600 - self.databases = discover_databases_(cursor, QUERIES['FIND_DATABASES'])
601 - is_superuser = check_if_superuser_(cursor, QUERIES['IF_SUPERUSER'])
602 - self.secondaries = discover_secondaries_(cursor, QUERIES['FIND_STANDBY'])
603 - self.server_version = detect_server_version(cursor, QUERIES['DETECT_SERVER_VERSION'])
604 - if self.server_version >= 94000:
605 - self.replication_slots = discover_replication_slots_(cursor, QUERIES['FIND_REPLICATION_SLOT'])
606 - cursor.close()
607 -
608 - if self.database_poll and isinstance(self.database_poll, str):
609 - self.databases = [dbase for dbase in self.databases if dbase in self.database_poll.split()] \
610 - or self.databases
611 -
612 - self.locks_zeroed = populate_lock_types(self.databases)
613 - self.add_additional_queries_(is_superuser)
614 - self.create_dynamic_charts_()
615 - return True
842 + self.check_queries()
843 except Exception as error:
617 - self.error(str(error))
844 + self.error(error)
845 return False
846
620 - def add_additional_queries_(self, is_superuser):
847 + self.populate_queries()
848 + self.create_dynamic_charts()
849
622 - if self.server_version >= 100000:
623 - wal = 'wal'
624 - lsn = 'lsn'
625 - else:
626 - wal = 'xlog'
627 - lsn = 'location'
628 - self.queries[QUERIES['BGWRITER']] = METRICS['BGWRITER']
629 - self.queries[QUERIES['DIFF_LSN'].format(wal, lsn)] = METRICS['WAL_WRITES']
630 - self.queries[QUERIES['STANDBY_DELTA'].format(wal, lsn)] = METRICS['STANDBY_DELTA']
631 -
632 - if self.index_stats:
633 - self.queries[QUERIES['INDEX_STATS']] = METRICS['INDEX_STATS']
634 - if self.table_stats:
635 - self.queries[QUERIES['TABLE_STATS']] = METRICS['TABLE_STATS']
636 - if is_superuser:
637 - self.queries[QUERIES['ARCHIVE'].format(wal)] = METRICS['ARCHIVE']
638 - if self.server_version >= 90400:
639 - self.queries[QUERIES['WAL'].format(wal, lsn)] = METRICS['WAL']
640 - if self.server_version >= 100000:
641 - self.queries[QUERIES['REPSLOT_FILES']] = METRICS['REPSLOT_FILES']
642 - if self.server_version >= 90400:
643 - self.queries[QUERIES['AUTOVACUUM']] = METRICS['AUTOVACUUM']
850 + return True
851
645 - def create_dynamic_charts_(self):
852 + def get_data(self):
853 + if not self.alive and not self.reconnect():
854 + return None
855
647 - for database_name in self.databases[::-1]:
648 - self.definitions['database_size']['lines'].append(
649 - [database_name + '_size', database_name, 'absolute', 1, 1024 * 1024])
650 - for chart_name in [name for name in self.order if name.startswith('db_stat')]:
651 - add_database_stat_chart_(order=self.order, definitions=self.definitions,
652 - name=chart_name, database_name=database_name)
856 + try:
857 + cursor = self.conn.cursor(cursor_factory=DictCursor)
858
654 - add_database_lock_chart_(order=self.order, definitions=self.definitions, database_name=database_name)
859 + self.data.update(zero_lock_types(self.databases))
860
656 - for application_name in self.secondaries[::-1]:
657 - add_replication_delta_chart_(
658 - order=self.order,
659 - definitions=self.definitions,
660 - name='standby_delta',
661 - application_name=application_name)
861 + for query, metrics in self.queries.items():
862 + self.query_stats(cursor, query, metrics)
863
663 - for slot_name in self.replication_slots[::-1]:
664 - add_replication_slot_chart_(
665 - order=self.order,
666 - definitions=self.definitions,
667 - name='replication_slot',
668 - slot_name=slot_name)
669 -
670 - def _get_data(self):
671 - result, _ = self._connect()
672 - if result:
673 - cursor = self.connection.cursor(cursor_factory=DictCursor)
674 - try:
675 - self.data.update(self.locks_zeroed)
676 - for query, metrics in self.queries.items():
677 - self.query_stats_(cursor, query, metrics)
678 -
679 - except OperationalError:
680 - self.connection = False
681 - cursor.close()
682 - return None
683 - else:
684 - cursor.close()
685 - return self.data
686 - else:
864 + except OperationalError:
865 + self.alive = False
866 return None
867
689 - def query_stats_(self, cursor, query, metrics):
868 + cursor.close()
869 +
870 + return self.data
871 +
872 + def query_stats(self, cursor, query, metrics):
873 cursor.execute(query, dict(databases=tuple(self.databases)))
874 +
875 for row in cursor:
876 for metric in metrics:
877 + # databases
878 if 'database_name' in row:
879 dimension_id = '_'.join([row['database_name'], metric])
880 + # secondaries
881 elif 'application_name' in row:
882 dimension_id = '_'.join([row['application_name'], metric])
883 + # replication slots
884 elif 'slot_name' in row:
885 dimension_id = '_'.join([row['slot_name'], metric])
886 + # other
887 else:
888 dimension_id = metric
889 +
890 if metric in row:
891 if row[metric] is not None:
892 self.data[dimension_id] = int(row[metric])
@@ -705,35 +894,105 @@ class Service(SimpleService):
894 if metric == row['mode']:
895 self.data[dimension_id] = row['locks_count']
896
897 + def check_queries(self):
898 + cursor = self.conn.cursor()
899
709 -def discover_databases_(cursor, query):
710 - cursor.execute(query)
711 - result = list()
712 - for db in [database[0] for database in cursor]:
713 - if db not in result:
714 - result.append(db)
715 - return result
900 + self.server_version = detect_server_version(cursor, query_factory(SERVER_VERSION))
901 + self.debug('server version: {0}'.format(self.server_version))
902
903 + self.is_superuser = check_if_superuser(cursor, query_factory(IF_SUPERUSER))
904 + self.debug('superuser: {0}'.format(self.is_superuser))
905
718 -def discover_secondaries_(cursor, query):
719 - cursor.execute(query)
720 - result = list()
721 - for sc in [standby[0] for standby in cursor]:
722 - if sc not in result:
723 - result.append(sc)
724 - return result
906 + self.databases = discover(cursor, query_factory(DATABASES))
907 + self.debug('discovered databases {0}'.format(self.databases))
908 + if self.databases_to_poll:
909 + to_poll = self.databases_to_poll.split()
910 + self.databases = [db for db in self.databases if db in to_poll] or self.databases
911 +
912 + self.secondaries = discover(cursor, query_factory(STANDBY))
913 + self.debug('discovered secondaries: {0}'.format(self.secondaries))
914 +
915 + if self.server_version >= 94000:
916 + self.replication_slots = discover(cursor, query_factory(REPLICATION_SLOT))
917 + self.debug('discovered replication slots: {0}'.format(self.replication_slots))
918 +
919 + cursor.close()
920 +
921 + def populate_queries(self):
922 + self.queries[query_factory(DATABASE)] = METRICS[DATABASE]
923 + self.queries[query_factory(BACKENDS)] = METRICS[BACKENDS]
924 + self.queries[query_factory(LOCKS)] = METRICS[LOCKS]
925 + self.queries[query_factory(BGWRITER)] = METRICS[BGWRITER]
926 + self.queries[query_factory(DIFF_LSN, self.server_version)] = METRICS[WAL_WRITES]
927 + self.queries[query_factory(STANDBY_DELTA, self.server_version)] = METRICS[STANDBY_DELTA]
928 +
929 + if self.do_index_stats:
930 + self.queries[query_factory(INDEX_STATS)] = METRICS[INDEX_STATS]
931 + if self.do_table_stats:
932 + self.queries[query_factory(TABLE_STATS)] = METRICS[TABLE_STATS]
933
934 + if self.is_superuser:
935 + self.queries[query_factory(ARCHIVE, self.server_version)] = METRICS[ARCHIVE]
936
727 -def discover_replication_slots_(cursor, query):
937 + if self.server_version >= 90400:
938 + self.queries[query_factory(WAL, self.server_version)] = METRICS[WAL]
939 +
940 + if self.server_version >= 100000:
941 + self.queries[query_factory(REPSLOT_FILES, self.server_version)] = METRICS[REPSLOT_FILES]
942 +
943 + if self.server_version >= 90400:
944 + self.queries[query_factory(AUTOVACUUM)] = METRICS[AUTOVACUUM]
945 +
946 + def create_dynamic_charts(self):
947 + for database_name in self.databases[::-1]:
948 + dim = [
949 + database_name + '_size',
950 + database_name,
951 + 'absolute',
952 + 1,
953 + 1024 * 1024,
954 + ]
955 + self.definitions['database_size']['lines'].append(dim)
956 + for chart_name in [name for name in self.order if name.startswith('db_stat')]:
957 + add_database_stat_chart(
958 + order=self.order,
959 + definitions=self.definitions,
960 + name=chart_name,
961 + database_name=database_name,
962 + )
963 + add_database_lock_chart(
964 + order=self.order,
965 + definitions=self.definitions,
966 + database_name=database_name,
967 + )
968 +
969 + for application_name in self.secondaries[::-1]:
970 + add_replication_delta_chart(
971 + order=self.order,
972 + definitions=self.definitions,
973 + name='standby_delta',
974 + application_name=application_name,
975 + )
976 +
977 + for slot_name in self.replication_slots[::-1]:
978 + add_replication_slot_chart(
979 + order=self.order,
980 + definitions=self.definitions,
981 + name='replication_slot',
982 + slot_name=slot_name,
983 + )
984 +
985 +
986 +def discover(cursor, query):
987 cursor.execute(query)
988 result = list()
730 - for slot in [replication_slot[0] for replication_slot in cursor]:
731 - if slot not in result:
732 - result.append(slot)
989 + for v in [value[0] for value in cursor]:
990 + if v not in result:
991 + result.append(v)
992 return result
993
994
736 -def check_if_superuser_(cursor, query):
995 +def check_if_superuser(cursor, query):
996 cursor.execute(query)
997 return cursor.fetchone()[0]
998
@@ -743,7 +1002,7 @@ def detect_server_version(cursor, query):
1002 return int(cursor.fetchone()[0])
1003
1004
746 -def populate_lock_types(databases):
1005 +def zero_lock_types(databases):
1006 result = dict()
1007 for database in databases:
1008 for lock_type in METRICS['LOCKS']:
@@ -753,7 +1012,11 @@ def populate_lock_types(databases):
1012 return result
1013
1014
756 -def add_database_lock_chart_(order, definitions, database_name):
1015 +def hide_password(config):
1016 + return dict((k, v if k != 'password' else '*****') for k, v in config.items())
1017 +
1018 +
1019 +def add_database_lock_chart(order, definitions, database_name):
1020 def create_lines(database):
1021 result = list()
1022 for lock_type in METRICS['LOCKS']:
@@ -770,7 +1033,7 @@ def add_database_lock_chart_(order, definitions, database_name):
1033 }
1034
1035
773 -def add_database_stat_chart_(order, definitions, name, database_name):
1036 +def add_database_stat_chart(order, definitions, name, database_name):
1037 def create_lines(database, lines):
1038 result = list()
1039 for line in lines:
@@ -787,7 +1050,7 @@ def add_database_stat_chart_(order, definitions, name, database_name):
1050 'lines': create_lines(database_name, chart_template['lines'])}
1051
1052
790 -def add_replication_delta_chart_(order, definitions, name, application_name):
1053 +def add_replication_delta_chart(order, definitions, name, application_name):
1054 def create_lines(standby, lines):
1055 result = list()
1056 for line in lines:
@@ -805,7 +1068,7 @@ def add_replication_delta_chart_(order, definitions, name, application_name):
1068 'lines': create_lines(application_name, chart_template['lines'])}
1069
1070
808 -def add_replication_slot_chart_(order, definitions, name, slot_name):
1071 +def add_replication_slot_chart(order, definitions, name, slot_name):
1072 def create_lines(slot, lines):
1073 result = list()
1074 for line in lines: