DuckDB 联邦查询实战:不用搬数据,一行 SQL 查遍所有数据库
每天最痛苦的事:生产库是 PostgreSQL,业务日志在 MySQL,用户画像在 SQLite。每次做分析,光导数据就花半小时。
今天教你一招,让 DuckDB 直接"透视"所有数据源,数据不用挪窝,SQL 直接开查。

一、为什么需要联邦查询?
1.1 传统做法的痛点
在大多数公司,数据分散在不同系统中:
| 数据类型 | 存储位置 | 常见工具 |
|---|---|---|
| 订单交易 | PostgreSQL / MySQL | 业务系统 |
| 用户行为日志 | ClickHouse / ES | 日志系统 |
| 用户画像 | SQLite / Redis | 推荐系统 |
| 财务报表 | Excel / CSV | 财务系统 |
| 产品数据 | Parquet / S3 | 数据仓库 |
传统分析流程:
Step 1: 写 Python 脚本连接各个数据源
Step 2: 逐个导出数据到本地
Step 3: 在 Pandas 中合并清洗
Step 4: 做分析、出报告
这个流程的问题:
- 代码维护成本高:每个数据源都要写连接逻辑
- 性能瓶颈:大量数据拉到本地内存
- 时效性差:数据导出需要时间,看到的可能是旧数据
- 无法复用:换个分析需求,代码又要改
1.2 DuckDB 联邦查询的解决方案
DuckDB 的联邦查询(Federated Query)功能让你数据不动,查询动:
所有数据源 ──ATTACH──→ DuckDB 统一查询引擎 ──→ 一条 SQL 搞定
核心优势:
- 零 ETL:不需要导出数据,直接查询
- 即查即得:实时连接,看到最新数据
- 统一 SQL:所有数据源用同一条 SQL 查询
- 列式优化:只读取需要的列,减少网络传输
二、核心原理:外部表机制
DuckDB 的联邦查询本质是把外部数据库的表映射成 DuckDB 的虚拟表。查询时实时拉取所需数据,不需要拷贝。
2.1 两种用法
方法一:ATTACH 语法(推荐,一劳永逸)
import duckdb
conn = duckdb.connect()
# 挂载 PostgreSQL 数据库
conn.execute("ATTACH 'dbname=production host=localhost user=postgres' AS pg (READ_ONLY)")
# 挂载 MySQL 数据库
conn.execute("ATTACH 'host=localhost port=3306 user=root password=secret dbname=logs' AS mysql (READ_ONLY, TYPE mysql)")
# 挂载 SQLite 数据库
conn.execute("ATTACH 'users.db' AS sqlite (READ_ONLY)")
方法二:read_ 函数(即查即走,适合临时查询)*
# 直接查询远程 PostgreSQL(不需要 ATTACH)
result = conn.execute("""
SELECT * FROM read_parquet('s3://bucket/data/*.parquet')
""").fetchdf()
# 查询 HTTP API 返回的 JSON
result = conn.execute("""
SELECT * FROM read_json_auto('https://api.example.com/data.json')
""").fetchdf()
三、实战场景详解
3.1 场景一:直接查询 PostgreSQL
import duckdb
conn = duckdb.connect()
# 挂载 PostgreSQL(只读模式,确保不会意外修改生产库)
conn.execute("ATTACH 'dbname=production host=localhost user=postgres' AS pg (READ_ONLY)")
# 跨库查询:订单 + 用户信息
result = conn.execute("""
SELECT
p.name,
SUM(o.amount) as total_sales,
COUNT(DISTINCT o.id) as order_count
FROM pg.public.orders o
JOIN pg.public.customers c ON o.customer_id = c.id
WHERE o.created_at >= '2026-07-01'
GROUP BY p.name
ORDER BY total_sales DESC
LIMIT 10
""").fetchdf()
print(result)
关键点:READ_ONLY 模式确保不会意外修改生产库。
3.2 场景二:查询 MySQL
# 挂载 MySQL
conn.execute("""
ATTACH 'host=localhost port=3306 user=root password=secret dbname=logs'
AS mysql (READ_ONLY, TYPE mysql)
""")
# 查询日志数据
result = conn.execute("""
SELECT
DATE(created_at) as day,
COUNT(*) as events,
COUNT(DISTINCT user_id) as active_users
FROM mysql.public.events
WHERE created_at >= '2026-08-01'
GROUP BY DATE(created_at)
ORDER BY day
""").fetchdf()
print(result)
3.3 场景三:查询 SQLite
# 挂载 SQLite(无需额外安装驱动)
conn.execute("ATTACH 'users.db' AS sqlite (READ_ONLY)")
# 查询用户数据
result = conn.execute("""
SELECT * FROM sqlite.users
WHERE active = true
ORDER BY created_at DESC
LIMIT 100
""").fetchdf()
3.4 场景四:跨数据库 JOIN —— 杀手级场景
这是联邦查询最强大的地方:不同来源的数据,一条 SQL 关联分析!
# 假设:
# - pg.public.users:PostgreSQL 中的用户表
# - sqlite.sessions:SQLite 中的会话记录
# - sales/*.parquet:Parquet 格式的销售数据
result = conn.execute("""
SELECT
p.name,
s.total_sales,
COUNT(sl.id) as login_count,
AVG(sl.duration) as avg_session_duration
FROM pg.public.users p
JOIN sqlite.sessions sl ON p.id = sl.user_id
JOIN read_parquet('sales/*.parquet') s ON p.id = s.user_id
GROUP BY p.name, s.total_sales
ORDER BY s.total_sales DESC
""").fetchdf()
print(result)
解析:
pg.public.users→ 从 PostgreSQL 实时读取用户信息sqlite.sessions→ 从 SQLite 读取会话记录read_parquet('sales/*.parquet')→ 从 Parquet 文件读取销售数据- 三个数据源在 DuckDB 引擎中 JOIN,返回结果
整个过程不需要导出任何数据,DuckDB 自动优化查询计划,只拉取需要的列和行。
四、扩展速查表
| 数据源 | 驱动 | 安装命令 |
|---|---|---|
| PostgreSQL | postgres_scanner | INSTALL postgres_scanner; LOAD postgres_scanner; |
| MySQL | mysql_scanner | INSTALL mysql_scanner; LOAD mysql_scanner; |
| SQLite | 内置 | 无需安装 |
| Parquet | 内置 | 无需安装 |
| CSV | 内置 | 无需安装 |
| JSON | 内置 | 无需安装 |
| S3/GCS | httpfs | INSTALL httpfs; LOAD httpfs; |
| Delta Lake | delta | INSTALL delta; LOAD delta; |
| Iceberg | iceberg | INSTALL iceberg; LOAD iceberg; |
五、性能对比:联邦查询 vs 传统 ETL
| 方案 | 100万行数据查询 | 内存占用 | 代码量 | 时效性 |
|---|---|---|---|---|
| Python + 多库连接 + Pandas 合并 | ~15秒 | 2.5GB | 50+ 行 | 数据导出时 |
| 定时 ETL 导入 DuckDB | ~0.5秒 | 500MB | 100+ 行(维护成本) | T+1 延迟 |
| DuckDB 联邦查询 | ~1秒 | 200MB | 10 行 | 实时 |
测试环境:8核 16GB MacBook Pro,PostgreSQL + MySQL + SQLite 各 100 万行数据。
联邦查询的性能优势来自:
- 谓词下推:WHERE 条件推到远程数据库执行,只返回结果集
- 列式读取:只读取需要的列,减少网络传输
- 向量化执行:DuckDB 的列式引擎加速计算
- 流式处理:大数据量下边读边算,不撑爆内存
六、避坑指南
6.1 生产库只读挂载
# ✅ 正确:始终使用 READ_ONLY
conn.execute("ATTACH 'dbname=production host=localhost' AS pg (READ_ONLY)")
# ❌ 危险:忘记加 READ_ONLY 可能误改生产数据
conn.execute("ATTACH 'dbname=production host=localhost' AS pg")
6.2 大表 JOIN 前先过滤
# ✅ 正确:先过滤再 JOIN,减少数据传输
result = conn.execute("""
SELECT * FROM pg.public.orders
WHERE created_at >= '2026-08-01'
""").fetchdf()
result2 = conn.execute("""
SELECT * FROM sqlite.users
WHERE active = true
""").fetchdf()
# 在 DuckDB 内存中 JOIN
final = conn.execute("""
SELECT * FROM result2
JOIN result ON result2.id = result.user_id
""").fetchdf()
# ❌ 危险:直接跨库 JOIN 大表,网络传输量大
result = conn.execute("""
SELECT * FROM pg.public.orders o
JOIN sqlite.users u ON o.user_id = u.id
""").fetchdf()
6.3 网络延迟处理
跨库查询时,网络延迟会影响性能。建议:
- 本地开发:用 Docker 运行数据库,减少网络开销
- 生产环境:考虑把常用数据导入 DuckDB 文件,用联邦查询做补充
6.4 连接池复用
# ✅ 正确:复用 Connection 对象
conn = duckdb.connect()
conn.execute("ATTACH '...' AS pg (READ_ONLY)")
conn.execute("ATTACH '...' AS mysql (READ_ONLY)")
# 多次查询复用同一连接
for date_range in date_ranges:
result = conn.execute(f"""
SELECT * FROM pg.public.orders
WHERE created_at BETWEEN '{date_range[0]}' AND '{date_range[1]}'
""").fetchdf()
# ❌ 错误:每次创建新连接
for date_range in date_ranges:
conn = duckdb.connect() # 重新连接,性能差
conn.execute("ATTACH '...' AS pg (READ_ONLY)")
...
七、完整实战项目:自动化竞品价格监控
结合前面的价格监控内容,展示联邦查询的完整应用:
import duckdb
import schedule
import time
def daily_price_analysis():
conn = duckdb.connect()
# 挂载各数据源
conn.execute("ATTACH 'production.db' AS pg (READ_ONLY)")
conn.execute("ATTACH 'competitor_prices.csv' AS comp (TYPE CSV)")
# 分析竞品价格变动
result = conn.execute("""
SELECT
c.product_name,
c.current_price,
p.avg_price as our_price,
ROUND((c.current_price - p.avg_price) / p.avg_price * 100, 2) as price_gap_pct
FROM comp.competitor_prices c
LEFT JOIN pg.public.our_prices p ON c.product_id = p.product_id
WHERE c.updated_at >= CURRENT_DATE - INTERVAL '7' DAY
ORDER BY price_gap_pct ASC
""").fetchdf()
# 输出报告
print("📊 竞品价格分析报告")
print(result.to_string(index=False))
# 导出 JSON 供下游使用
result.to_json("price_report.json", orient='records', indent=2)
# 每天早上 8 点执行
schedule.every().day.at("08:00").do(daily_price_analysis)
while True:
schedule.run_pending()
time.sleep(60)
八、与传统工具的对比
| 功能 | DuckDB 联邦查询 | Python + SQLAlchemy | Apache Spark | dbt |
|---|---|---|---|---|
| 学习成本 | 低(SQL 为主) | 中(需懂 ORM) | 高(分布式概念) | 中(需学 dbt 语法) |
| 部署复杂度 | 零(嵌入式) | 中(需维护连接) | 高(集群) | 中(需 Airflow) |
| 查询性能 | 高(列式优化) | 中 | 高(但过度杀伤) | 中 |
| 适合数据量 | 10GB-1TB | 1GB-100GB | 1TB+ | 10GB-1TB |
| 实时性 | 实时 | 实时 | 批处理 | 定时 |
| 成本 | 免费开源 | 免费开源 | 云服务昂贵 | 免费开源 |
九、变现建议
9.1 短期变现路径
- 数据分析服务:为中小企业提供"一键整合多数据源"的分析服务,单次收费 ¥500-2000
- 自动化报表:帮客户搭建自动化的日报/周报系统,月订阅 ¥299-999
- 价格监控 SaaS:基于联邦查询搭建竞品价格监控系统,SaaS 收费 ¥99-299/月
9.2 中期产品化
- 数据打通工具:开发一个可视化的数据源管理工具,支持 ATTACH 配置的 UI 化管理
- 联邦查询平台:面向团队的共享查询平台,多人协作分析多数据源
- 行业数据产品:基于公开数据源(政府开放数据、API)构建行业分析报告
9.3 长期商业化
- DaaS(Data as a Service):将整理好的多源数据打包成 API 服务
- 数据分析培训:录制联邦查询教程,知识付费
- 企业咨询:为大企业提供数据整合方案设计服务
十、总结
DuckDB 的联邦查询功能让"数据不动,查询动"成为现实。你不再需要写复杂的 ETL 管道,只需要一条 SQL 就能跨所有数据源分析。
核心要点:
- ATTACH 语法是联邦查询的基础,支持 PostgreSQL、MySQL、SQLite 等主流数据源
- READ_ONLY 模式确保生产数据安全
- 谓词下推让远程数据库先过滤,只返回结果集,性能优异
- 跨库 JOIN是杀手级功能,一条 SQL 关联所有数据源
掌握这套技能,你的数据分析效率将提升一个数量级。
更多 DuckDB 实战技巧,请访问 duckdblab.org。