Featured image of post 用 DuckDB + Cron 搭建全自动数据监控报警系统,睡后收入的新思路

用 DuckDB + Cron 搭建全自动数据监控报警系统,睡后收入的新思路

手把手教你用 DuckDB + Python + Telegram Bot 搭建全自动数据健康监测系统,自动发现异常、实时推送告警,打造可复售的监控 SaaS 产品。

用 DuckDB + Cron 搭建全自动数据监控报警系统,睡后收入的新思路

你有没有这样的痛点——业务数据每天都在涨,但出了问题却不知道。等客户来找你反馈,损失已经造成了。

今天教你用 DuckDB 搭建一套全自动的「数据健康监测系统」。每天定时跑查询,发现异常自动发 Telegram 通知。不用写复杂的 ETL,不用买昂贵的监控工具,一个 Python 脚本 + DuckDB 就能搞定。

DuckDB 数据监控报警系统架构图

一、为什么需要数据监控?

先说清楚:为什么要做数据监控?

  • 电商场景:GMV 突然下跌 30%,你要在客户下单前发现,而不是第二天看报表
  • 金融场景:某只股票放量突破关键位,自动推送给订阅用户
  • SaaS 场景:用户注册量连续 3 天下降,说明产品出了问题
  • 内容场景:某个视频播放量突然暴涨,及时跟进推广

这些场景的共同点是:发现要快,通知要准,系统要省

DuckDB 在这套系统里的角色:它不是数据库,而是「查询引擎」——直接从 CSV/Parquet/API 读数据,一行 SQL 查出异常,结果交给 Telegram 推送。

二、系统架构

CSV/Parquet 数据源
        ↓
   DuckDB(查询引擎)
        ↓
   Python 异常检测逻辑
        ↓
   Telegram Bot API
        ↓
   手机实时通知

这一层全部用 Python 完成,零外部依赖服务。数据存储用 .duckdb 单文件,增量追加即可。

三、核心监控逻辑

假设你是一个跨境电商卖家,关心三个核心指标:

  • 日订单量是否异常波动
  • 客单价是否持续下降
  • 退款率是否突然升高

数据源是每天的订单 CSV 导出。用 DuckDB 写监控查询:

import duckdb
import json
from datetime import datetime, timedelta

con = duckdb.connect("monitor.duckdb")

# 把历史订单合并成一个持久表(首次导入)
con.execute("""
    CREATE TABLE IF NOT EXISTS orders AS
    SELECT * FROM read_csv_auto('orders_2024*.csv')
""")

# 监控查询:计算过去 14 天的核心指标
def check_health():
    today = datetime.now().date()
    
    result = con.execute("""
        WITH daily_stats AS (
            SELECT
                DATE(order_date) AS stat_date,
                COUNT(*) AS order_count,
                ROUND(AVG(amount), 2) AS avg_order_value,
                ROUND(SUM(CASE WHEN status = 'refunded' THEN 1 ELSE 0 END) * 100.0 / COUNT(*), 2) AS refund_rate
            FROM orders
            WHERE DATE(order_date) >= CURRENT_DATE - INTERVAL '14' DAY
            GROUP BY DATE(order_date)
        )
        SELECT
            MAX(CASE WHEN stat_date = CURRENT_DATE - INTERVAL '1' DAY THEN order_count END) AS yesterday_orders,
            AVG(order_count) OVER (ORDER BY stat_date ROWS BETWEEN 6 PRECEDING AND 1 PRECEDING) AS ma7_orders,
            MAX(CASE WHEN stat_date = CURRENT_DATE - INTERVAL '1' DAY THEN refund_rate END) AS yesterday_refund_rate,
            AVG(refund_rate) OVER (ORDER BY stat_date ROWS BETWEEN 6 PRECEDING AND 1 PRECEDING) AS ma7_refund_rate
        FROM daily_stats
        ORDER BY stat_date DESC
        LIMIT 1
    """).fetchone()
    
    return {
        'yesterday_orders': result[0],
        'ma7_orders': result[1],
        'yesterday_refund_rate': result[2],
        'ma7_refund_rate': result[3]
    }

关键点:

  1. read_csv_auto('orders_2024*.csv') — DuckDB 支持通配符一次读多个文件,不用手动拼接
  2. 窗口函数 AVG(...) OVER (...) — 计算 7 日移动平均,不用预聚合
  3. 结果存在 .duckdb 文件里 — 增量追加而不是全量重建

四、异常检测算法

有了基础指标,接下来判断「是否异常」。

这里用简单的统计学方法:如果当前值偏离 7 日均值超过一定阈值,就认为是异常。

def detect_anomalies(stats):
    """返回异常列表,每项包含指标名、当前值、预期范围、异常程度"""
    anomalies = []
    
    # 订单量异常
    if stats['ma7_orders'] > 0:
        order_drop = (stats['yesterday_orders'] - stats['ma7_orders']) / stats['ma7_orders']
        if order_drop < -0.20:  # 下跌超过 20%
            anomalies.append({
                'metric': '订单量',
                'value': stats['yesterday_orders'],
                'expected': f"约 {stats['ma7_orders']:.0f}",
                'severity': '🔴 严重',
                'message': f"昨日订单量 {stats['yesterday_orders']},较7日均值下跌 {abs(order_drop)*100:.1f}%"
            })
        elif order_drop > 0.50:  # 暴涨超过 50%
            anomalies.append({
                'metric': '订单量',
                'value': stats['yesterday_orders'],
                'expected': f"约 {stats['ma7_orders']:.0f}",
                'severity': '🟡 关注',
                'message': f"昨日订单量 {stats['yesterday_orders']},较7日均值暴涨 {order_drop*100:.1f}%,请确认是否为正常活动"
            })
    
    # 退款率异常
    if stats['ma7_refund_rate'] > 0:
        refund_spike = (stats['yesterday_refund_rate'] - stats['ma7_refund_rate']) / stats['ma7_refund_rate']
        if refund_spike > 0.50:  # 退款率暴涨 50%
            anomalies.append({
                'metric': '退款率',
                'value': f"{stats['yesterday_refund_rate']:.2f}%",
                'expected': f"约 {stats['ma7_refund_rate']:.2f}%",
                'severity': '🔴 严重',
                'message': f"昨日退款率 {stats['yesterday_refund_rate']:.2f}%,较7日均值上涨 {refund_spike*100:.1f}%"
            })
    
    return anomalies

这套逻辑的好处:

  • 阈值可调-0.200.50),根据业务敏感度调整
  • 同时检测下跌和暴涨两种异常
  • 返回结构化数据,方便后续推送和记录

五、Telegram 推送模块

用 Python 的 requests 库调用 Telegram Bot API,把异常信息推送到频道或群组。

import requests

TELEGRAM_BOT_TOKEN = "YOUR_BOT_TOKEN"
CHAT_ID = "YOUR_CHAT_ID"  # 可以是群组ID或频道ID

def send_telegram(message):
    url = f"https://api.telegram.org/bot{TELEGRAM_BOT_TOKEN}/sendMessage"
    payload = {
        'chat_id': CHAT_ID,
        'text': message,
        'parse_mode': 'HTML'
    }
    try:
        r = requests.post(url, json=payload, timeout=10)
        return r.json().get('ok', False)
    except Exception as e:
        print(f"推送失败: {e}")
        return False

def notify_anomalies(anomalies):
    if not anomalies:
        send_telegram(
            f"✅ <b>每日数据巡检完成</b>\n\n"
            f"📅 {datetime.now().strftime('%Y-%m-%d')}\n"
            f"所有指标正常,无需处理。\n\n"
            f"— DuckDB 监控卫士"
        )
        return
    
    lines = [
        f"🚨 <b>数据异常预警</b>",
        f"📅 {datetime.now().strftime('%Y-%m-%d %H:%M')}",
        ""
    ]
    
    for a in anomalies:
        lines.append(f"{a['severity']} {a['metric']}")
        lines.append(f"   当前值:{a['value']}")
        lines.append(f"   预期范围:{a['expected']}")
        lines.append(f"   {a['message']}")
        lines.append("")
    
    lines.append("— DuckDB 监控卫士")
    
    send_telegram("\n".join(lines))

六、完整脚本与定时调度

现在把以上模块组合成完整的监控脚本:

import os
import glob
import json
from datetime import datetime

def incremental_update():
    """增量更新:只导入新增的 CSV 文件"""
    today_str = datetime.now().strftime('%Y%m%d')
    new_files = glob.glob(f'orders_{today_str}*.csv')
    
    if not new_files:
        print("今日无新数据文件")
        return
    
    for f in new_files:
        count = con.execute(f"SELECT COUNT(*) FROM read_csv_auto('{f}')").fetchone()[0]
        con.execute(f"""
            INSERT INTO orders
            SELECT * FROM read_csv_auto('{f}')
            WHERE order_id NOT IN (SELECT order_id FROM orders)
        """)
        print(f"已导入 {f}: {count} 条记录")

def main():
    # 1. 增量更新数据
    incremental_update()
    
    # 2. 执行健康检查
    stats = check_health()
    anomalies = detect_anomalies(stats)
    
    # 3. 推送结果
    notify_anomalies(anomalies)
    
    # 4. 记录今日日志
    log_entry = {
        'date': datetime.now().strftime('%Y-%m-%d'),
        'anomalies': len(anomalies),
        'stats': stats
    }
    with open('monitor_log.jsonl', 'a') as f:
        f.write(json.dumps(log_entry, ensure_ascii=False) + '\n')
    
    print(f"巡检完成,发现 {len(anomalies)} 条异常")

if __name__ == '__main__':
    main()

用 cron 定时执行(Linux):

# 每天早上 9 点自动运行
0 9 * * * cd /home/user/duckdb-monitor && python3 monitor.py >> monitor.log 2>&1

或者用 macOS 的 launchd,Windows 的计划任务,原理一样。

七、进阶:多数据源联合监控

真正的数据监控往往需要跨数据源。比如你的电商数据在 Shopify API,物流数据在快递 100,财务数据在 Excel。

DuckDB 的强项就在这里——不用先把数据搬到同一个数据库,直接 ATTACH 不同来源:

# 附件:Shopify 订单(CSV)
con.execute("ATTACH 'shopify_orders.csv' AS shopify (READ_ONLY)")

# 附件:物流数据(JSON)
con.execute("ATTACH 'logistics.json' AS log (READ_ONLY, TYPE JSON)")

# 附件:Excel 财务数据
con.execute("ATTACH 'finance.xlsx' AS fin (READ_ONLY)")

# 跨源联合查询:找出物流延误且退款率高的 SKU
con.execute("""
    SELECT 
        s.sku,
        s.product_name,
        s.order_count,
        AVG(l.transit_days) AS avg_transit,
        SUM(CASE WHEN f.status = 'refunded' THEN 1 ELSE 0 END) AS refund_count
    FROM shopify.orders s
    LEFT JOIN log.packages l ON s.tracking_number = l.tracking_id
    LEFT JOIN fin.refunds r ON s.order_id = r.order_id
    GROUP BY s.sku, s.product_name
    HAVING avg_transit > 7 AND refund_count > 5
    ORDER BY refund_count DESC
""").fetchdf()

ATTACH 是 DuckDB 最被低估的功能之一。它让你在一个查询里同时访问 CSV、JSON、Excel、Parquet、PostgreSQL、MySQL 等任意数据源,不需要 ETL 预处理。

八、与传统监控工具对比

功能DuckDB + Python 方案Grafana + Prometheus商业 SaaS(如 Datadog)
数据源接入CSV/JSON/Excel 直接读需 Exporter 采集API 对接,配置复杂
异常检测Python 自定义逻辑需配置告警规则内置 AI 异常检测
推送通知Telegram Bot APISlack/PagerDuty多渠道推送
部署成本0 元(本地/云服务器)需自建基础设施$20-100/节点/月
学习曲线Python + SQL
定制能力完全可控
适用场景中小团队、个人开发者大型基础设施企业级监控

九、变现建议

这套系统不只是自用工具,还可以变成收入来源:

🟢 低成本方案(¥0-5000 启动)

  • 模式:帮 10 个电商卖家部署这套系统,每个收 299 元/月,年费 2999 元
  • 具体步骤
    1. 把代码模板化,客户只需替换数据源路径
    2. 部署到便宜云服务器(如阿里云 ECS 最低配 ¥50/月)
    3. 每个客户单独 Telegram 群组,独立配置
  • 预期月收入:¥2990-5000
  • 适合人群:有 1-2 个客户资源的个人开发者

🟡 中成本方案(¥5000-50000 启动)

  • 模式:数据产品订阅。把监控结果做成每日早报,通过 Telegram 频道推送,定价 99 元/月
  • 具体步骤
    1. 选取特定行业(如跨境电商、股票)的数据源
    2. 搭建通用监控管道,自动生成行业分析报告
    3. 通过社交媒体引流,Telegram 频道付费订阅
  • 预期月收入:¥5000-20000(100-200 订阅用户)
  • 适合人群:有一定行业资源的内容创作者

🔴 高成本方案(¥50000+ 启动)

  • 模式:嵌入你的咨询/代运营服务。给客户提供「数据健康监测」作为增值项,提高客单价
  • 具体步骤
    1. 在现有咨询服务中加入数据监控模块
    2. 每月收取 ¥500-2000 的维护费
    3. 配合定期优化告警阈值和指标
  • 预期月收入:¥10000-50000(取决于客户数量)
  • 适合人群:已有客户基础的自由职业者或小型工作室

核心逻辑:你不是在卖代码,你是在卖「不再半夜被数据问题惊醒」的安心感

十、总结

组件技术选型成本
数据查询DuckDB(Python API)免费
数据存储.duckdb 单文件免费
异常检测Python 逻辑 + SQL 窗口函数免费
消息推送Telegram Bot API免费
定时调度Linux cron / 云服务器免费
总成本0 元

对比商业监控工具(Grafana + PagerDuty,年费数千到数万),这套方案在功能上不弱,成本为零。

记住:监控的本质不是技术,是「提前发现问题的意识」。DuckDB 只是让这件事变得极其简单。


📖 本文完整可运行代码(含多数据源 ATTACH 示例、Telegram 模板、cron 配置指南)已发布在 duckdblab.org,包含从 0 到部署的完整步骤。

💡 想系统学习 DuckDB 在监控、自动化、数据产品方向的应用?duckdblab.org 上有完整教程系列。

📺 Watch video tutorials → Olap Studio YouTube

Subscribe for more DuckDB & AI automation tutorials

使用 Hugo 构建
主题 StackJimmy 设计

⚠️ 本站为独立社区项目,与 DuckDB 基金会及 DuckDB 官方项目无任何从属、背书或赞助关系。

"DuckDB" 是 DuckDB 基金会的注册商标,本站仅以事实描述方式使用该名称。

本站内容仅供教育与社区推广用途,不构成任何商业服务。