Featured image of post 用 DuckDB 搭建数据异常检测与自动报警系统——中小企业 SaaS 原型

用 DuckDB 搭建数据异常检测与自动报警系统——中小企业 SaaS 原型

教你用 DuckDB 搭一套全自动数据异常检测系统:阈值检测、同比环比分析、品类细分异常,配合 Telegram/企微机器人自动报警。零成本,可卖给中小企业当 SaaS 产品。

DuckDB 数据异常检测系统架构图

痛点:企业缺的不是数据,是「第一时间感知」

很多中小企业每天都在产生大量业务数据——订单、流量、库存、用户行为……但管理者要么不看数据,要么等出问题了才发现。

“昨天销售额怎么突然掉了一半?” “这个月用户流失率怎么比上个月高这么多?” “供应链那边是不是出问题了?”

这些问题不是因为没有数据,而是因为没有实时感知。等到手动去查、去分析,损失已经发生了。

今天教你用 DuckDB 搭一个「数据异常检测与自动报警系统」——每分钟自动扫描数据,发现异常立即通知,整套系统成本不到 ¥50/月。


为什么选择 DuckDB?

能力DuckDBpandasSpark
100 万行数据加载0.2 秒3-5 秒需要集群
窗口函数计算一行 SQL需手写循环复杂配置
多源数据整合ATTACH + JOIN多次 mergeETL 预处理
部署复杂度pip install duckdbpandas/numpy 等多库Hadoop/Spark 集群
内存占用列式裁剪,按需读取全量加载到 RAM分布式内存

异常检测的核心是 SQL 聚合和窗口函数——而这正是 DuckDB 最擅长的领域。


第一步:搭建数据层

import duckdb

con = duckdb.connect("anomaly_detector.db")

# 创建订单表(模拟电商场景)
con.execute("""
CREATE TABLE IF NOT EXISTS orders (
    order_id BIGINT,
    order_time TIMESTAMP,
    amount DECIMAL(10,2),
    category VARCHAR,
    channel VARCHAR,
    customer_id BIGINT
)
""")

# 创建流量表
con.execute("""
CREATE TABLE IF NOT EXISTS pageviews (
    pv_id BIGINT,
    event_time TIMESTAMP,
    page VARCHAR,
    source VARCHAR,
    device VARCHAR
)
""")

print("✅ 数据表结构已创建")

💡 关键:真实场景中,你可以直接 ATTACH 客户的 PostgreSQL、MySQL、CSV 文件——DuckDB 的跨库查询能力让数据接入几乎为零成本。


第二步:注入模拟数据(含人工异常)

为了让系统有东西可检测,我们生成包含正常波动和人工异常的数据:

import random
from datetime import datetime, timedelta

random.seed(42)
con = duckdb.connect("anomaly_detector.db")

def generate_order_data(days=30):
    """生成 30 天的订单数据,含人工注入异常"""
    orders = []
    base_time = datetime(2026, 8, 1)
    order_id = 1

    for day in range(days):
        current_date = base_time + timedelta(days=day)
        is_weekend = current_date.weekday() >= 5
        base_orders = 80 if is_weekend else 120
        daily_orders = int(base_orders * random.uniform(0.8, 1.2))

        for _ in range(daily_orders):
            hour = random.randint(8, 23)
            minute = random.randint(0, 59)
            orders.append((
                order_id,
                current_date.replace(hour=hour, minute=minute),
                round(random.uniform(29, 599), 2),
                random.choice(['电子产品', '服装', '食品', '家居']),
                random.choice(['小程序', 'APP', '网页', '线下']),
                random.randint(1000, 9999)
            ))
            order_id += 1

    # 注入异常 1:第 15 天销售额突降 60%
    abnormal_date = base_time + timedelta(days=14)
    for i in range(10):
        orders.append((
            order_id, abnormal_date.replace(hour=10+i, minute=random.randint(0,59)),
            round(random.uniform(29, 599), 2), '电子产品', '小程序',
            random.randint(1000, 9999)
        ))
        order_id += 1

    # 注入异常 2:第 22 天某品类销量暴增 300%
    spike_date = base_time + timedelta(days=21)
    for i in range(50):
        orders.append((
            order_id, spike_date.replace(hour=random.randint(9,21), minute=random.randint(0,59)),
            round(random.uniform(29, 599), 2), '电子产品',
            random.choice(['小程序', 'APP']), random.randint(1000, 9999)
        ))
        order_id += 1

    return orders

orders_data = generate_order_data(30)
con.execute("INSERT INTO orders VALUES ?", orders_data)
print(f"✅ 插入了 {len(orders_data)} 条订单数据")

第三步:核心引擎——三类异常检测算法

检测 1:阈值异常(标准差检测)

def detect_threshold_anomaly(con):
    """基于标准差的阈值异常检测"""
    result = con.execute("""
        WITH daily_stats AS (
            SELECT
                DATE(order_time) AS dt,
                COUNT(*) AS order_count,
                SUM(amount) AS total_revenue,
                AVG(amount) AS avg_order_value
            FROM orders
            GROUP BY DATE(order_time)
        ),
        stats_with_window AS (
            SELECT
                dt,
                order_count,
                total_revenue,
                AVG(order_count) OVER (
                    ORDER BY dt
                    ROWS BETWEEN 6 PRECEDING AND 1 PRECEDING
                ) AS moving_avg_count,
                STDDEV(order_count) OVER (
                    ORDER BY dt
                    ROWS BETWEEN 6 PRECEDING AND 1 PRECEDING
                ) AS moving_std_count
            FROM daily_stats
        )
        SELECT
            dt,
            order_count,
            ROUND(moving_avg_count, 0) AS expected_count,
            ROUND(moving_std_count, 0) AS std_dev,
            CASE
                WHEN order_count < moving_avg_count - 2 * moving_std_count
                THEN '📉 严重偏低'
                WHEN order_count < moving_avg_count - 1 * moving_std_count
                THEN '⚠️ 轻微偏低'
                WHEN order_count > moving_avg_count + 2 * moving_std_count
                THEN '📈 严重偏高'
                WHEN order_count > moving_avg_count + 1 * moving_std_count
                THEN '🟡 轻微偏高'
                ELSE '✅ 正常'
            END AS status
        FROM stats_with_window
        ORDER BY dt DESC
        LIMIT 7
    """).fetchdf()
    return result

检测 2:同比环比异常

def detect_time_anomaly(con):
    """基于同比/环比的异常检测"""
    result = con.execute("""
        WITH today_stats AS (
            SELECT
                COUNT(*) AS order_count,
                SUM(amount) AS total_revenue
            FROM orders
            WHERE DATE(order_time) = (SELECT MAX(DATE(order_time)) FROM orders)
        ),
        yesterday_stats AS (
            SELECT COUNT(*) AS order_count
            FROM orders
            WHERE DATE(order_time) = (
                SELECT MAX(DATE(order_time)) - INTERVAL '1' DAY FROM orders
            )
        ),
        last_week_same_day AS (
            SELECT COUNT(*) AS order_count
            FROM orders
            WHERE DATE(order_time) = (
                SELECT MAX(DATE(order_time)) - INTERVAL '7' DAY FROM orders
            )
        )
        SELECT
            '日环比' AS compare_type,
            ROUND(
                (t.order_count - y.order_count) * 100.0 / NULLIF(y.order_count, 0),
                1
            ) AS change_pct,
            CASE
                WHEN (t.order_count - y.order_count) * 100.0 / NULLIF(y.order_count, 0) < -20
                THEN '🚨 断崖式下跌!'
                WHEN (t.order_count - y.order_count) * 100.0 / NULLIF(y.order_count, 0) > 20
                THEN '🚀 异常增长!'
                ELSE '✅ 正常波动'
            END AS alert
        FROM today_stats t, yesterday_stats y
        UNION ALL
        SELECT
            '周同比' AS compare_type,
            ROUND(
                (t.order_count - l.order_count) * 100.0 / NULLIF(l.order_count, 0),
                1
            ) AS change_pct,
            CASE
                WHEN (t.order_count - l.order_count) * 100.0 / NULLIF(l.order_count, 0) < -15
                THEN '🚨 同比大幅下降!'
                WHEN (t.order_count - l.order_count) * 100.0 / NULLIF(l.order_count, 0) > 15
                THEN '🚀 同比大幅增长!'
                ELSE '✅ 正常波动'
            END AS alert
        FROM today_stats t, last_week_same_day l
    """).fetchdf()
    return result

检测 3:品类异常(细分维度)

def detect_category_anomaly(con):
    """按品类检测异常"""
    result = con.execute("""
        WITH category_today AS (
            SELECT
                category,
                COUNT(*) AS today_count,
                SUM(amount) AS today_revenue
            FROM orders
            WHERE DATE(order_time) = (SELECT MAX(DATE(order_time)) FROM orders)
            GROUP BY category
        ),
        category_avg AS (
            SELECT
                category,
                AVG(daily_count) AS avg_daily_count,
                STDDEV(daily_count) AS std_daily_count
            FROM (
                SELECT
                    category,
                    DATE(order_time) AS dt,
                    COUNT(*) AS daily_count
                FROM orders
                WHERE DATE(order_time) >= (
                    SELECT MAX(DATE(order_time)) - INTERVAL '14' DAY FROM orders
                )
                GROUP BY category, DATE(order_time)
            ) sub
            GROUP BY category
        )
        SELECT
            c.category,
            c.today_count,
            ROUND(a.avg_daily_count, 0) AS expected_count,
            ROUND((c.today_count - a.avg_daily_count) * 100.0 / NULLIF(a.avg_daily_count, 0), 1) AS change_pct,
            CASE
                WHEN c.today_count < a.avg_daily_count - 2 * a.std_daily_count
                THEN '🚨 严重低于预期'
                WHEN c.today_count > a.avg_daily_count + 2 * a.std_daily_count
                THEN '🚀 严重高于预期'
                ELSE '✅ 正常'
            END AS status
        FROM category_today c
        JOIN category_avg a ON c.category = a.category
        ORDER BY ABS(change_pct) DESC
    """).fetchdf()
    return result

💡 关键洞察:这三种检测覆盖了最常见的异常类型——绝对值异常(阈值)、趋势异常(同比环比)、结构性异常(细分维度)。一个系统同时检测这三种,覆盖度远超单一方法。


第四步:汇总报告与自动推送

import requests
import json
from datetime import datetime

def generate_alert_report():
    """生成完整的异常检测报告"""
    con = duckdb.connect("anomaly_detector.db")

    threshold_alerts = detect_threshold_anomaly(con)
    time_alerts = detect_time_anomaly(con)
    category_alerts = detect_category_anomaly(con)

    all_alerts = []

    # 从阈值检测中提取异常
    for _, row in threshold_alerts.iterrows():
        if row['status'] != '✅ 正常':
            all_alerts.append({
                'type': '日订单量异常',
                'detail': f"{row['dt']} 订单量 {int(row['order_count'])},预期 {int(row['expected_count'])}±{int(row['std_dev'])}",
                'severity': '🚨 严重' if '严重' in row['status'] else '⚠️ 轻微',
                'source': '阈值检测'
            })

    # 从时间检测中提取异常
    for _, row in time_alerts.iterrows():
        if '🚨' in str(row['alert']) or '🚀' in str(row['alert']):
            all_alerts.append({
                'type': row['compare_type'],
                'detail': f"变化 {row['change_pct']}% —— {row['alert']}",
                'severity': '🚨 高优先级',
                'source': '时间序列检测'
            })

    # 从品类检测中提取异常
    for _, row in category_alerts.iterrows():
        if '🚨' in str(row['status']) or '🚀' in str(row['status']):
            all_alerts.append({
                'type': f"品类异常:{row['category']}",
                'detail': f"今日 {int(row['today_count'])} 单,预期 {int(row['expected_count'])} 单(变化 {row['change_pct']}%)",
                'severity': '🚨 高优先级' if '严重' in row['status'] else '⚠️ 中优先级',
                'source': '品类检测'
            })

    severity_order = {'🚨 高优先级': 0, '🚨 严重': 1, '⚠️ 中优先级': 2, '⚠️ 轻微': 3}
    all_alerts.sort(key=lambda x: severity_order.get(x['severity'], 99))

    return all_alerts

def send_telegram_alert(alerts, bot_token, chat_id):
    """发送 Telegram 报警消息"""
    if not alerts:
        message = f"✅ 数据监控正常({datetime.now().strftime('%Y-%m-%d %H:%M')})\n\n未发现异常指标。"
    else:
        lines = [f"🚨 数据异常报警({datetime.now().strftime('%Y-%m-%d %H:%M')})"]
        lines.append(f"共发现 {len(alerts)} 条告警:")
        lines.append("---")
        for a in alerts[:10]:
            lines.append(f"{a['severity']} {a['type']}")
            lines.append(f"  {a['detail']}")
        message = "\n".join(lines)

    url = f"https://api.telegram.org/bot{bot_token}/sendMessage"
    payload = {"chat_id": chat_id, "text": message, "parse_mode": "Markdown"}
    response = requests.post(url, json=payload)
    return response.json()

# 执行监控
alerts = generate_alert_report()
print(f"🔔 共发现 {len(alerts)} 条告警")
for a in alerts[:5]:
    print(f"  [{a['severity']}] {a['type']}: {a['detail']}")

第五步:定时调度——让系统自动运行

import schedule
import time

CHECK_INTERVAL_MINUTES = 5

def scheduled_monitoring():
    try:
        alerts = generate_alert_report()
        severe_count = sum(1 for a in alerts if '🚨' in a['severity'])
        if severe_count > 0:
            send_telegram_alert(alerts, BOT_TOKEN, CHAT_ID)
            print(f"⚠️ 发现 {severe_count} 条严重告警,已推送通知")
    except Exception as e:
        print(f"❌ 监控执行失败: {e}")

print("🚀 启动数据异常监控系统...")
scheduled_monitoring()

schedule.every(CHECK_INTERVAL_MINUTES).minutes.do(scheduled_monitoring)

try:
    while True:
        schedule.run_pending()
        time.sleep(30)
except KeyboardInterrupt:
    print("\n🛑 监控系统已停止")

用 cron 定时执行(生产环境推荐):

# 每 5 分钟执行一次
*/5 * * * * cd /home/user/anomaly-detector && python3 monitor.py >> monitor.log 2>&1

进阶:跨数据源联合检测

真实场景需要跨数据源。DuckDB 的 ATTACH 让你在一个查询里同时访问 CSV、JSON、Excel、PostgreSQL:

con = duckdb.connect("multi_source.db")

# 附件: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
result = 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_number
    LEFT JOIN fin.refunds f ON s.order_id = f.order_id
    GROUP BY s.sku, s.product_name
    HAVING avg_transit > 7 AND refund_count > 5
    ORDER BY refund_count DESC
""").fetchdf()

商业模式:从代码到收入

方案定价适合客户
基础版 SaaS¥500/月,监控 3 个指标,每日 1 次报告小型电商
专业版 SaaS¥1500/月,监控 10 个指标,实时报警 + 企微推送中型企业
企业版 SaaS¥3000/月,自定义指标 + 历史回溯 + API 对接大型企业
一次性交付¥3000-8000/项目定制化需求
按告警次数¥0.1/条告警,月费 ¥200 起步低频告警客户

一个自由分析师的可行性:

  • 同时服务 15 个中小企业 = ¥7500-22500/月收入
  • 系统自动化运行,每月维护时间 < 5 小时
  • 边际成本几乎为零(服务器 ¥50/月)

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


变现建议

  1. 从身边客户开始:找 3-5 个做电商的朋友,免费帮他们部署,换取案例和口碑。
  2. 打包成模板产品:把代码做成可配置模板,客户只需替换数据源路径和报警阈值。
  3. 增值服务:提供"每月数据健康报告"作为附加服务,提升客单价。
  4. 订阅制锁定:按月收费而非一次性,形成持续收入流。
  5. 内容营销:在 duckdblab.org 发布完整教程,吸引精准流量,转化为付费客户。
  6. API 化:将监控能力封装为 REST API,供其他开发者集成,拓展 B2B2C 模式。

本文完整可运行代码已发布在 duckdblab.org,包含从 0 到部署的完整步骤。

Olap 工作室 · 专注 DuckDB 实战 · 2026-09-28

📺 Watch video tutorials → Olap Studio YouTube

Subscribe for more DuckDB & AI automation tutorials

使用 Hugo 构建
主题 Stack 由 Jimmy 设计