#!/usr/bin/env /usr/bin/python3 # -*- coding:utf-8 -*- """ CDC 链路监控:源库复制槽 + YARN 采集作业存活。 DS 每 5 分钟调一次,无状态:每次独立判断。连续收到就是还没恢复,不再收到就是已恢复。 检查两类: 1. 复制槽——源库**必须是主库**,从库上 pg_current_wal_lsn() 报 "recovery is in progress",且 pg_replication_slots 只能看到本节点的槽 - cdc_ 前缀的槽 active = false → 采集作业可能已停 - 任意槽滞后超阈值 → 消费跟不上 2. YARN RUNNING 列表里缺 CDC 作业 → 作业挂了 连库失败 / YARN 查询失败本身也计入告警,不静默跳过——查不到状态和状态异常同样需要人介入。 Alerter 在开跑前构造:配置缺失或 key 未替换立刻暴露,不留到真出事才发现告警链路是哑的。 退出码:检测到异常仍返回 0(DS 任务不标红,异常靠企微通知); 只有告警推送失败或脚本自身异常才非 0——告警发不出去必须让调度看见。 CLI: python3 bin/cdc-monitor.py [-ds postgresql/prd-poyee-aliyun-cdc] [-channel default] [-jobs st-cdc-large,st-cdc-small] [-lag-gb 5] 参数: -ds 主库 datasource ref,解析项目同级 ../datasource/{ds}.ini -channel 告警通道名,对应 conf/alerter.ini [channels] 段 -jobs CDC 作业名,逗号分隔,逐个在 YARN RUNNING 列表里找 -lag-gb 槽滞后告警阈值(GB),正常应为 0 """ import argparse import os import subprocess import sys 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') _YARN_TIMEOUT_SEC = 60 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_yarn(jobs): """查 YARN RUNNING 列表,返回 [(项, 说明)]。""" try: running = subprocess.check_output( ['yarn', 'application', '-list', '-appStates', 'RUNNING'], stderr=subprocess.STDOUT, timeout=_YARN_TIMEOUT_SEC).decode('utf-8') except Exception as e: return [('YARN', '查询失败:%s' % e)] return [('作业 ' + job, '不在 YARN RUNNING 列表') for job in jobs if job not in running] def render(alerts, now=None): """拼企微 markdown。 检测时间用服务器本地时区(与 DS 告警消息同格式,同群里看着一致)。 now 留给单测注入固定值,默认取当前时间。 """ lines = ['### CDC 监控告警', '> 检测时间:%s' % (now or datetime.now()).strftime('%Y-%m-%d %H:%M:%S')] lines.extend('> %s:%s' % (k, v) for k, v in alerts) return '\n'.join(lines) def main(): parser = argparse.ArgumentParser(description='CDC 链路监控(复制槽 + YARN 作业存活)') 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('-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_yarn(jobs) for key, desc in alerts: print('%s: %s' % (key, desc)) # 进 DS 任务日志 if alerts: alerter.send_markdown(render(alerts)) if __name__ == '__main__': main()