Featured image of post DuckDB 实战:用数据流水线把公开数据变成可售卖的商品

DuckDB 实战:用数据流水线把公开数据变成可售卖的商品

用 DuckDB 搭建数据处理流水线,将公开数据加工成可售卖的城市商业热力指数数据集。从原始 CSV 到商品化数据产品全流程实战,附完整代码。

DuckDB 实战:用数据流水线把公开数据变成可售卖的商品

很多人做数据分析,第一步就错了——花大量时间写脚本、处理脏数据、重复清洗。但如果换一种思路:把数据处理本身产品化,你会发现自己可以批量生成可售卖的数据集,而不是按项目收费。

DuckDB 的 read_csv_auto、SQL 内联函数、以及它极快的窗口函数性能,让它成为搭建"数据流水线"的理想选择。下面用一个真实可复现的案例,带你走完从原始数据到商品化数据集的全流程。

数据流水线架构图


场景:用公开数据构建"城市商业热力指数"数据集

假设你发现了一个机会:很多中小企业(连锁餐饮、便利店、咖啡品牌)在做选址决策时,需要一份各城市商圈活跃度数据。这份数据市面上要几百元一份,但核心数据其实都是公开的——高德 POI、天气、人口。

我们用 DuckDB 在本地 5 分钟内搭完整个流水线。


传统方案 vs DuckDB 方案对比

在动手之前,先看看传统方案和 DuckDB 方案的差异:

维度传统方案 (Pandas + MySQL)DuckDB 方案
读取多 CSV循环 read_csv + concatread_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:

  1. poi_counts:按区域和品类分组,统计 POI 数量和唯一商户数
  2. district_stats:在每个区域内,计算综合指标,包括竞争强度(头部商户占比)
  3. 最终 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)

为什么这套方案能赚钱?

  1. 边际成本趋近于零:流水线跑完,生成 N 个城市数据只需要几秒,不需要人工干预
  2. 数据产品可以反复卖:一份数据可以做 CSV、JSON、API 三种格式,卖给不同客户
  3. DuckDB 让你摆脱服务器依赖:全程本地运行,不需要云数据库,成本为零
  4. 迭代极快:发现新的指标维度(比如加上天气数据),5 分钟改完 SQL 就能重新生成

很多数据卖家的瓶颈不是数据本身,而是"每次更新都要重新处理"。如果你先用 DuckDB 搭好这条流水线,你的竞争力就是:别人花 3 天更新,你花 3 分钟。


变现路径建议

这套流水线可以直接转化为以下商业模式:

模式定价策略目标客户
单次数据购买99-299 元/城市中小企业主、独立顾问
月度订阅199 元/月连锁品牌、投资分析师
API 服务按调用次数收费SaaS 平台、开发团队
定制化报告500-2000 元/份咨询公司、投资机构

具体操作步骤

  1. 第 1 周:用 DuckDB 搭好流水线,生成北京、上海、深圳 3 个城市的基准数据
  2. 第 2 周:在闲鱼、淘宝、数据堂等平台上架数据产品
  3. 第 3-4 周:收集客户反馈,优化指标体系,扩展到其他城市
  4. 第 2 个月:接入定时更新,推出订阅制服务
  5. 第 3 个月:封装成 API,接入企业客户

进阶方向

下一步可以进一步产品化:

  • 接入 Airflow 或 GitHub Actions,定时自动更新,做成订阅制数据服务
  • 用 DuckDB 的 parquet 写入能力,把结果存成列式存储,方便后续大数据量查询
  • 结合 FastAPI 封装成 API,直接给企业用户提供"查询即付费"的服务
  • 加入更多维度:天气数据、交通便利度、租金水平,提升数据价值

总结

DuckDB 的最大价值不在于"替代 pandas",而在于让你能用 SQL 完成端到端的数据产品流水线。从原始数据读取、清洗、聚合、到最终导出,全部在一个库内完成,无需配置数据库、无需管理连接池。

对于想通过数据变现的分析师和开发者来说,DuckDB 是一条被严重低估的利器。当你用 DuckDB 在本地 5 分钟搭完一条别人需要 3 天才能完成的流水线时,你的竞争优势就建立起来了。

本文的完整版已发布在 duckdblab.org,包含更详细的步骤、真实数据集下载链接和完整的流水线代码。学习更多 DuckDB 实战经验 → duckdblab.org

📺 Watch video tutorials → Olap Studio YouTube

Subscribe for more DuckDB & AI automation tutorials

使用 Hugo 构建
主题 StackJimmy 设计