如何存储dbt source freshness命令结果或仅发送异常记录告警邮件
DBT Source Freshness结果存储与告警邮件实现方案
一、将Freshness结果存储到数据库
核心思路
dbt source freshness支持输出JSON格式的结果,我们可以将JSON文件解析后写入数据库,实现结果的持久化存储。
方案1:Python脚本解析+数据库写入
- 生成JSON格式的Freshness结果
dbt source freshness --output json > freshness_results.json - 编写Python脚本解析并写入数据库(以PostgreSQL为例,适配其他数据库只需修改连接字符串)
import pandas as pd from sqlalchemy import create_engine import os # 读取JSON结果 freshness_data = pd.read_json("freshness_results.json") # 展平嵌套的source和tables字段 flattened_data = freshness_data.explode("sources").explode("sources.tables").reset_index(drop=True) tables_df = pd.json_normalize(flattened_data["sources.tables"]) # 从环境变量读取数据库凭证(避免硬编码) db_user = os.getenv("DB_USER") db_pass = os.getenv("DB_PASS") db_host = os.getenv("DB_HOST") db_port = os.getenv("DB_PORT") db_name = os.getenv("DB_NAME") # 连接数据库并写入数据 engine = create_engine(f"postgresql://{db_user}:{db_pass}@{db_host}:{db_port}/{db_name}") tables_df.to_sql( name="dbt_source_freshness_results", con=engine, if_exists="append", index=False ) - 集成到定时任务(如Airflow、Cron),定期执行上述命令和脚本。
方案2:DBT模型直接加载(适用于支持外部表的数据库)
如果你的数据库(如Snowflake、BigQuery)支持外部表,可以用DBT模型直接解析JSON文件:
- 在
sources.yml中定义外部数据源sources: - name: dbt_metadata schema: public loader: file external: location: "/path/to/freshness_results.json" format: json - 创建增量模型存储结果
{{ config(materialized='incremental', unique_key='table_name') }} with raw_freshness as ( select * from {{ source('dbt_metadata', 'freshness_results') }} ), flattened_results as ( select s.source_name, t.table_name, t.max_loaded_at::timestamp, t.snapshotted_at::timestamp, t.status, t.error, t.freshness.warning_after.count as warning_threshold_count, t.freshness.warning_after.unit as warning_threshold_unit, t.freshness.critical_after.count as critical_threshold_count, t.freshness.critical_after.unit as critical_threshold_unit from raw_freshness, lateral flatten(input => sources) s, lateral flatten(input => s.tables) t ) select * from flattened_results {% if is_incremental() %} where snapshotted_at > (select max(snapshotted_at) from {{ this }}) {% endif %}
二、针对Freshness警告/错误发送告警邮件
方案1:Python脚本解析+SMTP发送
- 执行
dbt source freshness生成JSON结果(同上文步骤) - 编写告警脚本
import json import smtplib from email.mime.text import MIMEText import os # 读取Freshness结果 with open("freshness_results.json", "r") as f: freshness_data = json.load(f) # 筛选警告/错误记录 alert_list = [] for source in freshness_data["sources"]: source_name = source["source_name"] for table in source["tables"]: if table["status"] in ["warning", "error"]: alert_detail = ( f"数据源: {source_name}\n" f"表名: {table['table_name']}\n" f"状态: {table['status'].upper()}\n" f"最后加载时间: {table['max_loaded_at']}\n" f"快照时间: {table['snapshotted_at']}\n" f"错误信息: {table.get('error', '无')}\n" "-------------------------\n" ) alert_list.append(alert_detail) # 发送告警邮件 if alert_list: email_content = "以下数据源新鲜度不符合要求:\n\n" + "".join(alert_list) msg = MIMEText(email_content, "plain", "utf-8") msg["Subject"] = "DBT数据源新鲜度告警" msg["From"] = os.getenv("ALERT_FROM_EMAIL") msg["To"] = os.getenv("ALERT_TO_EMAIL") # 连接SMTP服务器发送邮件 with smtplib.SMTP_SSL(os.getenv("SMTP_HOST"), os.getenv("SMTP_PORT")) as server: server.login(os.getenv("SMTP_USER"), os.getenv("SMTP_PASS")) server.send_message(msg)
方案2:Shell脚本+JQ解析+Mail命令
适合无Python环境的场景:
#!/bin/bash # 生成Freshness JSON结果 dbt source freshness --output json > /tmp/freshness_results.json # 用JQ筛选警告/错误记录 ALERT_CONTENT=$(jq -r ' .sources[] | .tables[] | select(.status == "warning" or .status == "error") | "数据源: \(.source_name)\n表名: \(.table_name)\n状态: \(.status)\n最后加载时间: \(.max_loaded_at)\n错误: \(.error // "无")\n---\n" ' /tmp/freshness_results.json) # 发送邮件(需系统已配置mail命令) if [ -n "$ALERT_CONTENT" ]; then echo -e "DBT数据源新鲜度告警:\n\n$ALERT_CONTENT" | mail -s "DBT Freshness Alert" alert-recipient@yourdomain.com fi
集成建议
将上述脚本与CI/CD工具(如GitHub Actions、GitLab CI)或调度工具(如Airflow、Cron)结合,实现定期执行和自动告警。
内容的提问来源于stack exchange,提问作者amit mishra
相关产品推荐
相关产品推荐

