30行代码用DuckDB搭建个人数据产品后端,月入过万不是梦
很多人想用 DuckDB 做数据产品,但卡在架构设计上。今天拆解一个完整案例——用 DuckDB + FastAPI 搭建一个「电商销售分析 API」,支持多用户隔离,按查询量收费。

一、项目背景:为什么这个方向能赚钱?
假设你接到了一个需求:为一个小商家提供每日销售数据 API。商家有 10 个店铺,每天产生约 50 万行销售记录。
传统方案:MySQL + 定时 ETL + 云服务器,开发周期 2 周,运维成本每月 500+ 元。
DuckDB 方案:直接对 Parquet 文件做 SQL 查询,零 ETL,毫秒级响应,单机就能扛住。
核心差异在于:DuckDB 是列式分析引擎,专门针对「读多写少」的分析场景优化。你的数据产品本质上是只读的分析服务,这恰恰是 DuckDB 的主场。
二、核心架构设计
Parquet 文件(每日分区)
↓
DuckDB 内存数据库(直接扫描 Parquet)
↓
FastAPI 路由
↓
付费用户 API 调用
这个架构的关键点:
- 数据用 Parquet 存储:列式格式,压缩比高,DuckDB 原生支持
- DuckDB 直接读取 Parquet:无需导入数据库,零拷贝
- FastAPI 做接口层:轻量、异步、自带 OpenAPI 文档
- 按店铺隔离数据:每个客户独立存储,天然数据隔离
三、第一步:准备测试数据
先用 Python 生成模拟数据,了解数据结构:
import duckdb
import pandas as pd
import numpy as np
from datetime import datetime, timedelta
def generate_daily_sales(days=30):
"""模拟生成30天销售数据"""
np.random.seed(42)
stores = [f"store_{i:03d}" for i in range(1, 11)]
products = [f"prod_{i:03d}" for i in range(1, 51)]
records = []
for day_offset in range(days):
date = datetime.now() - timedelta(days=day_offset)
n_records = np.random.randint(15000, 25000)
df = pd.DataFrame({
'date': np.random.choice(
pd.date_range(date, periods=1), n_records
),
'store_id': np.random.choice(stores, n_records),
'product_id': np.random.choice(products, n_records),
'quantity': np.random.randint(1, 20, n_records),
'unit_price': np.round(np.random.uniform(10, 500, n_records), 2),
'region': np.random.choice(['CN', 'US', 'EU'], n_records, p=[0.6, 0.3, 0.1])
})
df['revenue'] = df['quantity'] * df['unit_price']
records.append(df)
return pd.concat(records, ignore_index=True)
# 生成 150 万行模拟数据
df = generate_daily_sales(30)
print(f"生成 {len(df)} 行数据")
print(df.head())
生成的数据结构:
date:销售日期store_id:店铺 IDproduct_id:商品 IDquantity:销量unit_price:单价region:区域(CN/US/EU)revenue:营收(销量×单价)
四、第二步:构建查询引擎
DuckDB 的核心优势是直接对 Parquet 做 SQL 查询,无需加载到内存:
import duckdb
class SalesAnalyticsEngine:
def __init__(self, parquet_path="s3://bucket/sales/"):
self.path = parquet_path
self.con = duckdb.connect("memory:")
# 直接扫描 Parquet,零数据拷贝
self.con.execute(f"""
CREATE TABLE sales AS
SELECT * FROM read_parquet('{parquet_path}*.parquet')
""")
def daily_revenue(self, store_id: str, days: int = 30) -> list:
"""某店铺近N天每日营收"""
result = self.con.execute(f"""
SELECT date::DATE as dt,
SUM(revenue) as revenue,
SUM(quantity) as orders
FROM sales
WHERE store_id = '{store_id}'
AND date >= CURRENT_DATE - INTERVAL '{days} DAYS'
GROUP BY dt
ORDER BY dt DESC
LIMIT {days}
""").fetchall()
return [{"date": r[0].isoformat(), "revenue": r[1], "orders": r[2]} for r in result]
def top_products(self, store_id: str, n: int = 10) -> list:
"""TOP N 畅销商品"""
result = self.con.execute(f"""
SELECT product_id,
SUM(quantity) as total_qty,
SUM(revenue) as total_rev,
COUNT(DISTINCT date) as active_days
FROM sales
WHERE store_id = '{store_id}'
AND date >= CURRENT_DATE - INTERVAL '30 DAYS'
GROUP BY product_id
ORDER BY total_rev DESC
LIMIT {n}
""").fetchall()
return [{"product_id": r[0], "qty": r[1], "revenue": r[2], "days": r[3]} for r in result]
def regional_breakdown(self, store_id: str) -> dict:
"""区域销售分布"""
result = self.con.execute(f"""
SELECT region,
COUNT(*) as txn_count,
SUM(revenue) as revenue,
AVG(unit_price) as avg_price
FROM sales
WHERE store_id = '{store_id}'
GROUP BY region
""").fetchall()
return {r[0]: {"txns": r[1], "revenue": r[2], "avg_price": r[3]} for r in result}
关键理解:read_parquet() 不会把数据加载到内存,而是扫描 Parquet 文件的元数据,按需读取需要的列。这意味着即使 Parquet 文件有 10GB,你的内存占用也可能只有几十 MB。
五、第三步:封装为 REST API
from fastapi import FastAPI, HTTPException
from fastapi.middleware.cors import CORSMiddleware
import uvicorn
app = FastAPI(title="电商销售分析 API")
engine = SalesAnalyticsEngine()
@app.middleware("http")
async def auth_middleware(request, call_next):
"""简单 API Key 鉴权"""
api_key = request.headers.get("X-API-Key")
if api_key != "your-premium-key":
raise HTTPException(status_code=401)
return await call_next(request)
@app.get("/api/v1/{store_id}/daily")
async def get_daily(store_id: str, days: int = 30):
return {"store_id": store_id, "data": engine.daily_revenue(store_id, days)}
@app.get("/api/v1/{store_id}/top-products")
async def get_top_products(store_id: str, n: int = 10):
return {"store_id": store_id, "data": engine.top_products(store_id, n)}
@app.get("/api/v1/{store_id}/regional")
async def get_regional(store_id: str):
return {"store_id": store_id, "data": engine.regional_breakdown(store_id)}
if __name__ == "__main__":
uvicorn.run(app, host="0.0.0.0", port=8000)
启动后访问 http://localhost:8000/docs 可以看到自动生成的 API 文档。
六、性能对比:DuckDB vs 传统方案
| 维度 | DuckDB + Parquet | MySQL + ETL |
|---|---|---|
| 开发周期 | 1-2 天 | 1-2 周 |
| 数据存储成本 | 本地/S3,几乎为零 | 云服务器 + 数据库实例 |
| 查询响应(50万行) | < 50ms | 100-500ms |
| 运维复杂度 | 低(无数据库服务) | 高(需维护 MySQL) |
| 扩展路径 | DuckDB Cloud 无缝迁移 | 需重新架构 |
| 月运维成本 | ~100元(服务器) | ~500元+ |
七、进阶:加入缓存层
对于高频查询,可以加一个轻量 Redis 缓存:
import json, redis
r = redis.Redis(host='localhost', port=6379, db=0)
def get_with_cache(store_id: str, query_type: str, ttl: int = 300):
key = f"sales:{store_id}:{query_type}"
cached = r.get(key)
if cached:
return json.loads(cached)
if query_type == "daily":
data = engine.daily_revenue(store_id)
elif query_type == "top_products":
data = engine.top_products(store_id)
else:
data = engine.regional_breakdown(store_id)
r.setex(key, ttl, json.dumps(data))
return data
5 分钟 TTL 意味着:同一店铺的查询在 5 分钟内直接返回缓存,DuckDB 只在首次或缓存过期时才执行查询。
八、落地建议:从 MVP 到上线
- Day 1:搭好 API + 基础查询
- Day 2:加认证和限流
- Day 3:加监控(DuckDB 内置
duckdb.query_stats()可以看执行计划) - Day 4-7:找 3-5 个种子用户免费试用,收集反馈后定价上架
九、变现建议
这种数据产品单次开发成本约 2-3 天,按月订阅制定价 99-299 元/店铺,100 个店铺就是月入 1-3 万。
定价策略:
- 基础版(99元/月):每日营收趋势 + TOP 10 商品
- 进阶版(199元/月):基础功能 + 区域分析 + API 调用量 1000次/月
- 高级版(299元/月):全部功能 + 自定义查询 + API 调用量 5000次/月
获客渠道:
- 小红书/抖音发 DuckDB 教程,引流到私域
- 在 Twitter/X 上分享 DuckDB 实战案例
- 在独立开发者社区(Indie Hackers、V2EX)发帖
本文的完整版已发布在 duckdblab.org,包含更详细的步骤和更多案例。