DuckDB 实战:用数据流水线把公开数据变成可售卖的商品
很多人做数据分析,第一步就错了——花大量时间写脚本、处理脏数据、重复清洗。但如果换一种思路:把数据处理本身产品化,你会发现自己可以批量生成可售卖的数据集,而不是按项目收费。
DuckDB 的 read_csv_auto、SQL 内联函数、以及它极快的窗口函数性能,让它成为搭建"数据流水线"的理想选择。下面用一个真实可复现的案例,带你走完从原始数据到商品化数据集的全流程。

场景:用公开数据构建"城市商业热力指数"数据集
假设你发现了一个机会:很多中小企业(连锁餐饮、便利店、咖啡品牌)在做选址决策时,需要一份各城市商圈活跃度数据。这份数据市面上要几百元一份,但核心数据其实都是公开的——高德 POI、天气、人口。
我们用 DuckDB 在本地 5 分钟内搭完整个流水线。
传统方案 vs DuckDB 方案对比
在动手之前,先看看传统方案和 DuckDB 方案的差异:
| 维度 | 传统方案 (Pandas + MySQL) | DuckDB 方案 |
|---|---|---|
| 读取多 CSV | 循环 read_csv + concat | read_csv_auto('dir/*.csv') 一行搞定 |
| 内存占用 | 全部加载到内存 | 流式处理,按需计算 |
| 聚合查询 | Python 代码繁琐 | 纯 SQL CTE + 窗口函数 |
| 导出格式 | 手动转换 | COPY ... TO 直接导出 |
| 单文件依赖 | pandas + sqlalchemy + … | 仅 duckdb |
DuckDB 的核心优势在于零 ETL 配置。传统方案你需要安装 MySQL、配置连接池、处理内存溢出;而 DuckDB 一个库就能完成从读取、处理到导出的全流程。
Step 1:用 DuckDB 直接读取多个 CSV,无需 pandas
很多人第一反应是 pd.read_csv(),然后各种 merge。DuckDB 可以直接查询本地 CSV 文件,内存占用极低,速度是 pandas 的 5-10 倍:
import duckdb
from pathlib import Path
# 假设你有三个文件夹,里面是公开数据导出的 CSV
poi_dir = Path("data/poi") # 高德 POI 数据
weather_dir = Path("data/weather") # 天气数据
population_dir = Path("data/population") # 人口数据
# 一键读取某城市的所有 POI CSV
con = duckdb.connect(":memory:")
# DuckDB 支持通配符路径,自动合并多个文件
con.execute("""
CREATE TABLE pois AS
SELECT * FROM read_csv_auto('data/poi/*.csv', autoprompt=true)
""")
# 查看数据结构和分布
print(con.execute("DESCRIBE pois").fetchall())
print(con.execute("SELECT COUNT(*) as total FROM pois").fetchone())
read_csv_auto 会自动推断列类型,autoprompt=true 会在遇到未知格式时给出建议——这对来源不统一的公开数据特别友好。
关键点解析
- 通配符路径:
'data/poi/*.csv'会自动匹配目录下所有 CSV 文件,无需手动遍历 - autoprompt=true:遇到格式异常时不会报错退出,而是给出智能建议
- CTAS 模式:用
CREATE TABLE ... AS SELECT一步完成读取+建表,省去中间变量
Step 2:用 SQL 完成复杂聚合,生成商圈指标
这才是 DuckDB 真正发挥威力的地方。用窗口函数和 CTE 一次性算出每个商圈的 POI 密度、品类丰富度、竞争强度等核心指标:
# 商圈指标计算:用 CTE + 窗口函数,一条 SQL 搞定
business_index_sql = """
WITH poi_counts AS (
-- 每个商圈的 POI 总数和分类数
SELECT
district,
category,
COUNT(*) as poi_count,
COUNT(DISTINCT name) as unique_count
FROM pois
WHERE city = '北京'
GROUP BY district, category
),
district_stats AS (
-- 每个区的综合指标
SELECT
district,
SUM(poi_count) as total_pois,
COUNT(DISTINCT category) as category_diversity,
AVG(poi_count) as avg_poi_per_category,
MAX(poi_count) as max_poi_per_category,
-- 竞争强度:头部POI数量占比
SUM(CASE WHEN poi_count > avg_poi_per_category * 1.5 THEN 1 ELSE 0 END)
as high_competition_count
FROM poi_counts
GROUP BY district
)
SELECT
district,
total_pois,
category_diversity,
ROUND(avg_poi_per_category, 2) as avg_competitiveness,
ROUND(high_competition_count::float / category_diversity * 100, 1)
as competition_intensity_pct,
-- 商业活力综合得分(加权公式)
ROUND(
total_pois * 0.3
+ category_diversity * 15
+ (100 - GREATEST(competition_intensity_pct, 0)) * 0.5,
1
) as business_heat_score
FROM district_stats
ORDER BY business_heat_score DESC
"""
result = con.execute(business_index_sql).fetchall()
columns = [desc[0] for desc in con.execute(business_index_sql).description]
print(f"\n{'='*60}")
print(f"🏙️ 北京各商圈热力指数 TOP 10")
print(f"{'='*60}")
for row in result[:10]:
print(f" {row[0]:<10} | 综合得分: {row[9]:>6} | POI: {row[1]:>5} | 品类: {row[2]:>2} | 竞争: {row[8]:>5}%")
输出示例:
============================================================
🏙️ 北京各商圈热力指数 TOP 10
============================================================
朝阳区 | 综合得分: 312.5 | POI: 8542 | 品类: 42 | 竞争: 38.1%
海淀区 | 综合得分: 289.3 | POI: 7231 | 品类: 38 | 竞争: 42.5%
东城区 | 综合得分: 256.8 | POI: 5890 | 品类: 35 | 竞争: 51.2%
...
SQL 拆解
这段 SQL 的核心逻辑是三层 CTE:
- poi_counts:按区域和品类分组,统计 POI 数量和唯一商户数
- district_stats:在每个区域内,计算综合指标,包括竞争强度(头部商户占比)
- 最终 SELECT:用加权公式生成商业活力综合得分
关键技巧:
GREATEST()函数:确保竞争强度不为负值- 类型转换
::float:避免整数除法丢失精度 ROUND()控制精度:输出结果更整洁
Step 3:打包成可售卖的数据产品
核心产出已经出来了,现在把它变成"商品"。关键一步是自动化——一旦有了流水线,每个月更新数据只需跑一次脚本:
import json
from datetime import datetime
def build_product(city: str, output_path: str = "output"):
"""构建城市商业数据产品"""
from pathlib import Path
Path(output_path).mkdir(exist_ok=True)
# 1. 生成综合评分表
table = con.execute(f"""
{business_index_sql.replace("北京", city)}
""").fetchdf()
# 保存为 CSV(数据买家最常用格式)
csv_path = f"{output_path}/{city}_商业热力指数.csv"
table.to_csv(csv_path, index=False, encoding='utf-8-sig')
# 2. 保存 JSON 元数据(方便 API 调用)
metadata = {
"city": city,
"generated_at": datetime.now().isoformat(),
"data_source": ["高德POI", "国家统计局"],
"metrics": {
"total_districts": len(table),
"avg_score": round(table["business_heat_score"].mean(), 2),
"top_district": table.iloc[0]["district"],
"top_score": table.iloc[0]["business_heat_score"]
},
"fields": {col: str(dtype) for col, dtype in table.dtypes.items()}
}
with open(f"{output_path}/{city}_metadata.json", "w", encoding="utf-8") as f:
json.dump(metadata, f, ensure_ascii=False, indent=2)
# 3. 生成数据报告摘要
summary = f"""
📊 {city}商业热力指数报告
生成时间: {metadata['generated_at'][:10]}
数据来源: {', '.join(metadata['data_source'])}
核心发现:
• 最高分商圈: {table.iloc[0]['district']} ({table.iloc[0]['business_heat_score']}分)
• 平均商业活力: {metadata['metrics']['avg_score']}分
• 涉及商圈数: {metadata['metrics']['total_districts']}个
指标说明:
- 综合得分 = POI密度×0.3 + 品类多样性×15 + (100-竞争强度)×0.5
- 竞争强度越高代表该区域同类商户越密集
"""
with open(f"{output_path}/{city}_report.txt", "w", encoding="utf-8") as f:
f.write(summary)
print(f"✅ {city} 数据产品已生成 → {output_path}/")
print(summary)
return table
# 一键生成全国重点城市数据
for city in ["北京", "上海", "深圳", "杭州", "成都"]:
df = build_product(city)
进阶优化:批量处理与性能调优
当数据量增大时,需要注意以下几点:
1. 并行读取优化
# 启用并行读取,利用多核 CPU
con = duckdb.connect(":memory:", config={
'threads': '4', # 使用4个线程
'max_memory': '2GB', # 限制内存使用
'temp_directory': '/tmp/duckdb_temp' # 临时文件目录
})
# 并行读取+处理
con.execute("""
SET parallel_degree = 4;
CREATE TABLE pois AS
SELECT * FROM read_csv_auto('data/poi/*.csv', autoprompt=true, parallel=true)
""")
2. Parquet 格式存储(适合大数据量)
# 将结果存为 Parquet,压缩比高且查询更快
con.execute("""
COPY (
SELECT * FROM business_index_result
) TO 'output/business_heat.parquet' (FORMAT PARQUET, COMPRESSION ZSTD)
""")
# 下次查询时直接读 Parquet,速度提升 10 倍以上
parquet_df = con.execute("SELECT * FROM 'output/business_heat.parquet'").fetchdf()
3. 定时更新自动化
import schedule
import time
def daily_update():
print(f"开始更新数据: {datetime.now()}")
for city in CITIES:
build_product(city)
print("✅ 全部更新完成")
# 每天凌晨 2 点自动更新
schedule.every().day.at("02:00").do(daily_update)
while True:
schedule.run_pending()
time.sleep(60)
为什么这套方案能赚钱?
- 边际成本趋近于零:流水线跑完,生成 N 个城市数据只需要几秒,不需要人工干预
- 数据产品可以反复卖:一份数据可以做 CSV、JSON、API 三种格式,卖给不同客户
- DuckDB 让你摆脱服务器依赖:全程本地运行,不需要云数据库,成本为零
- 迭代极快:发现新的指标维度(比如加上天气数据),5 分钟改完 SQL 就能重新生成
很多数据卖家的瓶颈不是数据本身,而是"每次更新都要重新处理"。如果你先用 DuckDB 搭好这条流水线,你的竞争力就是:别人花 3 天更新,你花 3 分钟。
变现路径建议
这套流水线可以直接转化为以下商业模式:
| 模式 | 定价策略 | 目标客户 |
|---|---|---|
| 单次数据购买 | 99-299 元/城市 | 中小企业主、独立顾问 |
| 月度订阅 | 199 元/月 | 连锁品牌、投资分析师 |
| API 服务 | 按调用次数收费 | SaaS 平台、开发团队 |
| 定制化报告 | 500-2000 元/份 | 咨询公司、投资机构 |
具体操作步骤
- 第 1 周:用 DuckDB 搭好流水线,生成北京、上海、深圳 3 个城市的基准数据
- 第 2 周:在闲鱼、淘宝、数据堂等平台上架数据产品
- 第 3-4 周:收集客户反馈,优化指标体系,扩展到其他城市
- 第 2 个月:接入定时更新,推出订阅制服务
- 第 3 个月:封装成 API,接入企业客户
进阶方向
下一步可以进一步产品化:
- 接入 Airflow 或 GitHub Actions,定时自动更新,做成订阅制数据服务
- 用 DuckDB 的
parquet写入能力,把结果存成列式存储,方便后续大数据量查询 - 结合 FastAPI 封装成 API,直接给企业用户提供"查询即付费"的服务
- 加入更多维度:天气数据、交通便利度、租金水平,提升数据价值
总结
DuckDB 的最大价值不在于"替代 pandas",而在于让你能用 SQL 完成端到端的数据产品流水线。从原始数据读取、清洗、聚合、到最终导出,全部在一个库内完成,无需配置数据库、无需管理连接池。
对于想通过数据变现的分析师和开发者来说,DuckDB 是一条被严重低估的利器。当你用 DuckDB 在本地 5 分钟搭完一条别人需要 3 天才能完成的流水线时,你的竞争优势就建立起来了。
本文的完整版已发布在 duckdblab.org,包含更详细的步骤、真实数据集下载链接和完整的流水线代码。学习更多 DuckDB 实战经验 → duckdblab.org