| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143 |
- #!/usr/bin/env /usr/bin/python3
- # -*- coding:utf-8 -*-
- """
- CDC 链路监控:源库复制槽 + Flink 采集作业存活。
- DS 每 10 分钟调一次,无状态:每次独立判断。连续收到就是还没恢复,不再收到就是已恢复。
- 检查两类:
- 1. 复制槽——源库**必须是主库**,从库上 pg_current_wal_lsn() 报
- "recovery is in progress",且 pg_replication_slots 只能看到本节点的槽
- - cdc_ 前缀的槽 active = false → 采集作业可能已停
- - 任意槽滞后超阈值 → 消费跟不上
- 2. Flink session 集群上 CDC 作业缺失,或状态非 RUNNING → 作业挂了
- 连库失败 / Flink 查询失败本身也计入告警,不静默跳过——查不到状态和状态异常同样需要人介入。
- Alerter 在开跑前构造:配置缺失或 key 未替换立刻暴露,不留到真出事才发现告警链路是哑的。
- 退出码:检测到异常仍返回 0(DS 任务不标红,异常靠企微通知);
- 只有告警推送失败或脚本自身异常才非 0——告警发不出去必须让调度看见。
- CLI:
- python3 bin/cdc-monitor.py [-ds postgresql/prd-poyee-aliyun-cdc] [-channel default]
- [-jm cdhmaster02:8081] [-jobs st-cdc-large,st-cdc-small]
- [-lag-gb 5]
- 参数:
- -ds 主库 datasource ref,解析项目同级 ../datasource/{ds}.ini
- -channel 告警通道名,对应 conf/alerter.ini [channels] 段
- -jm Flink JobManager REST 地址
- -jobs CDC 作业名,逗号分隔,逐个在 Flink 作业列表里找
- -lag-gb 槽滞后告警阈值(GB),正常应为 0
- """
- import argparse
- import json
- import os
- import sys
- import urllib.request
- from datetime import datetime
- project_root = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
- sys.path.append(project_root)
- from dw_base.alerter.alerter import Alerter
- from dw_base.io.db import postgresql as pgdb
- # 不加 slot_name 过滤:顺带覆盖财务库那条 PG→PG 的槽
- SLOT_SQL = ('SELECT slot_name, active, '
- 'pg_wal_lsn_diff(pg_current_wal_lsn(), confirmed_flush_lsn) AS lag_bytes '
- 'FROM pg_replication_slots')
- _FLINK_TIMEOUT_SEC = 20
- def check_slots(ds_ref, lag_crit_bytes):
- """查复制槽,返回 [(项, 说明)]。连库失败本身也算一项。"""
- alerts = []
- conn = None
- try:
- conn = pgdb.connect(ds_ref)
- for slot_name, active, lag_bytes in pgdb.query(conn, SLOT_SQL):
- # active 只对 cdc_ 槽判:PG→PG 那条 apply 卡住时会反复重连,会误报
- if slot_name.startswith('cdc_') and not active:
- alerts.append(('槽 ' + slot_name, '空闲,采集作业可能已停'))
- lag = int(lag_bytes or 0)
- if lag > lag_crit_bytes:
- alerts.append(('槽 ' + slot_name, '滞后 %.1f GB' % (lag / 1024.0 ** 3)))
- except Exception as e:
- alerts.append(('主库', '复制槽查询失败:%s' % e))
- finally:
- if conn is not None:
- conn.close()
- return alerts
- def check_flink(jm, jobs):
- """查 Flink session 集群的作业状态,返回 [(项, 说明)]。
- 必须按 state 过滤:作业取消/失败后仍留在 /jobs/overview 列表里,
- 只判名字存在会漏报——重启死循环时状态是 RESTARTING,名字一直在。
- """
- try:
- with urllib.request.urlopen('http://%s/jobs/overview' % jm,
- timeout=_FLINK_TIMEOUT_SEC) as resp:
- overview = json.loads(resp.read().decode('utf-8'))
- except Exception as e:
- return [('Flink 集群', '查询失败:%s' % e)]
- # 同名作业会留下历史记录,RUNNING 优先
- states = {}
- for j in overview.get('jobs', []):
- if j['name'] not in states or j['state'] == 'RUNNING':
- states[j['name']] = j['state']
- alerts = []
- for job in jobs:
- state = states.get(job)
- if state is None:
- alerts.append(('作业 ' + job, '不在 Flink 集群上'))
- elif state != 'RUNNING':
- alerts.append(('作业 ' + job, '状态 %s' % state))
- return alerts
- def render(alerts, now=None):
- """拼企微 markdown。
- 检测时间用服务器本地时区(与 DS 告警消息同格式,同群里看着一致)。
- now 留给单测注入固定值,默认取当前时间。
- """
- lines = ['### <font color="warning">CDC 监控告警</font>',
- '> 检测时间:%s' % (now or datetime.now()).strftime('%Y-%m-%d %H:%M:%S')]
- lines.extend('> %s:<font color="warning">%s</font>' % (k, v) for k, v in alerts)
- return '\n'.join(lines)
- def main():
- parser = argparse.ArgumentParser(description='CDC 链路监控(复制槽 + Flink 作业存活)')
- parser.add_argument('-ds', default='postgresql/prd-poyee-aliyun-cdc',
- help='主库 datasource ref(默认 postgresql/prd-poyee-aliyun-cdc)')
- parser.add_argument('-channel', default='default',
- help='告警通道名,对应 conf/alerter.ini [channels](默认 default)')
- parser.add_argument('-jm', default='cdhmaster02:8081',
- help='Flink JobManager REST 地址(默认 cdhmaster02:8081)')
- parser.add_argument('-jobs', default='st-cdc-large,st-cdc-small',
- help='CDC 作业名,逗号分隔(默认 st-cdc-large,st-cdc-small)')
- parser.add_argument('-lag-gb', type=float, default=5, dest='lag_gb',
- help='槽滞后告警阈值 GB(默认 5,正常应为 0)')
- args = parser.parse_args()
- alerter = Alerter(channel=args.channel)
- jobs = [j.strip() for j in args.jobs.split(',') if j.strip()]
- alerts = check_slots(args.ds, int(args.lag_gb * 1024 ** 3)) + check_flink(args.jm, jobs)
- for key, desc in alerts:
- print('%s: %s' % (key, desc)) # 进 DS 任务日志
- if alerts:
- alerter.send_markdown(render(alerts))
- if __name__ == '__main__':
- main()
|