cdc-monitor.py 8.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200
  1. #!/usr/bin/env /usr/bin/python3
  2. # -*- coding:utf-8 -*-
  3. """
  4. CDC 链路监控:源库复制槽 + Flink 采集作业存活。
  5. 两种模式,同一份检查逻辑与配置:
  6. - 缺省(DS 每 10 分钟调一次):只在检测到异常时推企微。无状态,每次独立判断,
  7. 连续收到就是还没恢复,不再收到就是已恢复。
  8. - -report(DS 手动触发,上线核对用):无条件推当前状态,正常也发。
  9. 检查两类:
  10. 1. 复制槽——源库**必须是主库**,从库上 pg_current_wal_lsn() 报
  11. "recovery is in progress",且 pg_replication_slots 只能看到本节点的槽
  12. - cdc_ 前缀的槽 active = false → 采集作业可能已停
  13. - 任意槽滞后超阈值 → 消费跟不上
  14. 2. Flink session 集群上 CDC 作业缺失,或状态非 RUNNING → 作业挂了
  15. 连库失败 / Flink 查询失败本身也计入告警,不静默跳过——查不到状态和状态异常同样需要人介入。
  16. 两个 check 都返回 (facts, alerts):facts 是当前值,供 -report 无条件列出;
  17. alerts 是判定结果,供缺省模式决定发不发。判定口径两种模式共用一份。
  18. Alerter 在开跑前构造:配置缺失或 key 未替换立刻暴露,不留到真出事才发现告警链路是哑的。
  19. 退出码:检测到异常仍返回 0(DS 任务不标红,异常靠企微通知);
  20. 只有告警推送失败或脚本自身异常才非 0——告警发不出去必须让调度看见。
  21. CLI:
  22. python3 bin/cdc-monitor.py [-ds postgresql/prd-poyee-aliyun-cdc] [-channel default]
  23. [-jm cdhmaster02:8081] [-jobs st-cdc-large,st-cdc-small]
  24. [-lag-gb 5] [-report]
  25. 参数:
  26. -ds 主库 datasource ref,解析项目同级 ../datasource/{ds}.ini
  27. -channel 告警通道名,对应 conf/alerter.ini [channels] 段
  28. -jm Flink JobManager REST 地址
  29. -jobs CDC 作业名,逗号分隔,逐个在 Flink 作业列表里找
  30. -lag-gb 槽滞后告警阈值(GB),正常应为 0
  31. -report 无条件推当前状态,正常也发;不带则只在异常时告警
  32. """
  33. import argparse
  34. import json
  35. import os
  36. import sys
  37. import urllib.request
  38. from datetime import datetime
  39. project_root = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
  40. sys.path.append(project_root)
  41. from dw_base.alerter.alerter import Alerter
  42. from dw_base.io.db import postgresql as pgdb
  43. # 不加 slot_name 过滤:顺带覆盖财务库那条 PG→PG 的槽
  44. SLOT_SQL = ('SELECT slot_name, active, '
  45. 'pg_wal_lsn_diff(pg_current_wal_lsn(), confirmed_flush_lsn) AS lag_bytes '
  46. 'FROM pg_replication_slots')
  47. _FLINK_TIMEOUT_SEC = 20
  48. def check_slots(ds_ref, lag_crit_bytes):
  49. """查复制槽,返回 (facts, alerts)。
  50. facts: [(槽名, active, 滞后字节)],连库失败时为空。
  51. alerts: [(项, 说明)],连库失败本身也算一项。
  52. """
  53. facts = []
  54. alerts = []
  55. conn = None
  56. try:
  57. conn = pgdb.connect(ds_ref)
  58. for slot_name, active, lag_bytes in pgdb.query(conn, SLOT_SQL):
  59. lag = int(lag_bytes or 0)
  60. facts.append((slot_name, active, lag))
  61. # active 只对 cdc_ 槽判:PG→PG 那条 apply 卡住时会反复重连,会误报
  62. if slot_name.startswith('cdc_') and not active:
  63. alerts.append(('槽 ' + slot_name, '空闲,采集作业可能已停'))
  64. if lag > lag_crit_bytes:
  65. alerts.append(('槽 ' + slot_name, '滞后 %.1f GB' % (lag / 1024.0 ** 3)))
  66. except Exception as e:
  67. alerts.append(('主库', '复制槽查询失败:%s' % e))
  68. finally:
  69. if conn is not None:
  70. conn.close()
  71. return facts, alerts
  72. def check_flink(jm, jobs):
  73. """查 Flink session 集群的作业状态,返回 (facts, alerts)。
  74. facts: [(作业名, state)],按 -jobs 顺序;集群上没有的 state 为 None。
  75. 只报 -jobs 点名的作业——这是 CDC 链路状态,不是集群状态。
  76. 必须按 state 过滤:作业取消/失败后仍留在 /jobs/overview 列表里,
  77. 只判名字存在会漏报——重启死循环时状态是 RESTARTING,名字一直在。
  78. """
  79. try:
  80. with urllib.request.urlopen('http://%s/jobs/overview' % jm,
  81. timeout=_FLINK_TIMEOUT_SEC) as resp:
  82. overview = json.loads(resp.read().decode('utf-8'))
  83. except Exception as e:
  84. return [], [('Flink 集群', '查询失败:%s' % e)]
  85. # 同名作业会留下历史记录,RUNNING 优先
  86. states = {}
  87. for j in overview.get('jobs', []):
  88. if j['name'] not in states or j['state'] == 'RUNNING':
  89. states[j['name']] = j['state']
  90. facts = []
  91. alerts = []
  92. for job in jobs:
  93. state = states.get(job)
  94. facts.append((job, state))
  95. if state is None:
  96. alerts.append(('作业 ' + job, '不在 Flink 集群上'))
  97. elif state != 'RUNNING':
  98. alerts.append(('作业 ' + job, '状态 %s' % state))
  99. return facts, alerts
  100. def _line(key, value, is_bad):
  101. """一行正文:字段名默认色,只给值上色(对齐 kb/41 §3.5)。"""
  102. return '> %s:<font color="%s">%s</font>' % (key, 'warning' if is_bad else 'info', value)
  103. def _head(title, color, now):
  104. """标题 + 检测时间。时间用服务器本地时区(与 DS 告警消息同格式,同群里看着一致)。"""
  105. return ['### <font color="%s">%s</font>' % (color, title),
  106. '> 检测时间:%s' % (now or datetime.now()).strftime('%Y-%m-%d %H:%M:%S')]
  107. def render(alerts, now=None):
  108. """拼告警 markdown:只列异常项。now 留给单测注入固定值。"""
  109. lines = _head('CDC 监控告警', 'warning', now)
  110. lines.extend(_line(k, v, True) for k, v in alerts)
  111. return '\n'.join(lines)
  112. def render_status(slot_facts, flink_facts, alerts, now=None):
  113. """拼状态报告 markdown:无条件列出每一项当前值。
  114. 异常项的值上 warning,正常项上 info;标题按有无异常切色。
  115. 查询失败时没有事实可列,把失败原因单独补在末尾,否则报告会是空的。
  116. """
  117. bad = set(k for k, _ in alerts)
  118. lines = _head('CDC 链路状态', 'warning' if alerts else 'info', now)
  119. fact_keys = set()
  120. for slot_name, active, lag in slot_facts:
  121. key = '槽 ' + slot_name
  122. fact_keys.add(key)
  123. lines.append(_line(key, '%s,滞后 %.1f GB' % (
  124. '活跃' if active else '空闲', lag / 1024.0 ** 3), key in bad))
  125. for job, state in flink_facts:
  126. key = '作业 ' + job
  127. fact_keys.add(key)
  128. lines.append(_line(key, state or '不在集群上', key in bad))
  129. lines.extend(_line(k, v, True) for k, v in alerts if k not in fact_keys)
  130. return '\n'.join(lines)
  131. def main():
  132. parser = argparse.ArgumentParser(description='CDC 链路监控(复制槽 + Flink 作业存活)')
  133. parser.add_argument('-ds', default='postgresql/prd-poyee-aliyun-cdc',
  134. help='主库 datasource ref(默认 postgresql/prd-poyee-aliyun-cdc)')
  135. parser.add_argument('-channel', default='default',
  136. help='告警通道名,对应 conf/alerter.ini [channels](默认 default)')
  137. parser.add_argument('-jm', default='cdhmaster02:8081',
  138. help='Flink JobManager REST 地址(默认 cdhmaster02:8081)')
  139. parser.add_argument('-jobs', default='st-cdc-large,st-cdc-small',
  140. help='CDC 作业名,逗号分隔(默认 st-cdc-large,st-cdc-small)')
  141. parser.add_argument('-lag-gb', type=float, default=5, dest='lag_gb',
  142. help='槽滞后告警阈值 GB(默认 5,正常应为 0)')
  143. parser.add_argument('-report', action='store_true',
  144. help='无条件推当前状态,正常也发;不带则只在异常时告警')
  145. args = parser.parse_args()
  146. alerter = Alerter(channel=args.channel)
  147. jobs = [j.strip() for j in args.jobs.split(',') if j.strip()]
  148. slot_facts, slot_alerts = check_slots(args.ds, int(args.lag_gb * 1024 ** 3))
  149. flink_facts, flink_alerts = check_flink(args.jm, jobs)
  150. alerts = slot_alerts + flink_alerts
  151. for key, desc in alerts:
  152. print('%s: %s' % (key, desc)) # 进 DS 任务日志
  153. if args.report:
  154. content = render_status(slot_facts, flink_facts, alerts)
  155. print(content) # 发了什么,DS 任务日志里留一份
  156. alerter.send_markdown(content)
  157. elif alerts:
  158. alerter.send_markdown(render(alerts))
  159. if __name__ == '__main__':
  160. main()