Featured image of post Build a Data Anomaly Detection & Auto-Alert System with DuckDB — A SME SaaS Prototype

Build a Data Anomaly Detection & Auto-Alert System with DuckDB — A SME SaaS Prototype

Learn to build a fully automated data anomaly detection system with DuckDB: threshold detection, YoY/MoM analysis, category-level anomalies, plus Telegram/WeChat auto-alerts. Zero cost, sellable as a SaaS to SMEs.

DuckDB Anomaly Detection System Architecture

The Pain: Businesses Don’t Lack Data — They Lack “Real-Time Awareness”

Many small and medium enterprises generate massive amounts of business data every day — orders, traffic, inventory, user behavior — but managers either ignore the data or only discover problems after damage is done.

“Why did sales drop by half yesterday?” “Why is the churn rate higher this month?” “Is there a problem with the supply chain?”

The issue isn’t a lack of data — it’s a lack of real-time awareness. By the time someone manually checks and analyzes, the loss has already occurred.

Today I’ll show you how to build a Data Anomaly Detection & Auto-Alert System with DuckDB — scanning data every minute, detecting anomalies automatically, and sending instant notifications. The entire system costs less than $7/month.


Why DuckDB?

CapabilityDuckDBpandasSpark
Load 1M rows0.2 seconds3-5 secondsNeeds cluster
Window functionsOne SQL lineManual loopsComplex config
Multi-source integrationATTACH + JOINMultiple mergesETL preprocessing
Deployment complexitypip install duckdbpandas/numpy/etc.Hadoop/Spark cluster
Memory usageColumnar pruning, on-demandFull load to RAMDistributed memory

Anomaly detection的核心是 SQL 聚合和窗口函数——而这正是 DuckDB 最擅长的领域。

(The core of anomaly detection is SQL aggregation and window functions — exactly what DuckDB excels at.)


Step 1: Build the Data Layer

import duckdb

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

# Create orders table (e-commerce scenario)
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
)
""")

# Create pageviews table
con.execute("""
CREATE TABLE IF NOT EXISTS pageviews (
    pv_id BIGINT,
    event_time TIMESTAMP,
    page VARCHAR,
    source VARCHAR,
    device VARCHAR
)
""")

print("✅ Data schema created")

💡 Key insight: In real scenarios, you can directly ATTACH a client’s PostgreSQL, MySQL, or CSV files — DuckDB’s cross-database query capability makes data integration virtually zero-cost.


Step 2: Generate Simulated Data (with Injected Anomalies)

import random
from datetime import datetime, timedelta

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

def generate_order_data(days=30):
    """Generate 30 days of order data with injected anomalies"""
    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(['Electronics', 'Clothing', 'Food', 'Home']),
                random.choice(['MiniProgram', 'APP', 'Web', 'Offline']),
                random.randint(1000, 9999)
            ))
            order_id += 1

    # Inject anomaly 1: Day 15 — 60% sales drop
    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), 'Electronics', 'MiniProgram',
            random.randint(1000, 9999)
        ))
        order_id += 1

    # Inject anomaly 2: Day 22 — 300% spike in one category
    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), 'Electronics',
            random.choice(['MiniProgram', '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"✅ Inserted {len(orders_data)} order records")

Step 3: Core Engine — Three Anomaly Detection Algorithms

Detection 1: Threshold Anomaly (Standard Deviation)

def detect_threshold_anomaly(con):
    """Threshold-based anomaly detection using standard deviation"""
    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,
                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 '📉 Severe Low'
                WHEN order_count < moving_avg_count - 1 * moving_std_count
                THEN '⚠️ Slight Low'
                WHEN order_count > moving_avg_count + 2 * moving_std_count
                THEN '📈 Severe High'
                WHEN order_count > moving_avg_count + 1 * moving_std_count
                THEN '🟡 Slight High'
                ELSE '✅ Normal'
            END AS status
        FROM stats_with_window
        ORDER BY dt DESC
        LIMIT 7
    """).fetchdf()
    return result

Detection 2: Year-over-Year / Month-over-Month

def detect_time_anomaly(con):
    """YoY/MoM anomaly detection"""
    result = con.execute("""
        WITH today_stats AS (
            SELECT COUNT(*) AS order_count
            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
            'Day-over-Day' 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 '🚨 Cliff-like Drop!'
                WHEN (t.order_count - y.order_count) * 100.0 / NULLIF(y.order_count, 0) > 20
                THEN '🚀 Abnormal Surge!'
                ELSE '✅ Normal Fluctuation'
            END AS alert
        FROM today_stats t, yesterday_stats y
        UNION ALL
        SELECT
            'Week-over-Week' 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 '🚨 Significant YoY Drop!'
                WHEN (t.order_count - l.order_count) * 100.0 / NULLIF(l.order_count, 0) > 15
                THEN '🚀 Significant YoY Surge!'
                ELSE '✅ Normal Fluctuation'
            END AS alert
        FROM today_stats t, last_week_same_day l
    """).fetchdf()
    return result

Detection 3: Category-Level Anomaly

def detect_category_anomaly(con):
    """Anomaly detection by product category"""
    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 '🚨 Severely Below Expected'
                WHEN c.today_count > a.avg_daily_count + 2 * a.std_daily_count
                THEN '🚀 Severely Above Expected'
                ELSE '✅ Normal'
            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

💡 Key insight: These three detections cover the most common anomaly types — absolute value anomalies (threshold), trend anomalies (YoY/MoM), and structural anomalies (segment-level). A system that detects all three simultaneously far exceeds the coverage of any single method.


Step 4: Alert Report Generation & Push

import requests
import json
from datetime import datetime

def generate_alert_report():
    """Generate a complete anomaly detection 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 = []

    # Extract anomalies from threshold detection
    for _, row in threshold_alerts.iterrows():
        if row['status'] != '✅ Normal':
            all_alerts.append({
                'type': 'Daily Order Anomaly',
                'detail': f"{row['dt']} orders: {int(row['order_count'])}, expected {int(row['expected_count'])}±{int(row['std_dev'])}",
                'severity': '🚨 Severe' if 'Severe' in row['status'] else '⚠️ Slight',
                'source': 'Threshold Detection'
            })

    # Extract anomalies from time series detection
    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"Change {row['change_pct']}% — {row['alert']}",
                'severity': '🚨 High Priority',
                'source': 'Time Series Detection'
            })

    # Extract anomalies from category detection
    for _, row in category_alerts.iterrows():
        if '🚨' in str(row['status']) or '🚀' in str(row['status']):
            all_alerts.append({
                'type': f"Category Anomaly: {row['category']}",
                'detail': f"Today: {int(row['today_count'])} orders, expected {int(row['expected_count'])} ({row['change_pct']}% change)",
                'severity': '🚨 High Priority' if 'Severe' in row['status'] else '⚠️ Medium',
                'source': 'Category Detection'
            })

    severity_order = {'🚨 High Priority': 0, '🚨 Severe': 1, '⚠️ Medium': 2, '⚠️ Slight': 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):
    """Send Telegram alert message"""
    if not alerts:
        message = f"✅ Data Monitoring Normal ({datetime.now().strftime('%Y-%m-%d %H:%M')})\n\nNo anomalies detected."
    else:
        lines = [f"🚨 Data Anomaly Alert ({datetime.now().strftime('%Y-%m-%d %H:%M')})"]
        lines.append(f"Total {len(alerts)} 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()

# Run monitoring
alerts = generate_alert_report()
print(f"🔔 Found {len(alerts)} alerts")
for a in alerts[:5]:
    print(f"  [{a['severity']}] {a['type']}: {a['detail']}")

Step 5: Scheduled Execution

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} severe alerts — notification sent")
    except Exception as e:
        print(f"❌ Monitoring failed: {e}")

print("🚀 Starting anomaly detection system...")
scheduled_monitoring()

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

try:
    while True:
        schedule.run_pending()
        time.sleep(30)
except KeyboardInterrupt:
    print("\n🛑 System stopped")

Cron configuration (production recommended):

# Run every 5 minutes
*/5 * * * * cd /home/user/anomaly-detector && python3 monitor.py >> monitor.log 2>&1

Advanced: Cross-Source Joint Detection

Real scenarios require cross-source data. DuckDB’s ATTACH lets you query CSV, JSON, Excel, PostgreSQL simultaneously in one query:

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

# Attach: Shopify orders (CSV)
con.execute("ATTACH 'shopify_orders.csv' AS shopify (READ_ONLY)")

# Attach: Logistics data (JSON)
con.execute("ATTACH 'logistics.json' AS log (READ_ONLY, TYPE JSON)")

# Attach: Finance data (Excel)
con.execute("ATTACH 'finance.xlsx' AS fin (READ_ONLY)")

# Cross-source query: Find SKUs with delayed shipping AND high refund rate
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()

Business Model: From Code to Revenue

PlanPricingTarget Customer
Basic SaaS$70/mo, 3 metrics, daily reportSmall e-commerce
Pro SaaS$210/mo, 10 metrics, real-time alerts + WeChat pushMedium enterprise
Enterprise$420/mo, custom metrics + history + APILarge enterprise
One-time Deployment$400-1100/projectCustom needs
Pay-per-Alert$0.01/alert, $28/mo minimumLow-frequency clients

Freelance analyst feasibility:

  • Serve 15 SMEs simultaneously = $1,050-3,150/month revenue
  • System runs automatically, < 5 hours maintenance/month
  • Near-zero marginal cost (server ~$7/month)

Core logic: You’re not selling code — you’re selling “peace of mind, no more waking up to data emergencies at midnight.”


Monetization Advice

  1. Start with nearby clients: Find 3-5 e-commerce friends, deploy for free in exchange for case studies and referrals.
  2. Package as template product: Make the code configurable — clients only need to swap data source paths and alert thresholds.
  3. Value-added services: Offer “monthly data health reports” as an upsell to increase average revenue per customer.
  4. Subscription lock-in: Charge monthly rather than one-time to build recurring revenue.
  5. Content marketing: Publish full tutorials on duckdblab.org to attract qualified traffic and convert to paying customers.
  6. API transformation: Encapsulate monitoring as a REST API for other developers to integrate, expanding into B2B2C.

Full runnable code is available at duckdblab.org, from setup to production deployment.

Olap Studio · Focused on DuckDB Practical Skills · 2026-09-28

📺 Watch video tutorials → Olap Studio YouTube

Subscribe for more DuckDB & AI automation tutorials

Built with Hugo
Theme Stack designed by Jimmy