用 DuckDB + Cron 搭建全自动数据监控报警系统,睡后收入的新思路
你有没有这样的痛点——业务数据每天都在涨,但出了问题却不知道。等客户来找你反馈,损失已经造成了。
今天教你用 DuckDB 搭建一套全自动的「数据健康监测系统」。每天定时跑查询,发现异常自动发 Telegram 通知。不用写复杂的 ETL,不用买昂贵的监控工具,一个 Python 脚本 + 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]
}
关键点:
read_csv_auto('orders_2024*.csv')— DuckDB 支持通配符一次读多个文件,不用手动拼接- 窗口函数
AVG(...) OVER (...)— 计算 7 日移动平均,不用预聚合 - 结果存在
.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.20、0.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 API | Slack/PagerDuty | 多渠道推送 |
| 部署成本 | 0 元(本地/云服务器) | 需自建基础设施 | $20-100/节点/月 |
| 学习曲线 | Python + SQL | 高 | 中 |
| 定制能力 | 完全可控 | 中 | 低 |
| 适用场景 | 中小团队、个人开发者 | 大型基础设施 | 企业级监控 |
九、变现建议
这套系统不只是自用工具,还可以变成收入来源:
🟢 低成本方案(¥0-5000 启动)
- 模式:帮 10 个电商卖家部署这套系统,每个收 299 元/月,年费 2999 元
- 具体步骤:
- 把代码模板化,客户只需替换数据源路径
- 部署到便宜云服务器(如阿里云 ECS 最低配 ¥50/月)
- 每个客户单独 Telegram 群组,独立配置
- 预期月收入:¥2990-5000
- 适合人群:有 1-2 个客户资源的个人开发者
🟡 中成本方案(¥5000-50000 启动)
- 模式:数据产品订阅。把监控结果做成每日早报,通过 Telegram 频道推送,定价 99 元/月
- 具体步骤:
- 选取特定行业(如跨境电商、股票)的数据源
- 搭建通用监控管道,自动生成行业分析报告
- 通过社交媒体引流,Telegram 频道付费订阅
- 预期月收入:¥5000-20000(100-200 订阅用户)
- 适合人群:有一定行业资源的内容创作者
🔴 高成本方案(¥50000+ 启动)
- 模式:嵌入你的咨询/代运营服务。给客户提供「数据健康监测」作为增值项,提高客单价
- 具体步骤:
- 在现有咨询服务中加入数据监控模块
- 每月收取 ¥500-2000 的维护费
- 配合定期优化告警阈值和指标
- 预期月收入:¥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 上有完整教程系列。