---
title: "DQC 结果的实时告警投递：用 binlog 把配置热更新做进调度器"
description: "数据质量告警投递层的工程复盘。内容包括配置热更新的三次迭代、MySQL binlog 监听、Quartz 与 APScheduler 的 cron 差异、一条多层 CTE 的批量查询、告警的沉默歧义，以及共享数据库集成的代价。"
pubDate: 2025-11-02
tags: ["工程复盘", "MySQL", "Binlog", "CDC", "Python", "数据质量"]
---

把一条监控规则的报警时间从 14:59 改成 15:30，回车。不到一秒，调度器重建完这个任务。没有发布，没有重启，没有接口调用。

实现方式不是消息队列，也不是配置中心：调度进程伪装成一台 MySQL 从库，从 binlog 里读配置变更。

下面是这套机制的设计过程，以及它带来的新问题。

## 边界：只做最后一公里

数据仓库里有一批玩家标签字段，供下游 BI 报表和运营看板使用。一个字段算错，下游一整串跟着错，而且往往几天后才被发现。典型事故有这么几类：

- 上游 ETL 没跑完或跑挂，字段整片是 NULL
- 计算逻辑出错，累充次数全变成 0
- 上游加了新枚举值，下游没兼容
- 字段类型被改，下游解析失败
- 同一个标签在两张表里对不上

规则字典表只有五行，对应的正是这五类：空值校验、0 值校验、枚举值校验、类型校验，以及跨实体或同实体标签校验。最后一条是唯一需要跨表比对的，后面那条大查询有一半的复杂度是为它准备的。

**校验引擎不在本文范围内。** 它每天跑完把结论写进 MySQL，投递层从那里往后接。现有证据表明它出自 Java 生态——表里的 JSON 字段全是 camelCase，配置表的报警时间字段存的是 Quartz 六段 cron。这个格式后面会成为一个坑。

投递层对四张表全程只读，没有任何 `INSERT` 或 `UPDATE`。

四张表构成一条由粗到细的漏斗：

| 表 | 规模 | 回答什么 |
|---|---|---|
| `规则字典表` | 5 行 | 有哪几种校验 |
| `标签监控规则配置表` | 几十条 | 盯什么字段、几点报、报给哪个群 |
| `每日执行结果汇总表` | 万级 | 今天这个字段正不正常 |
| `异常明细表` | 百万级 | 问题落在哪几行 |

前三张回答「有没有事」，最后一张回答「事在哪」。

配置表的字段按职责分四组，这个划分决定了整个系统的形状：

| 组 | 字段 | 作用 |
|---|---|---|
| 盯什么 | 游戏 ID、发行主体、表名、字段名、目标表名 | 定位监控对象 |
| 怎么判 | 规则 ID、用户阈值、默认阈值 | 用哪条规则、阈值多少 |
| 怎么报 | 是否飞书报警、报警频率、报警时间、报警 webhook | 报不报、几点报、报给哪个群 |
| 状态审计 | 是否启用、是否删除、创建更新时间与操作人 | 软删除与操作追溯 |

一条配置的实际形态：某个游戏、某个发行主体，账号属性宽表里的「截至昨日累计充值次数」字段，0 值校验，每天 14:59 报到指定飞书群。表名那一列存的不是固定库名，而是带 `{{game_id}}` 占位符的模板，运行时按游戏展开成各自的库——配置表本身是个小 DSL。

投递层只做三件事：读结论、组卡片、发对群。难点不在这三件事本身，在于配置是活的——这也是下面全部设计的起点。

## 配置搬进数据库之后

监控规则最初写在代码里。加一条规则要改代码、重启服务；QA 或运营发现某个字段需要盯，自己改不了，得排期找开发。响应周期以小时计。

改造方向是把规则搬进数据库配置表，让他们自己增删改。这解决了一半问题。

另一半随之出现：**配置改完了，正在跑的调度器怎么知道？**

需求的形状是：报警时间按数据库里的时间点去发飞书消息，而这个时间不固定，可能中途修改，也可能新增或减少一条记录。

「修改、新增、减少」对应的正是 `UPDATE`、`INSERT`、`DELETE`。需求从一开始就是变更数据捕获（CDC）的形状。

## 三次迭代：轮询、时间戳、binlog

第一版用定时轮询：

```python
# 伪代码
while True:
    time.sleep(60)
    sync_jobs()      # 全表 SELECT，重建所有任务
```

每分钟扫一次全表，不管有没有变化都要重建 CronTrigger、重新比对增删改。配置表一天可能一次都不改，这些查询全是浪费。

第二版加短路：每轮先 `SELECT MAX(update_time)` 取任务表的最大更新时间，与上一轮缓存的值比对，两者相同就直接返回，连任务列表都不查。

轮询开销降到可以忽略，数据库压力接近于零。但**延迟没变**：最坏情况下，配置改完要等满一个轮询周期才生效。本质仍是轮询，只是轮得便宜了。

第三版换成 binlog 监听：

```python
# 伪代码，表名已脱敏
BinLogStreamReader(
    connection_settings=db_settings,
    server_id=<一个固定常量>,                     # 必须全局唯一
    only_events=[WriteRowsEvent, UpdateRowsEvent, DeleteRowsEvent],
    only_tables=["标签监控规则配置表"],           # 只订阅这一张表
    blocking=True,
    resume_stream=True,
    freeze_schema=True,
)
```

延迟从一个轮询周期降到亚秒级，数据库压力仍接近于零——不再有任何主动查询，事件由 MySQL 推送。

**为什么不用 Canal、Debezium 或 Maxwell。** 这三者是中台级方案，适合多库多表的变更同步。现有业务只需要监听一张配置表。引 Canal 要搭 JVM，引 Debezium 要搭 Kafka Connect，为一张配置表付这个基础设施成本不成比例。`python-mysql-replication` 是纯 Python 库，装上即可用。

这个选择反过来约束了数据库驱动：放弃性能更好的 C 扩展驱动，改用纯 Python 的 pymysql，理由之一是它有配套的 replication 实现。

**binlog 监听的本质是一条伪装成从库的长连接。** 它走 MySQL 的 replication 协议，服务端会建一个 dump 线程并按 `server_id` 登记，`SHOW PROCESSLIST` 里能直接看到。

`server_id` 因此必须全局唯一：**新连接带着已存在的 ID 注册时，MySQL 会主动杀掉持有该 ID 的旧连接。** 这个行为本是为了清理从库重连留下的僵尸连接，副作用是两个进程用同一个 ID 会互相踢下线。

这条约束直接影响了部署方式。投递层曾考虑跟 Web 服务合并进同一个 gunicorn，当时的配置是两个 worker，那意味着每个进程都会带着同一个 `server_id` 去注册，后连上的踢掉先连上的，被踢的重连再踢回去，监听会时断时续。

而且问题不止 binlog：每个 worker 还会各注册一份定时任务，同一条告警被发送多次。这种重复不报错、不留异常日志，只是群里多出一条一模一样的消息，**比监听断连更难被发现**。

所以投递层最终是独立进程，与 Web 服务分开部署。单实例之下 `server_id` 用固定常量即可——**这个问题是靠部署形态规避的，不是靠代码兜住的**。

选值本身还有个没解决的隐患。`server_id` 若与某台生产从库相同，被踢下线的就是那台从库，复制直接断掉。实际用的是一个三位数的低位常量，和常规从库的编号挤在同一区间；事后沉淀下来的结论是取 100000 以上，但代码里没有改。

资源消耗方面，binlog 写入是异步 IO，不阻塞事务提交，对主库影响很小。风险在另一头：ROW 格式下大表高频变更会让 binlog 文件迅速膨胀，消费端一旦阻塞，日志堆积可能占满磁盘。监听一张低频变更的配置表不触及这些风险，但上生产前得先算这笔账。

**第二版的时间戳短路保留了下来。** binlog 已经能精确告知变更内容，`MAX(update_time)` 看似冗余。留着它有两个用途：启动时的全量同步复用这段逻辑；监听器一旦挂掉，轮询是唯一的降级路径。

另有一处差异值得注意：`MAX(update_time)` 检测不到物理 `DELETE`，删掉一行后最大更新时间可能反而变小。配置表走软删除，删除本质是 `UPDATE`，所以没踩到；如果配置表用物理删除，这个短路会漏掉整行的消失。binlog 路径没有这个问题，`DeleteRowsEvent` 是独立的事件类型。

## 第一次启动，四条配置全挂

监听器跑起来了，任务一个都没注册上。四条配置，四条全部解析失败，日志里刷的是 APScheduler 自己抛的那句 `Unrecognized expression "?" for field "day_of_week"`。换一条配置，报错里的字段名跟着换成 `day`——`?` 写在哪一段，它就在哪一段卡住。

根因是两套 cron 方言的差异。配置表里的表达式是六段的 Quartz 格式（秒 分 时 日 月 周），用 `?` 表示「此位不限定」。APScheduler 的 `CronTrigger.from_crontab()` 只认标准 Unix crontab 的五段，也不识别 `?`。

修复需要两步，第一次只做了一半。

第一步绕开 `from_crontab()`，改用 `CronTrigger` 构造函数——它原生支持秒字段，六段表达式可以按 `second` / `minute` / `hour` / `day` / `month` / `day_of_week` 逐个传进去，顺带把时区钉死在 `Asia/Shanghai`。秒的问题解决了，`?` 还在，于是又挂了一轮。

第二步才是完整修复：解析前先把整串里的 `?` 全部替换成 `*`，再按空格切分——六段走构造函数，五段回落到 `from_crontab()`，其余长度直接抛错。

十九分钟后重启，四个任务全部注册成功，监听器正常接管。

`from_crontab` 是类 Unix 风格，`CronTrigger` 构造函数是类 Quartz 风格。跨生态搬配置时，**格式不兼容比想象中更容易被漏掉——它不在代码里，在数据里**。

## 合并调度，拆开投递

<picture>
  <source media="(max-width: 640px)" srcset="./dqc-assets/two-lanes-mobile.svg">
  <img src="./dqc-assets/two-lanes.svg" alt="配置热更新与到点执行两条链路：链路 A 由改配置经 MySQL binlog、监听器到重建调度表，链路 B 由到点触发经批量查询、按 webhook 分组到发飞书卡片，两条链路的交点是调度表" loading="lazy" decoding="async">
</picture>

*图1：两条链路各自独立，唯一的交点是调度表。链路 A 负责维护它，链路 B 负责消费它。*

系统由两条独立链路组成。链路 A 事件驱动：配置变更经 binlog 进来，重建调度表。链路 B 时间驱动：到点触发，查数据，发卡片。两者唯一的交点是 APScheduler 的任务表——A 维护它，B 消费它，**改配置和跑任务互不阻塞**。

链路 B 里有两次方向相反的分组。

**第一次是合并。** 最初一条规则对应一个定时任务，N 条规则就是 N 个 CronTrigger。实际情况是大量规则的报警时间相同——同一批监控往往被设在同一个时间点。改成按报警时间分组后，相同时间点的规则合并为一个任务：

```
4 条规则  →  2 个 CronTrigger
```

触发时整组一起处理，数据库查询从 N 次变成一次。这个改动在发送函数的签名上留下了痕迹：参数从单个任务变成了任务列表，函数开头还留着一道类型检查——收到非列表就记一条 warning、直接返回失败，提示调用方更新用法。那句 warning 是写给还没迁移的旧调用点的，也是合并确实上过线的证据。

**第二次是拆开。** 每条配置可以指定自己的报警 webhook，同一时间点触发的规则可能要发往不同的群。所以查完数据后，再按 webhook 重新分组，分别投递。

合并是为了少建任务，拆开是为了发对人。两个维度正交。

## 让百万行的明细表算得动

飞书卡片里有一列是异常值的去重计数，只对跨实体校验这一类规则输出。同一个异常值可能对应很多行，而要看的是有多少种不同的异常值，不是总行数。

最初用相关子查询实现：主查询每出一行，就去异常明细表做一次 `COUNT(DISTINCT)`。明细表有几百万行，这个写法等于让它被反复扫描，查询时间不可控。

重构后是一条 CTE，一次查完整批规则。`tasks` 那一层只是把本批次的规则参数化拼成一张临时表，真正做事的是后面四层：

```sql
-- 表名与业务字段名已脱敏
WITH
tasks     AS ( /* 本批次要查的规则，UNION ALL 拼成一张临时表 */ ),
latest_dt AS ( SELECT s.config_id, MAX(s.dt) AS dt FROM 每日执行结果汇总表 s
               JOIN tasks t ON s.config_id = t.config_id
               GROUP BY s.config_id ),
ranked    AS ( SELECT s.*, ROW_NUMBER() OVER (
                   PARTITION BY s.config_id, s.game_id, ...
                   ORDER BY s.update_time DESC, s.id DESC) AS rn
               FROM 每日执行结果汇总表 s JOIN tasks t ON ...
               JOIN latest_dt d ON s.config_id = d.config_id AND s.dt = d.dt
               WHERE s.status = 'ERROR' ),
picked    AS ( SELECT * FROM ranked WHERE rn = 1 ),
dd        AS ( SELECT ..., COUNT(DISTINCT d.异常值) AS 去重个数
               FROM 异常明细表 d JOIN picked p ON ...
               GROUP BY ... )
SELECT ... FROM picked p LEFT JOIN dd ON ...
```

关键在 `dd` 这一层：原本散在每一行的去重计数，改成一次性按维度分组聚合，明细表只扫一遍。后来又为这个查询建了一条联合索引，把选择性最强的规则 ID 放在最左列。

这条 SQL 里更值得注意的不是性能，是 `latest_dt`。

它取的是每个配置**最近一次有数据**的日期，而不是**最近一次 ERROR** 的日期。两者只差一个筛选条件，行为完全不同。

取最近一次 ERROR 日期，会导致这样的结果：某个配置的问题已经修好，但只要历史上出现过 ERROR，系统就会一直把那条旧记录捞出来报警。**告警持续响，问题早已不存在。**

改成「取最近一次有数据的日期，再从中筛出 ERROR」后，逻辑自洽：配置恢复正常的第二天，最新日期那批数据里没有 ERROR，自然不再报。

这个改动配了一组测试，覆盖六个场景：

| 场景 | 该不该报 |
|---|---|
| 最近日期有 ERROR | 报 |
| **最近日期正常，历史有 ERROR** | **不报** |
| 连续多天 ERROR | 继续报 |
| 今天刚出现 ERROR | 报 |
| 同一天 ERROR 与正常混合 | 只报 ERROR 那条 |
| 真实数据回放 | 发到测试群验证 |

第二个场景就是这次重构的起因。告警系统一旦开始报「已经不存在的问题」，人就会开始无视所有告警。这比漏报更危险，因为它是慢性的。

## 群里没消息，是好事还是坏事

后期加了一个功能：**查询结果没有异常时，也发一张卡片。**

群里一整天没有消息，无法判断是数据没问题，还是服务挂了。这两种情况表现一致，含义完全相反。加上正常状态通知之后，沉默只剩一种解释：进程死了。这条通知的作用是**把系统存活状态变成可观测的**。

实现分两种情况：整批任务都正常时发一张汇总卡；部分正常部分异常时，异常告警照常发到各自的群，再补一张正常任务的汇总。

发送环节有两个细节。

**HTTP 200 不代表发送成功。** 飞书 webhook 接口即使卡片结构有问题也返回 200，真正的状态在响应体里。判断条件因此落在业务状态码上，而不是 HTTP 状态码：POST 回来先尝试解析 JSON，解析不出来直接算失败、把响应文本截断记进日志；解析得出来，再看 `code` 字段是不是 0。两道都过了才算发送成功。

**一次触发可能发出多条消息。** 一个时间点的任务组会按 webhook 拆成若干条异常告警，可能再加一条全局抄送，最后加一条正常汇总。一次普通的定时触发发出四条消息是常态。

## 没做到的部分

**发送失败没有重试。** 发送函数里失败只记一条错误日志和推送失败的卡片消息到运维群，卡片就此丢失，没有重投也没有死信队列。对告警系统而言这是个明显缺口——最需要送达的时刻，往往正是系统不稳定的时刻。

**binlog 位点没有持久化。** 监听器用了 `resume_stream=True`，但没有保存 binlog 文件名和位点。进程重启后从当前最新位置开始读，宕机期间的配置变更会丢。

这一条不算疏漏：启动顺序兜住了它——程序先做一次全量同步，再启动监听。全量兜底加增量提速是 CDC 的常见模式。配置表数据量极小，全量同步成本接近于零，这个取舍成立。换成业务数据同步，位点持久化就是必需的。

## 共享数据库集成的代价

上游的校验引擎与 Python 投递层之间**只通过 MySQL 表交接，没有接口调用**。

好处很实在：上线快，不需要对齐发版节奏，不需要维护接口契约，也不需要处理服务间调用的超时与重试。数据团队里这种形态很常见，因为数据本来就在库里。

代价是耦合从接口转移到了表结构。这不是假想的风险，库里就有现成的例子：「发行主体」这个概念，配置表用的是完整拼写，汇总表和明细表用的是它的缩短形式，两个列名只差中间几个字母。三张表要 JOIN 在一起，投递层就得在规范化任务的那一步专门写一段兼容代码，两个名字都试一遍——这种东西没有任何契约能提前告诉你，只能撞上之后补。

上游加一列、改一个 JSON 字段名、调整一次状态枚举，投递层都得跟着改，而且没有任何契约能提前发现，通常是飞书卡片里突然出现一片空值才知道。卡片的表格模板也是这种耦合下的自我防御：上游写进来的阈值 JSON 有两种结构，空值和 0 值校验一种，跨实体校验另一种，代码里没法拿到契约，只能备两套模板，靠字段是否存在来判断该用哪一套。

重做一次的话，可以在两个系统之间加一层视图或物化表，把上游的物理表结构和下游查询隔开。视图的列名与类型由下游定义，上游改动时先在视图层适配，投递层代码不动。这只是把隐式契约变成显式契约。

开头那个「改完就生效」，我认为价值不在延迟数字本身。生效延迟从一个轮询周期降到亚秒级之后，配置的使用方式会跟着变：从「攒一批一起改」变成「想到就改」。
