You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何存储dbt source freshness命令结果或仅发送异常记录告警邮件

DBT Source Freshness结果存储与告警邮件实现方案

一、将Freshness结果存储到数据库

核心思路

dbt source freshness支持输出JSON格式的结果,我们可以将JSON文件解析后写入数据库,实现结果的持久化存储。

方案1:Python脚本解析+数据库写入

  1. 生成JSON格式的Freshness结果
    dbt source freshness --output json > freshness_results.json
    
  2. 编写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
    )
    
  3. 集成到定时任务(如Airflow、Cron),定期执行上述命令和脚本。

方案2:DBT模型直接加载(适用于支持外部表的数据库)

如果你的数据库(如Snowflake、BigQuery)支持外部表,可以用DBT模型直接解析JSON文件:

  1. 在sources.yml中定义外部数据源
    sources:
      - name: dbt_metadata
        schema: public
        loader: file
        external:
          location: "/path/to/freshness_results.json"
          format: json
    
  2. 创建增量模型存储结果
    {{ 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发送

  1. 执行dbt source freshness生成JSON结果(同上文步骤)
  2. 编写告警脚本
    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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.22 18:45:32