
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?
| Capability | DuckDB | pandas | Spark |
|---|---|---|---|
| Load 1M rows | 0.2 seconds | 3-5 seconds | Needs cluster |
| Window functions | One SQL line | Manual loops | Complex config |
| Multi-source integration | ATTACH + JOIN | Multiple merges | ETL preprocessing |
| Deployment complexity | pip install duckdb | pandas/numpy/etc. | Hadoop/Spark cluster |
| Memory usage | Columnar pruning, on-demand | Full load to RAM | Distributed 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
| Plan | Pricing | Target Customer |
|---|---|---|
| Basic SaaS | $70/mo, 3 metrics, daily report | Small e-commerce |
| Pro SaaS | $210/mo, 10 metrics, real-time alerts + WeChat push | Medium enterprise |
| Enterprise | $420/mo, custom metrics + history + API | Large enterprise |
| One-time Deployment | $400-1100/project | Custom needs |
| Pay-per-Alert | $0.01/alert, $28/mo minimum | Low-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
- Start with nearby clients: Find 3-5 e-commerce friends, deploy for free in exchange for case studies and referrals.
- Package as template product: Make the code configurable — clients only need to swap data source paths and alert thresholds.
- Value-added services: Offer “monthly data health reports” as an upsell to increase average revenue per customer.
- Subscription lock-in: Charge monthly rather than one-time to build recurring revenue.
- Content marketing: Publish full tutorials on duckdblab.org to attract qualified traffic and convert to paying customers.
- 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