master
go 91 lines 2.47 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package postgres
4
5 import (
6 "fmt"
7 )
8
9 func (c *Collector) doQueryReplicationMetrics() error {
10 if err := c.doQueryReplStandbyAppWALDelta(); err != nil {
11 return fmt.Errorf("querying replication standby app wal delta error: %v", err)
12 }
13
14 if c.pgVersion >= pgVersion10 {
15 if err := c.doQueryReplStandbyAppWALLag(); err != nil {
16 return fmt.Errorf("querying replication standby app wal lag error: %v", err)
17 }
18 }
19
20 if c.pgVersion >= pgVersion10 && c.canQueryReplicationSlotFiles() {
21 if err := c.doQueryReplSlotFiles(); err != nil {
22 return fmt.Errorf("querying replication slot files error: %v", err)
23 }
24 }
25
26 return nil
27 }
28
29 func (c *Collector) doQueryReplStandbyAppWALDelta() error {
30 q := queryReplicationStandbyAppDelta(c.pgVersion)
31
32 var app string
33 return c.doQuery(q, func(column, value string, _ bool) {
34 switch column {
35 case "application_name":
36 app = value
37 c.getReplAppMetrics(app).updated = true
38 default:
39 // TODO: delta calculation was changed in https://github.com/netdata/netdata/go/plugins/plugin/go.d/pull/1039
40 // - 'replay_delta' (probably other deltas too?) can be negative
41 // - Also, WAL delta != WAL lag after that PR
42 v := max(parseInt(value), 0)
43 switch column {
44 case "sent_delta":
45 c.getReplAppMetrics(app).walSentDelta += v
46 case "write_delta":
47 c.getReplAppMetrics(app).walWriteDelta += v
48 case "flush_delta":
49 c.getReplAppMetrics(app).walFlushDelta += v
50 case "replay_delta":
51 c.getReplAppMetrics(app).walReplayDelta += v
52 }
53 }
54 })
55 }
56
57 func (c *Collector) doQueryReplStandbyAppWALLag() error {
58 q := queryReplicationStandbyAppLag()
59
60 var app string
61 return c.doQuery(q, func(column, value string, _ bool) {
62 switch column {
63 case "application_name":
64 app = value
65 c.getReplAppMetrics(app).updated = true
66 case "write_lag":
67 c.getReplAppMetrics(app).walWriteLag += parseInt(value)
68 case "flush_lag":
69 c.getReplAppMetrics(app).walFlushLag += parseInt(value)
70 case "replay_lag":
71 c.getReplAppMetrics(app).walReplayLag += parseInt(value)
72 }
73 })
74 }
75
76 func (c *Collector) doQueryReplSlotFiles() error {
77 q := queryReplicationSlotFiles(c.pgVersion)
78
79 var slot string
80 return c.doQuery(q, func(column, value string, _ bool) {
81 switch column {
82 case "slot_name":
83 slot = value
84 c.getReplSlotMetrics(slot).updated = true
85 case "replslot_wal_keep":
86 c.getReplSlotMetrics(slot).walKeep += parseInt(value)
87 case "replslot_files":
88 c.getReplSlotMetrics(slot).files += parseInt(value)
89 }
90 })
91 }