
为什么你需要增量更新?
每天早上 9 点,你都要花 1-2 小时重新跑一遍昨天的销售数据分析——导 CSV、清洗、聚合、生成报表、发微信群。日复一日,年复一年。
这不是"勤奋",是工具效率太低。
假设你每天处理 10 万条订单数据,每个月就是 300 万条。全量重跑的问题:
- 每天都要读取全部历史数据,浪费 90%+ 的计算资源
- 报表生成时间随数据量线性增长,半年后可能要从 30 秒变成 10 分钟
- 无法做到"实时"——只能在每天固定时间点跑一次
增量更新的核心思路:不是每次都全量重跑,而是只处理新增的数据,然后增量合并到主表中。
初始化数据库与表结构
import duckdb
from datetime import datetime, timedelta
import random
# 创建数据库(首次运行时)
con = duckdb.connect("daily_report.db")
# 创建订单主表
con.execute("""
CREATE TABLE IF NOT EXISTS orders (
order_id VARCHAR PRIMARY KEY,
order_date DATE,
product_id VARCHAR,
product_name VARCHAR,
quantity INTEGER,
unit_price DECIMAL(10,2),
region VARCHAR,
channel VARCHAR,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
""")
# 创建每日聚合报表表
con.execute("""
CREATE TABLE IF NOT EXISTS daily_summary (
report_date DATE PRIMARY KEY,
total_orders INTEGER,
total_revenue DECIMAL(12,2),
avg_order_value DECIMAL(10,2),
top_product VARCHAR,
top_product_revenue DECIMAL(12,2),
region_revenue JSON,
channel_revenue JSON,
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
""")
# 创建增量数据临时表
con.execute("""
CREATE TABLE IF NOT EXISTS daily_staging (
order_id VARCHAR,
order_date DATE,
product_id VARCHAR,
product_name VARCHAR,
quantity INTEGER,
unit_price DECIMAL(10,2),
region VARCHAR,
channel VARCHAR
);
""")
print("✅ 数据库初始化完成")
模拟每日新增数据
def generate_daily_orders(days=30):
"""生成指定天数的模拟订单数据"""
products = [
('P001', 'iPhone 15', 7999),
('P002', 'MacBook Pro', 14999),
('P003', 'AirPods Pro', 1899),
('P004', 'iPad Air', 4799),
('P005', 'Apple Watch', 2999),
('P006', 'AirTag', 229),
('P007', 'Magic Keyboard', 999),
('P008', 'Studio Display', 11999),
]
regions = ['华东', '华南', '华北', '西南', '华中', '东北', '西北']
channels = ['天猫', '京东', '拼多多', '抖音', '自营APP']
all_orders = []
for day_offset in range(days):
date = datetime.now() - timedelta(days=day_offset)
num_orders = random.randint(80, 200)
for _ in range(num_orders):
product = random.choice(products)
quantity = random.randint(1, 5)
order_id = f"ORD{date.strftime('%Y%m%d')}{random.randint(1000, 9999)}"
all_orders.append({
'order_id': order_id,
'order_date': date,
'product_id': product[0],
'product_name': product[1],
'quantity': quantity,
'unit_price': product[2],
'region': random.choice(regions),
'channel': random.choice(channels),
})
return all_orders
# 生成 30 天历史数据
print("📊 生成历史数据...")
all_orders = generate_daily_orders(30)
print(f"✅ 共生成 {len(all_orders)} 条订单记录")
核心:MERGE UPSERT 幂等加载
第一次运行时,用 MERGE 确保幂等性——即使重复运行也不会产生重复数据:
def initial_load(orders):
"""首次全量加载数据"""
# 批量插入临时表
for order in orders:
con.execute("""
INSERT INTO daily_staging
(order_id, order_date, product_id, product_name, quantity, unit_price, region, channel)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
""", [
order['order_id'], order['order_date'], order['product_id'],
order['product_name'], order['quantity'], order['unit_price'],
order['region'], order['channel']
])
# 使用 MERGE 进行幂等加载(UPSERT)
con.execute("""
MERGE INTO orders AS target
USING daily_staging AS source
ON target.order_id = source.order_id
WHEN NOT MATCHED THEN
INSERT (order_id, order_date, product_id, product_name,
quantity, unit_price, region, channel)
VALUES (source.order_id, source.order_date, source.product_id,
source.product_name, source.quantity, source.unit_price,
source.region, source.channel)
WHEN MATCHED THEN
UPDATE SET
quantity = source.quantity,
unit_price = source.unit_price,
updated_at = CURRENT_TIMESTAMP
""")
# 清空 staging 表
con.execute("TRUNCATE TABLE daily_staging")
loaded = con.execute("SELECT COUNT(*) FROM orders").fetchone()[0]
print(f"✅ 首次加载完成,共 {loaded} 条订单")
return loaded
增量更新:只处理新增数据
这是整个系统的核心——每天只处理新增数据,然后用 MERGE 增量更新报表:
def incremental_update(new_orders):
"""增量更新:只处理新增数据,更新报表"""
# 1. 将新数据插入 staging
for order in new_orders:
con.execute("""
INSERT INTO daily_staging
(order_id, order_date, product_id, product_name, quantity, unit_price, region, channel)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
""", [
order['order_id'], order['order_date'], order['product_id'],
order['product_name'], order['quantity'], order['unit_price'],
order['region'], order['channel']
])
# 2. 增量加载到主表(只插入不存在的记录)
con.execute("""
MERGE INTO orders AS target
USING daily_staging AS source
ON target.order_id = source.order_id
WHEN NOT MATCHED THEN
INSERT (order_id, order_date, product_id, product_name,
quantity, unit_price, region, channel)
VALUES (source.order_id, source.order_date, source.product_id,
source.product_name, source.quantity, source.unit_price,
source.region, source.channel)
""")
con.execute("TRUNCATE TABLE daily_staging")
# 3. 增量更新日报表(关键步骤)
con.execute("""
MERGE INTO daily_summary AS target
USING (
SELECT
order_date AS report_date,
COUNT(*) AS total_orders,
ROUND(SUM(quantity * unit_price), 2) AS total_revenue,
ROUND(AVG(quantity * unit_price), 2) AS avg_order_value,
(SELECT product_name FROM orders o2
WHERE o2.order_date = orders.order_date
GROUP BY product_name
ORDER BY SUM(quantity * unit_price) DESC
LIMIT 1) AS top_product,
(SELECT ROUND(SUM(quantity * unit_price), 2)
FROM orders o2
WHERE o2.order_date = orders.order_date
GROUP BY product_name
ORDER BY SUM(quantity * unit_price) DESC
LIMIT 1) AS top_product_revenue,
(SELECT json_group_object(region, ROUND(SUM(quantity * unit_price), 2))
FROM orders o3 WHERE o3.order_date = orders.order_date
) AS region_revenue,
(SELECT json_group_object(channel, ROUND(SUM(quantity * unit_price), 2))
FROM orders o4 WHERE o4.order_date = orders.order_date
) AS channel_revenue
FROM orders
WHERE order_date = (SELECT MAX(order_date) FROM orders)
GROUP BY order_date
) AS source
ON target.report_date = source.report_date
WHEN NOT MATCHED THEN
INSERT (report_date, total_orders, total_revenue, avg_order_value,
top_product, top_product_revenue, region_revenue, channel_revenue)
VALUES (source.report_date, source.total_orders, source.total_revenue,
source.avg_order_value, source.top_product, source.top_product_revenue,
source.region_revenue, source.channel_revenue)
WHEN MATCHED THEN
UPDATE SET
total_orders = source.total_orders,
total_revenue = source.total_revenue,
avg_order_value = source.avg_order_value,
top_product = source.top_product,
top_product_revenue = source.top_product_revenue,
region_revenue = source.region_revenue,
channel_revenue = source.channel_revenue,
updated_at = CURRENT_TIMESTAMP
""")
print("✅ 增量更新完成")
性能对比
| 指标 | 全量重跑方案 | 增量更新方案 |
|---|---|---|
| 每日处理数据量 | 30,000 条(全量) | 1,000 条(新增) |
| 报表生成时间 | 2 小时 → 随时间增长 | < 30 秒 → 恒定 |
| 资源消耗 | 高(每次都全量读取) | 低(只处理新增) |
| 数据一致性 | 每次重新计算 | 增量合并 + 每周全量校验 |
| 适用场景 | 数据量 < 1万条 | 数据量持续增长 |
对于月处理 30 万条数据的场景,增量更新可以将每日报表生成时间从 30 秒压缩到 2 秒以内。
进阶:物化视图加速
对于需要频繁查询的报表,可以创建物化视图缓存结果:
-- 创建物化视图缓存近 30 天的日报
CREATE MATERIALIZED VIEW IF NOT EXISTS mv_30d_summary AS
SELECT
report_date,
total_orders,
total_revenue,
avg_order_value,
top_product,
top_product_revenue,
region_revenue,
channel_revenue
FROM daily_summary
WHERE report_date >= CURRENT_DATE - INTERVAL '30' DAY;
查询时直接从物化视图读取,速度提升 10 倍以上:
SELECT * FROM mv_daily_trend ORDER BY report_date DESC LIMIT 7;
自动化调度
#!/bin/bash
# daily_report_incremental.sh — 每天凌晨 3 点自动增量更新
cd ~/daily-report-system
# 1. 拉取昨日新增数据(从 API 或数据仓库同步)
python3 fetch_yesterday_orders.py
# 2. 执行增量更新
python3 incremental_update.py
# 3. 生成今日简报并推送
python3 send_report.py
# 4. 记录运行日志
echo "[$(date)] 增量更新完成" >> /var/log/daily_report.log
crontab 配置:
# 工作日凌晨 3 点自动增量更新
0 3 * * 1-5 /home/user/daily-report-system/daily_incremental.sh
# 每周日凌晨 4 点执行全量校验(确保数据一致性)
0 4 * * 0 /home/user/daily-report-system/daily_full_verify.py
💰 变现建议
卖报表服务:很多中小企业没有数据团队,每月花 2000-5000 元请人跑报表。你这套系统一次部署,每月收 500-1000 元维护费,10 个客户 = 月入 5000-10000 元。
卖自动化系统:把这套系统打包成"日报自动生成器",一次性收费 3000-8000 元,面向有数据需求的中小企业。
做 SaaS 产品:把增量更新 + 物化视图的方案做成 SaaS,按月订阅收费,客户自己上传数据,系统自动更新报表。
核心卖点就一句话:“你的日报,从 2 小时变成 30 秒。”
🎯 今晚行动
- 创建一个 DuckDB 数据库,按照上面的代码搭建基础框架
- 用你自己的真实数据替换模拟数据,测试增量更新效果
- 配上 crontab,让它每天自动运行
- 把生成的报表发到一个微信群或企业微信,感受"自动化"的爽感
下一步可以探索:把增量更新的结果通过 FastAPI 暴露成 API,做成实时数据产品。
📺 更多 DuckDB 实战教程,订阅 YouTube 频道 → youtube.com/@duckdblab