Browse Source

refactor(monitor): cdc-monitor 作业存活改查 Flink REST,按 state 判非只判作业名

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HAeFesmMXaZSrNTdFXE8Hu
tianyu.chu 1 week ago
parent
commit
97ff4bf50d
2 changed files with 133 additions and 45 deletions
  1. 40 18
      bin/cdc-monitor.py
  2. 93 27
      tests/unit/alerter/test_cdc_monitor.py

+ 40 - 18
bin/cdc-monitor.py

@@ -1,18 +1,18 @@
 #!/usr/bin/env /usr/bin/python3
 # -*- coding:utf-8 -*-
 """
-CDC 链路监控:源库复制槽 + YARN 采集作业存活。
+CDC 链路监控:源库复制槽 + Flink 采集作业存活。
 
-DS 每 5 分钟调一次,无状态:每次独立判断。连续收到就是还没恢复,不再收到就是已恢复。
+DS 每 10 分钟调一次,无状态:每次独立判断。连续收到就是还没恢复,不再收到就是已恢复。
 
 检查两类:
   1. 复制槽——源库**必须是主库**,从库上 pg_current_wal_lsn() 报
      "recovery is in progress",且 pg_replication_slots 只能看到本节点的槽
      - cdc_ 前缀的槽 active = false → 采集作业可能已停
      - 任意槽滞后超阈值 → 消费跟不上
-  2. YARN RUNNING 列表里缺 CDC 作业 → 作业挂了
+  2. Flink session 集群上 CDC 作业缺失,或状态非 RUNNING → 作业挂了
 
-连库失败 / YARN 查询失败本身也计入告警,不静默跳过——查不到状态和状态异常同样需要人介入。
+连库失败 / Flink 查询失败本身也计入告警,不静默跳过——查不到状态和状态异常同样需要人介入。
 
 Alerter 在开跑前构造:配置缺失或 key 未替换立刻暴露,不留到真出事才发现告警链路是哑的。
 
@@ -21,18 +21,21 @@ Alerter 在开跑前构造:配置缺失或 key 未替换立刻暴露,不留
 
 CLI:
   python3 bin/cdc-monitor.py [-ds postgresql/prd-poyee-aliyun-cdc] [-channel default]
-                             [-jobs st-cdc-large,st-cdc-small] [-lag-gb 5]
+                             [-jm cdhmaster02:8081] [-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 列表里找
+  -jm       Flink JobManager REST 地址
+  -jobs     CDC 作业名,逗号分隔,逐个在 Flink 作业列表里找
   -lag-gb   槽滞后告警阈值(GB),正常应为 0
 """
 import argparse
+import json
 import os
-import subprocess
 import sys
+import urllib.request
 from datetime import datetime
 
 project_root = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
@@ -46,7 +49,7 @@ 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
+_FLINK_TIMEOUT_SEC = 20
 
 
 def check_slots(ds_ref, lag_crit_bytes):
@@ -70,16 +73,33 @@ def check_slots(ds_ref, lag_crit_bytes):
     return alerts
 
 
-def check_yarn(jobs):
-    """查 YARN RUNNING 列表,返回 [(项, 说明)]。"""
+def check_flink(jm, jobs):
+    """查 Flink session 集群的作业状态,返回 [(项, 说明)]。
+
+    必须按 state 过滤:作业取消/失败后仍留在 /jobs/overview 列表里,
+    只判名字存在会漏报——重启死循环时状态是 RESTARTING,名字一直在。
+    """
     try:
-        running = subprocess.check_output(
-            ['yarn', 'application', '-list', '-appStates', 'RUNNING'],
-            stderr=subprocess.STDOUT, timeout=_YARN_TIMEOUT_SEC).decode('utf-8')
+        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 [('YARN', '查询失败:%s' % e)]
-    return [('作业 ' + job, '不在 YARN RUNNING 列表')
-            for job in jobs if job not in running]
+        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):
@@ -95,11 +115,13 @@ def render(alerts, now=None):
 
 
 def main():
-    parser = argparse.ArgumentParser(description='CDC 链路监控(复制槽 + YARN 作业存活)')
+    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',
@@ -109,7 +131,7 @@ def main():
     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)
+    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 任务日志

+ 93 - 27
tests/unit/alerter/test_cdc_monitor.py

@@ -1,12 +1,13 @@
 # -*- coding:utf-8 -*-
 """
-bin/cdc-monitor.py 单测:复制槽判定 / YARN 存活 / 消息渲染 / main 串联。
+bin/cdc-monitor.py 单测:复制槽判定 / Flink 作业存活 / 消息渲染 / main 串联。
 
-不连真 PG(pgdb.connect + query 走替身)、不跑真 yarn(subprocess 走替身)、
+不连真 PG(pgdb.connect + query 走替身)、不打真 JobManager(urlopen 走替身)、
 不发真消息(Alerter 走替身)。
 脚本路径含连字符,用 importlib.util 动态加载为模块。
 """
 import importlib.util
+import json
 import os
 import re
 import sys
@@ -20,6 +21,7 @@ SCRIPT_PATH = os.path.join(PROJECT_ROOT, 'bin', 'cdc-monitor.py')
 
 GB = 1024 ** 3
 LAG_CRIT = 5 * GB
+JM = 'cdhmaster02:8081'
 
 
 def _load_script():
@@ -41,6 +43,22 @@ def _stub_pg(monkeypatch, rows, conn=None):
     return conn
 
 
+def _stub_flink(monkeypatch, jobs):
+    """替身 JobManager:/jobs/overview 返回给定作业列表。返回捕获的请求参数。"""
+    captured = {}
+
+    def _open(url, timeout=None):
+        captured['url'] = url
+        captured['timeout'] = timeout
+        resp = MagicMock()
+        resp.read.return_value = json.dumps({'jobs': jobs}).encode('utf-8')
+        resp.__enter__.return_value = resp
+        return resp
+
+    monkeypatch.setattr(MON.urllib.request, 'urlopen', _open)
+    return captured
+
+
 # ---------- 复制槽 ----------
 
 def test_cdc_slot_inactive_alerts(monkeypatch):
@@ -104,32 +122,62 @@ def test_db_failure_becomes_alert(monkeypatch):
     assert '数据源 ini 不存在' in alerts[0][1]
 
 
-# ---------- YARN ----------
+# ---------- Flink 作业存活 ----------
 
 def test_all_jobs_running_no_alert(monkeypatch):
-    monkeypatch.setattr(MON.subprocess, 'check_output',
-                        lambda *a, **kw: b'app-1 st-cdc-large RUNNING\napp-2 st-cdc-small RUNNING')
-    assert MON.check_yarn(['st-cdc-large', 'st-cdc-small']) == []
+    captured = _stub_flink(monkeypatch, [
+        {'name': 'st-cdc-large', 'state': 'RUNNING'},
+        {'name': 'st-cdc-small', 'state': 'RUNNING'},
+    ])
+    assert MON.check_flink(JM, ['st-cdc-large', 'st-cdc-small']) == []
+    assert captured['url'] == 'http://cdhmaster02:8081/jobs/overview'
 
 
 def test_missing_job_alerts(monkeypatch):
-    monkeypatch.setattr(MON.subprocess, 'check_output',
-                        lambda *a, **kw: b'app-1 st-cdc-large RUNNING')
-    assert MON.check_yarn(['st-cdc-large', 'st-cdc-small']) == [
-        ('作业 st-cdc-small', '不在 YARN RUNNING 列表')]
+    _stub_flink(monkeypatch, [{'name': 'st-cdc-large', 'state': 'RUNNING'}])
+    assert MON.check_flink(JM, ['st-cdc-large', 'st-cdc-small']) == [
+        ('作业 st-cdc-small', '不在 Flink 集群上')]
+
+
+def test_restarting_job_alerts(monkeypatch):
+    """2026-09-07 的故障形态:作业名一直在列表里,只判存在会漏报 7 小时。"""
+    _stub_flink(monkeypatch, [
+        {'name': 'st-cdc-large', 'state': 'RESTARTING'},
+        {'name': 'st-cdc-small', 'state': 'RESTARTING'},
+    ])
+    assert MON.check_flink(JM, ['st-cdc-large', 'st-cdc-small']) == [
+        ('作业 st-cdc-large', '状态 RESTARTING'),
+        ('作业 st-cdc-small', '状态 RESTARTING'),
+    ]
+
+
+def test_canceled_job_alerts(monkeypatch):
+    """取消的作业仍留在 /jobs/overview 里。"""
+    _stub_flink(monkeypatch, [{'name': 'st-cdc-large', 'state': 'CANCELED'}])
+    assert MON.check_flink(JM, ['st-cdc-large']) == [
+        ('作业 st-cdc-large', '状态 CANCELED')]
+
+
+def test_stale_record_does_not_mask_running(monkeypatch):
+    """同名作业的历史记录不能盖掉在跑的那条,RUNNING 优先。"""
+    _stub_flink(monkeypatch, [
+        {'name': 'st-cdc-large', 'state': 'CANCELED'},
+        {'name': 'st-cdc-large', 'state': 'RUNNING'},
+    ])
+    assert MON.check_flink(JM, ['st-cdc-large']) == []
 
 
-def test_yarn_failure_becomes_single_alert(monkeypatch):
-    def _boom(*a, **kw):
-        raise OSError('yarn: command not found')
-    monkeypatch.setattr(MON.subprocess, 'check_output', _boom)
-    alerts = MON.check_yarn(['st-cdc-large', 'st-cdc-small'])
+def test_flink_failure_becomes_single_alert(monkeypatch):
+    def _boom(url, timeout=None):
+        raise OSError('Connection refused')
+    monkeypatch.setattr(MON.urllib.request, 'urlopen', _boom)
+    alerts = MON.check_flink(JM, ['st-cdc-large', 'st-cdc-small'])
     assert len(alerts) == 1          # 查不到就是一条,不按作业数量刷屏
-    assert alerts[0][0] == 'YARN'
-    assert 'command not found' in alerts[0][1]
+    assert alerts[0][0] == 'Flink 集群'
+    assert 'Connection refused' in alerts[0][1]
 
 
-# ---------- 渲染:走一圈全部五类告警 ----------
+# ---------- 渲染:走一圈全部类告警 ----------
 
 FIXED_NOW = datetime(2026, 8, 26, 15, 21, 27)
 
@@ -151,26 +199,27 @@ def test_render_all_alert_kinds():
     alerts = [
         ('槽 cdc_large', '空闲,采集作业可能已停'),
         ('槽 cdc_small', '滞后 6.0 GB'),
-        ('作业 st-cdc-large', '不在 YARN RUNNING 列表'),
+        ('作业 st-cdc-large', '不在 Flink 集群上'),
+        ('作业 st-cdc-small', '状态 RESTARTING'),
         ('主库', '复制槽查询失败:connection refused'),
-        ('YARN', '查询失败:timeout'),
+        ('Flink 集群', '查询失败:timeout'),
     ]
     out = MON.render(alerts, now=FIXED_NOW)
     lines = out.split('\n')
     assert lines[0] == '### <font color="warning">CDC 监控告警</font>'
     assert lines[1] == '> 检测时间:2026-08-26 15:21:27'
-    assert len(lines) == 7
+    assert len(lines) == 8
     for (key, desc), line in zip(alerts, lines[2:]):
         assert line == '> {0}:<font color="warning">{1}</font>'.format(key, desc)
 
 
 # ---------- main 串联 ----------
 
-def _run_main(monkeypatch, slot_alerts, yarn_alerts):
+def _run_main(monkeypatch, slot_alerts, flink_alerts):
     """跑 main,返回替身 alerter 实例供断言。"""
     monkeypatch.setattr(sys, 'argv', ['cdc-monitor.py'])
     monkeypatch.setattr(MON, 'check_slots', lambda ds, lag: list(slot_alerts))
-    monkeypatch.setattr(MON, 'check_yarn', lambda jobs: list(yarn_alerts))
+    monkeypatch.setattr(MON, 'check_flink', lambda jm, jobs: list(flink_alerts))
     instance = MagicMock()
     monkeypatch.setattr(MON, 'Alerter', MagicMock(return_value=instance))
     MON.main()
@@ -186,7 +235,7 @@ def test_main_sends_once_with_all_alerts(monkeypatch):
     instance = _run_main(
         monkeypatch,
         [('槽 cdc_large', '空闲,采集作业可能已停')],
-        [('作业 st-cdc-small', '不在 YARN RUNNING 列表')])
+        [('作业 st-cdc-small', '状态 RESTARTING')])
     instance.send_markdown.assert_called_once()
     content = instance.send_markdown.call_args[0][0]
     assert '槽 cdc_large' in content
@@ -197,21 +246,37 @@ def test_main_builds_alerter_before_checks(monkeypatch):
     """健康时也要构造 Alerter:配置坏了要立刻暴露,不能等真出事才发现发不出去。"""
     monkeypatch.setattr(sys, 'argv', ['cdc-monitor.py'])
     monkeypatch.setattr(MON, 'check_slots', lambda ds, lag: [])
-    monkeypatch.setattr(MON, 'check_yarn', lambda jobs: [])
+    monkeypatch.setattr(MON, 'check_flink', lambda jm, jobs: [])
     fake_cls = MagicMock()
     monkeypatch.setattr(MON, 'Alerter', fake_cls)
     MON.main()
     fake_cls.assert_called_once_with(channel='default')
 
 
+def test_main_defaults_cover_both_jobs(monkeypatch):
+    """默认值必须带上小表组,否则它挂了不报警。"""
+    seen = {}
+    monkeypatch.setattr(sys, 'argv', ['cdc-monitor.py'])
+    monkeypatch.setattr(MON, 'check_slots', lambda ds, lag: [])
+    monkeypatch.setattr(MON, 'check_flink',
+                        lambda jm, jobs: seen.update(jm=jm, jobs=jobs) or [])
+    monkeypatch.setattr(MON, 'Alerter', MagicMock())
+
+    MON.main()
+
+    assert seen['jm'] == 'cdhmaster02:8081'
+    assert seen['jobs'] == ['st-cdc-large', 'st-cdc-small']
+
+
 def test_main_passes_cli_overrides(monkeypatch):
     seen = {}
     monkeypatch.setattr(sys, 'argv', [
         'cdc-monitor.py', '-ds', 'postgresql/other', '-channel', 'realtime',
-        '-jobs', 'a, b ,', '-lag-gb', '2'])
+        '-jm', 'other-host:9081', '-jobs', 'a, b ,', '-lag-gb', '2'])
     monkeypatch.setattr(MON, 'check_slots',
                         lambda ds, lag: seen.update(ds=ds, lag=lag) or [])
-    monkeypatch.setattr(MON, 'check_yarn', lambda jobs: seen.update(jobs=jobs) or [])
+    monkeypatch.setattr(MON, 'check_flink',
+                        lambda jm, jobs: seen.update(jm=jm, jobs=jobs) or [])
     fake_cls = MagicMock()
     monkeypatch.setattr(MON, 'Alerter', fake_cls)
 
@@ -219,5 +284,6 @@ def test_main_passes_cli_overrides(monkeypatch):
 
     assert seen['ds'] == 'postgresql/other'
     assert seen['lag'] == 2 * GB
+    assert seen['jm'] == 'other-host:9081'
     assert seen['jobs'] == ['a', 'b']      # 逗号分隔去空白、丢空项
     fake_cls.assert_called_once_with(channel='realtime')