云计算百科
云计算领域专业知识百科平台

SQLMesh 告警实战:只用 YAML 配置和 SQL 模型,搭一套数据管道告警

这一篇不写一行 Python。全程只用两样东西:config.yml 配置文件,和 .sql 模型文件。

目标很实在:SQLMesh 定时任务半夜跑挂了、上游数据出现脏值了,让告警自动发到值班同学的邮箱里。

文中所有 SQLMesh 行为(事件名、YAML 字段、审计函数签名)均来自官方 Notifications guide、Auditing 文档,以及源码 notification_target.py 逐条核对。


一、先搞清楚:SQLMesh 能告什么警

告警不是「跑完发个消息」这么笼统。SQLMesh 的通知是事件驱动的——你在 notify_on 里声明订阅哪些事件,SQLMesh 就只在那些时刻发通知。

从源码的 NotificationEvent 枚举看,实际有 10 个事件(官方文档页的表格只列了 7 个,migration_* 那三个没写进去):

事件notify_on 值触发时机消息内容
Plan 开始 apply_start sqlmesh plan 开始应用 Plan \\{plan_id}` apply started for environment `{environment}`.`
Plan 结束 apply_end plan 应用成功 Plan \\{plan_id}` apply finished for environment `{environment}`.`
Plan 失败 apply_failure plan 应用抛异常 Plan \\{plan_id}` in environment `{environment}` apply failed.` + 异常栈
Run 开始 run_start sqlmesh run 开始 SQLMesh run started for environment \\{environment}`.`
Run 结束 run_end run 成功完成 SQLMesh run finished for environment \\{environment}`.`
Run 失败 run_failure run 抛异常 SQLMesh run failed. + 异常栈
审计失败 audit_failure 数据质量校验不通过 Audit failure. + 审计错误详情
Migration 开始 migration_start 内部状态表迁移开始 SQLMesh migration started.
Migration 结束 migration_end 迁移完成 SQLMesh migration finished.
Migration 失败 migration_failure 迁移抛异常 SQLMesh migration failed. + 异常栈

📌 migration_* 值得单独说一句。 它是 SQLMesh 升级版本时对内部状态表做结构迁移的事件。迁移失败意味着元数据可能处于不一致状态,是比单次 run 失败更严重的问题。官方文档页没列这三个事件,但源码里有,可以正常订阅。

实践建议:日常生产只订阅加粗的 4 个失败类事件。run_start / run_end 在多环境项目里会为每个环境各发一条,很容易把邮箱刷成噪音,然后所有人对告警免疫——那是比没有告警更糟的状态。成功与否用日报汇总就够了。 在这里插入图片描述


二、最小可用:run 失败自动发邮件

2.1 config.yml

SQLMesh 项目根目录的 config.yml 里加一段 notification_targets:

model_defaults:
dialect: duckdb
start: 2024-01-01

notification_targets:
– type: smtp
notify_on:
– run_failure
host: smtp.example.com
port: 465
user: data–alert@example.com
password: "你的邮箱授权码"
sender: data–alert@example.com
recipients:
– data–team@example.com
– oncall@example.com

就这么多。sqlmesh run 一旦失败,上面两个收件人就会收到邮件。

2.2 关键字段说明

字段说明
type 邮件用 smtp。内置还支持 slack_webhook / slack_api / console / generic
notify_on 订阅的事件列表,值见上表
host SMTP 服务器地址
port 默认 465(源码里 port: int = 465)
user / password 登录凭据。国内邮箱这里填授权码,不是登录密码
sender 发件人地址,通常需要和 user 一致,否则会被服务器拒绝
recipients 收件人列表,可多个
subject 邮件主题,默认 SQLMesh Notification(源码支持,官方文档页没提)

⚠️ 一个必须知道的限制:源码里邮件发送用的是 smtplib.SMTP_SSL,也就是只支持隐式 SSL,不支持 STARTTLS。如果你的企业邮箱(自建 Exchange、部分内网 MTA)只开了 25 或 587 端口 + STARTTLS,内置 smtp 目标用不了。这种情况的解法见本系列第二篇。

2.3 凭据不要硬编码进配置文件

config.yml 是要提交到 Git 的,密码写进去等于公开。SQLMesh 的 YAML 支持 env_var() 模板函数:

notification_targets:
– type: smtp
notify_on:
– run_failure
– apply_failure
– audit_failure
host: "{{ env_var('SMTP_HOST') }}"
port: 465
user: "{{ env_var('SMTP_USER') }}"
password: "{{ env_var('SMTP_PASSWORD') }}"
sender: data–alert@example.com
recipients:
– data–team@example.com
subject: "[数据仓库] SQLMesh 管道告警"

设置环境变量:

# Linux / macOS
export SMTP_HOST="smtp.example.com"
export SMTP_USER="data-alert@example.com"
export SMTP_PASSWORD="你的授权码"

# Windows PowerShell
$env:SMTP_HOST = "smtp.example.com"
$env:SMTP_USER = "data-alert@example.com"
$env:SMTP_PASSWORD = "你的授权码"

定时任务里记得把这些变量写进调度器(Airflow 的 Connection、crontab 的环境段、systemd 的 Environment=),否则 cron 环境读不到你 shell 里 export 的变量——这是「本地能发、上了调度就发不出去」最常见的原因。


三、失败长什么样:邮件正文预览

邮件正文由两部分拼成:msg(固定文案)+ 异常栈。以 run 失败为例,你会收到:

Subject: [数据仓库] SQLMesh 管道告警
To: data-team@example.com
From: data-alert@example.com

SQLMesh run failed.
Traceback (most recent call last):
File "…/sqlmesh/core/scheduler.py", line 214, in run
…
sqlglot.errors.ExecuteError: Binder Error: Referenced column "order_dt" not found
in FROM clause! Candidate bindings: "sushi.orders.start_ts", "sushi.orders.id"
LINE 12: WHERE order_dt >= '2024-03-01'

审计失败的正文则是 Audit failure. 加上审计错误的字符串形式,其中包含审计名、模型名、违规行数和实际执行的审计 SQL。


四、SQL 模型 + Audits:把「数据质量告警」也管起来

前面是「管道跑挂了」的告警。但更常见、也更危险的情况是:管道跑成功了,数据却是错的——上游漏推了一批、金额出现负数、主键重复了。任务全绿,报表全错,没人知道。

SQLMesh 的 audits 就是干这个的:它是一段 SQL 查询,返回任何行就代表失败。每次 run/plan 时自动执行。这一节全部是纯 SQL。

4.1 项目结构

my_project/
├── config.yml
├── models/
│ ├── orders.sql
│ ├── items.sql
│ └── customers.sql
└── audits/
└── generic.sql ← 自定义审计放这里

4.2 用内置审计(最省事,覆盖 80% 场景)

SQLMesh 内置了一大批通用审计,不需要你写任何 SQL,只要在模型的 MODEL 语句里声明即可。

models/orders.sql:

MODEL (
name sushi.orders,
kind INCREMENTAL_BY_TIME_RANGE (
time_column ds
),
owner data–eng@example.com,
audits (
not_null(columns := (id, customer_id, start_ts, end_ts)),
unique_values(columns := (id)),
number_of_rows(threshold := 100),
accepted_range(column := item_id, min_v := 1, max_v := 100000),
accepted_values(column := status, is_in := ('placed', 'paid', 'shipped', 'done', 'cancelled')),
forall(criteria := (
end_ts >= start_ts,
quantity > 0
))
)
);

SELECT
id,
customer_id,
item_id,
quantity,
status,
start_ts,
end_ts,
CAST(start_ts AS DATE) AS ds
FROM sushi.raw_orders
WHERE ds BETWEEN @start_ds AND @end_ds

这几行 audits (…) 声明的含义:

审计作用参数
not_null 指定列不允许有 NULL columns := (…)
unique_values 指定列值必须唯一(无重复) columns := (…)
number_of_rows 行数必须超过阈值——上游漏推数据的经典哨兵 threshold := N
accepted_range 数值必须落在区间内 column、min_v、max_v、可选 inclusive := false
accepted_values 枚举值白名单 column、is_in := (…)
forall 万能款:任意布尔表达式必须对所有行为 TRUE criteria := (expr1, expr2, …)

models/items.sql 再来一组,演示字符串和统计类审计:

MODEL (
name sushi.items,
kind FULL,
owner data–eng@example.com,
audits (
not_empty_string(column := name),
string_length_between(column := name, min_v := 2, max_v := 50),
accepted_range(column := price, min_v := 0.01, max_v := 1000, inclusive := false),
mean_in_range(column := price, min_v := 5, max_v := 50),
z_score(column := price, threshold := 4),
unique_combination_of_columns(columns := (id, ds))
)
);

SELECT id, name, price, CAST(ds AS DATE) AS ds
FROM sushi.seed

📌 z_score(column := price, threshold := 4) 是个很好用的异常值哨兵:任何一行的 z 分数绝对值超过 4 就告警。价格里混进一个 999999 之类的脏值,它会立刻抓住。

📌 统计类审计(mean_in_range / stddev_in_range / z_score)的阈值几乎必然需要反复试错调整,官方文档也明确这么提醒。别指望一次配对。

4.3 内置审计全清单(按用途分组)

通用断言

  • forall(criteria := (…)) — 任意布尔表达式

行数与空值

  • number_of_rows(threshold := N)
  • not_null(columns := (…))
  • at_least_one(column := c) — 该列至少有一个非 NULL 值
  • not_null_proportion(column := c, threshold := 0.8) — NULL 占比不超过阈值

具体取值

  • not_constant(column := c) — 至少有两个不同的非 NULL 值(能发现「整列变成同一个值」)
  • unique_values(columns := (…))
  • unique_combination_of_columns(columns := (…))
  • accepted_values(column := c, is_in := (…))
  • not_accepted_values(column := c, is_in := (…))

数值分布

  • sequential_values(column := c, interval := 1) — 值必须逐个递增指定步长(查断号很好用)
  • accepted_range(column := c, min_v := , max_v := , inclusive := )
  • mutually_exclusive_ranges(lower_bound_column := , upper_bound_column := ) — 区间不得重叠

字符串

  • not_empty_string(column := c)
  • string_length_equal(column := c, v := N)
  • string_length_between(column := c, min_v := , max_v := , inclusive := )
  • valid_uuid(column := c)
  • valid_email(column := c)
  • valid_url(column := c)
  • valid_http_method(column := c)
  • match_regex_pattern_list(column := c, patterns := (…))
  • not_match_regex_pattern_list(column := c, patterns := (…))
  • match_like_pattern_list(column := c, patterns := (…))
  • not_match_like_pattern_list(column := c, patterns := (…))

统计

  • mean_in_range(column := c, min_v := , max_v := )
  • stddev_in_range(column := c, min_v := , max_v := )
  • z_score(column := c, threshold := N)
  • kl_divergence(column := c, target_column := ref, threshold := 0.1) — 两列分布差异(PSI)
  • chi_square(column := c, target_column := ref, critical_value := 6.635)

⚠️ 两个官方明确提醒的 NULL 陷阱:

  • accepted_values:NULL 行在多数数据库里会通过。要挡 NULL 得另配 not_null。
  • not_accepted_values:不支持拒绝 NULL。同上。

4.4 自定义审计:写你自己的 SQL 规则

内置的不够用时,在 audits/ 目录建 .sql 文件自己写。规则很简单:查询返回行 = 审计失败。

audits/generic.sql:

— 规则 1:订单金额不能为负
AUDIT (
name assert_amount_not_negative
);
SELECT *
FROM @this_model
WHERE amount < 0;

— 规则 2:带参数的通用审计——某列不得超过阈值
AUDIT (
name does_not_exceed_threshold,
defaults (
threshold = 10
)
);
SELECT *
FROM @this_model
WHERE @column >= @threshold;

— 规则 3:增量模型专用——当天的数据必须覆盖到最新一小时
— (能发现「任务成功但只跑了一半」的情况)
AUDIT (
name assert_partition_is_fresh
);
SELECT *
FROM @this_model
WHERE ds = @end_ds
AND EXTRACT(HOUR FROM MAX(event_time)) < 20;

— 规则 4:跨模型一致性——订单里的 customer_id 必须在客户表里存在
AUDIT (
name assert_customer_exists
);
SELECT o.customer_id
FROM @this_model AS o
LEFT JOIN sushi.customers AS c
ON o.customer_id = c.id
WHERE o.ds BETWEEN @start_ds AND @end_ds
AND c.id IS NULL;

@this_model 是个特殊宏,指向正在被审计的模型;对增量模型它还保证只查相关的时间区间。@column、@threshold 是参数,在模型声明时传入。@start_ds / @end_ds 是当前处理区间。

在模型里引用:

MODEL (
name sushi.orders,
kind INCREMENTAL_BY_TIME_RANGE (time_column ds),
owner data–eng@example.com,
audits (
assert_amount_not_negative,
does_not_exceed_threshold(column := quantity, threshold := 500),
does_not_exceed_threshold(column := amount, threshold := 100000),
assert_partition_is_fresh,
assert_customer_exists
)
);

注意同一个通用审计可以用不同参数重复应用(上面 does_not_exceed_threshold 用了两次)。

💡 如果审计的参数名撞上 SQL 关键字(比如叫 values),调用时要加引号:my_audit(column := a, "values" := (1,2,3))。

4.5 内联审计:规则和模型写在一起

不想单独建 audits/ 文件的话,审计可以直接写在模型文件末尾:

MODEL (
name sushi.customers,
kind FULL,
audits (
not_null(columns := (id)),
assert_email_or_phone
)
);

SELECT id, name, email, phone FROM sushi.raw_customers;

AUDIT (
name assert_email_or_phone
);
SELECT * FROM @this_model
WHERE email IS NULL AND phone IS NULL;

4.6 全局审计:所有模型统一兜底

有些规则应该对每个模型都生效(比如「不允许空表」)。写在 config.yml 的 model_defaults 里,不用每个模型重复声明:

model_defaults:
dialect: duckdb
start: 2024-01-01
audits:
– number_of_rows(threshold := 0)
– assert_amount_not_negative

这是 YAML 配置,不是 Python——model_defaults.audits 直接支持带参数的写法。


五、blocking vs non-blocking:告警要不要拦住管道

这是审计配置里最重要的一个决定。

默认情况下,所有审计都是 blocking(阻塞)的:审计失败会中止 plan 或 run,防止脏数据往下游传播。

但 plan 和 run 的后果差别很大,必须分清:

sqlmesh plansqlmesh run
何时审计 提升到生产之前 直接对生产环境审计
审计失败后果 plan 停止,生产表完全没被碰过,脏数据留在隔离表里 run 停止,但脏数据已经在生产表里了
修复方式 改逻辑后重新 plan 必须修生产表:修上游 → 跑 restatement plan

所以:在 CI 里用 plan 阶段就拦住问题,成本远低于在 run 阶段发现。

5.1 什么时候该改成 non-blocking

有些审计是「值得关注但不该停管道」的——比如统计类阈值(mean_in_range、z_score)本身就需要长期调参,误报率高。让它阻塞管道会造成大量无谓的中断。

两种改法:

方法一:在调用处加 blocking := false

MODEL (
name sushi.items,
audits (
not_null(columns := (id)), — 阻塞:主键必须有值
accepted_range(column := price, min_v := 0.01, max_v := 1000), — 阻塞:硬规则
mean_in_range(column := price, min_v := 5, max_v := 50, blocking := false), — 只告警
z_score(column := price, threshold := 4, blocking := false) — 只告警
)
);

方法二:直接用内置的 _non_blocking 变体

每个内置审计都有一个同名 + _non_blocking 后缀的版本:

audits (
mean_in_range_non_blocking(column := price, min_v := 5, max_v := 50),
z_score_non_blocking(column := price, threshold := 4)
)

方法三:在审计定义里声明

AUDIT (
name assert_item_price_is_not_null,
blocking false
);
SELECT * FROM sushi.items
WHERE ds BETWEEN @start_ds AND @end_ds
AND price IS NULL;

⚠️ 关键:non-blocking 审计失败时仍然会触发 audit_failure 通知事件。这正是我们想要的——「不停管道,但要让人知道」。

5.2 临时跳过某个审计

调试期间想完全关掉一个审计,用 skip:

AUDIT (
name assert_item_price_is_not_null,
skip true
);
SELECT * FROM sushi.items WHERE price IS NULL;


六、审计失败告警的 5 个前提条件(最容易「配了没反应」的地方)

这是本文最实用的一节。官方文档明确写了:审计失败通知要按模型发给特定负责人,必须同时满足 5 个条件:

  • 模型的 owner 字段已填写
  • 该模型执行了一个或多个审计
  • 该 owner 配置了用户级通知目标(不是全局的)
  • 该 owner 的通知目标 notify_on 里包含审计失败事件
  • 审计在 prod 环境失败
  • 少任何一个,你就收不到通知。逐条对应到配置:

    条件 1 + 2:模型里要有 owner 和 audits

    MODEL (
    name sushi.orders,
    kind INCREMENTAL_BY_TIME_RANGE (time_column ds),
    owner zhangsan@example.com, — ← 条件 1
    audits ( — ← 条件 2
    not_null(columns := (id)),
    assert_amount_not_negative
    )
    );

    SELECT ... FROM sushi.raw_orders WHERE ds BETWEEN @start_ds AND @end_ds

    条件 3 + 4:config.yml 里要有对应用户名的用户级通知目标

    users:
    – username: zhangsan
    notification_targets:
    – type: smtp
    notify_on:
    – audit_failure # ← 条件 4
    host: "{{ env_var('SMTP_HOST') }}"
    port: 465
    user: "{{ env_var('SMTP_USER') }}"
    password: "{{ env_var('SMTP_PASSWORD') }}"
    sender: data–alert@example.com
    recipients:
    – zhangsan@example.com

    🔥 条件 3 是最常踩的坑:这里 users[].username 的值必须和模型 owner 字段的值完全一致(zhangsan@example.com),SQLMesh 就是靠这个字符串把「哪个模型出问题」映射到「该通知谁」的。写成 zhangsan 而 owner 写 zhangsan@example.com,通知就静默失效。

    条件 5:只在 prod 触发。开发环境(虚拟环境)里审计失败不会发这类定向通知。

    6.1 完整示例:按负责人分发 + 团队兜底

    model_defaults:
    dialect: duckdb
    start: 2024-01-01
    audits:
    – number_of_rows(threshold := 0)

    # ── 全局目标:管道级故障发给整个团队 ──
    notification_targets:
    – type: smtp
    notify_on:
    – run_failure
    – apply_failure
    – migration_failure
    host: "{{ env_var('SMTP_HOST') }}"
    port: 465
    user: "{{ env_var('SMTP_USER') }}"
    password: "{{ env_var('SMTP_PASSWORD') }}"
    sender: data–alert@example.com
    recipients:
    – data–team@example.com
    – oncall@example.com
    subject: "[数据仓库] SQLMesh 管道故障"

    # ── 用户级目标:数据质量告警发给对应 owner ──
    users:
    – username: zhangsan@example.com
    notification_targets:
    – type: smtp
    notify_on:
    – audit_failure
    host: "{{ env_var('SMTP_HOST') }}"
    port: 465
    user: "{{ env_var('SMTP_USER') }}"
    password: "{{ env_var('SMTP_PASSWORD') }}"
    sender: data–alert@example.com
    recipients:
    – zhangsan@example.com
    subject: "[数据质量] 你负责的模型审计失败"

    – username: lisi@example.com
    notification_targets:
    – type: smtp
    notify_on:
    – audit_failure
    host: "{{ env_var('SMTP_HOST') }}"
    port: 465
    user: "{{ env_var('SMTP_USER') }}"
    password: "{{ env_var('SMTP_PASSWORD') }}"
    sender: data–alert@example.com
    recipients:
    – lisi@example.com
    subject: "[数据质量] 你负责的模型审计失败"

    这样的分工很清晰:管道挂了找值班,数据脏了找 owner。


    七、开发期防刷屏:username 开关

    开发时会反复 sqlmesh plan,每次都往团队邮箱发通知会迅速招人烦。SQLMesh 提供了一个顶层 username 开关:

    # 只有 zhangsan 的通知目标会生效,其他全部静音
    username: zhangsan

    users:
    – username: zhangsan
    notification_targets:
    – type: smtp
    notify_on:
    – apply_start
    – apply_end
    – run_start
    – run_end
    host: "{{ env_var('SMTP_HOST') }}"
    port: 465
    user: "{{ env_var('SMTP_USER') }}"
    password: "{{ env_var('SMTP_PASSWORD') }}"
    sender: data–alert@example.com
    recipients:
    – zhangsan@example.com
    subject: "[开发环境] SQLMesh 通知(仅本人)"

    ⚠️ 上面这段是完整可复制的配置。注意 YAML 里单独一行的 … 是文档结束符,不是省略号——写在缩进结构里会直接导致解析报错(could not find expected ':')。文档里常见的 … 只是为了示意「此处省略」,实际配置必须补全字段。

    源码逻辑(NotificationTargetManager.notify)是:如果设置了 username,就只调用该用户的通知目标,全局目标全部跳过。

    更好的做法:把它放在机器级配置 ~/.sqlmesh/config.yml 里,而不是项目配置里。这样:

    • 你的开发机 → 只通知你自己
    • 生产调度机 → 没有这个配置 → 全员正常通知
    • 项目配置文件保持干净,不会有人不小心把自己的 username 提交上去导致生产告警全哑

    # ~/.sqlmesh/config.yml (只在你自己的开发机上)
    username: zhangsan

    ⚠️ 这条值得反复强调:如果误把 username: 某人 提交到了项目的 config.yml,生产环境的所有告警都只会发给那一个人,团队里其他人和值班组完全收不到——而且不会有任何报错。排查告警缺失时,这是第一个要检查的地方。


    八、验证:怎么确认告警真的能发出来

    我没有可运行的 SQLMesh 环境和真实邮箱凭据,本文配置未做端到端实测。 下面是从低成本到高成本的验证顺序:

    第 1 步:用 console 目标确认事件被触发

    源码里有个 console 类型(注释说是留给测试用的),它把通知打到控制台,不需要任何外部依赖:

    notification_targets:
    – type: console
    notify_on:
    – run_start
    – run_end
    – run_failure
    – apply_start
    – apply_end
    – apply_failure
    – audit_failure

    这一步能把「SQLMesh 没触发事件」和「邮件发送环节有问题」彻底分开。建议先做这步。

    第 2 步:单独测审计规则

    审计可以脱离 run 单独执行:

    sqlmesh -p ./my_project audit –start 2024-03-01 –end 2024-03-02

    失败时的输出长这样,信息量很足:

    Found 1 audit(s).
    assert_item_price_is_not_null FAIL.

    Finished with 1 audit error(s).

    Failure in audit assert_item_price_is_not_null for model sushi.items (audits/items.sql).
    Got 3 results, expected 0.
    SELECT * FROM sqlmesh.sushi__items__1836721418_83893210 WHERE ds BETWEEN '2024-03-01' AND '2024-03-02' AND price IS NULL

    Got 3 results, expected 0 直接告诉你有 3 行脏数据,还会打印实际执行的 SQL(含物理表名),可以复制出来单独查这 3 行。

    第 3 步:人为制造一次失败,验证邮件通路

    • 测 run_failure:临时把某个模型的 SQL 改成必然报错的(比如引用一个不存在的列),跑 sqlmesh run,看邮件到没到。测完记得改回来。
    • 测 audit_failure:临时把某个审计的阈值调到必然不通过(比如 number_of_rows(threshold := 999999999)),跑一次。

    第 4 步:在调度环境里再验一次

    第 3 步在本地通过 ≠ 生产通过。最常见的两个断点:

  • 环境变量没传进调度器——cron / Airflow worker 的环境和你的 shell 不一样。
  • 网络策略——生产机器可能不允许直连外部 SMTP 的 465 端口。

  • 九、小结

    只用 YAML + SQL,能做到的事:

    • 管道故障告警:notification_targets + type: smtp + notify_on: [run_failure, apply_failure, migration_failure]
    • 数据质量告警:模型的 audits (…) 声明 + 内置审计 + audits/*.sql 自定义规则
    • 告警分级:blocking 拦管道,blocking := false / _non_blocking 只告警
    • 责任到人:模型 owner + users[].notification_targets + notify_on: [audit_failure]
    • 全局兜底:model_defaults.audits 统一规则,顶层 notification_targets 通知团队
    • 开发期静音:~/.sqlmesh/config.yml 里的 username

    做不到的事:内置 IM 通道只有 Slack,国内团队用不上;内置 SMTP 只支持 SSL 465。这两条都需要自定义通知目标——见本系列第二篇《自定义通知渠道:把 SQLMesh 告警接到钉钉和任意 Webhook》。

    赞(0)
    未经允许不得转载:网硕互联帮助中心 » SQLMesh 告警实战:只用 YAML 配置和 SQL 模型,搭一套数据管道告警
    分享到: 更多 (0)

    评论 抢沙发

    评论前必须登录!