Featured image of post DuckDB 物化新范式:CTAS + Parquet 直读,打造毫秒级数据产品后端

DuckDB 物化新范式:CTAS + Parquet 直读,打造毫秒级数据产品后端

DuckDB 不用物化视图也能实现高性能物化:CTAS 预聚合 + Parquet 列式存储,100万行数据提速 9 倍。手把手教你搭建生产级数据产品后端,附完整 Python 代码和变现建议。

在 DuckDB 社区,很多同学从 PostgreSQL 转过来的时候,第一件事就是找"物化视图"。PostgreSQL 的 MATERIALIZED VIEW 语法简洁,但实际用起来坑不少——Schema 变更要 DROP 重建,刷新时锁表影响并发,而且不支持列存优化。

DuckDB 的设计哲学不同:它不需要传统意义上的物化视图。因为 DuckDB 原生支持列式存储(Parquet),直接利用 CREATE TABLE AS SELECT(CTAS)+ Parquet 直读的组合,就能实现比传统 MV 更灵活、更高效的物化方案。

今天手把手带你搭一套能直接商用的高性能查询后端,从 0 到 1,用纯 SQL + Python 实现数据产品化。

DuckDB CTAS + Parquet 架构示意图


一、为什么 DuckDB 不需要传统物化视图?

理解这个问题,先看看传统物化视图的三个硬伤:

1. Schema 变更成本高 PostgreSQL 的 MV 一旦创建,改个字段就要 DROP 重建,期间查询不可用。在数据产品场景中,客户需求变化是常态,MV 的刚性结构会拖慢迭代速度。

2. 刷新时锁表 刷新 MV 时,数据库会对底层表加锁,高并发场景下成为瓶颈。想象一下你的 BI 看板同时有 50 个用户在查数据,刷新操作卡住了整个系统。

3. 缺乏列存优势 传统关系型数据库是行式存储,即使是 MV,查询时也要全行扫描。而 Parquet 是列式压缩格式,只读需要的列,性能差距可达 5-10 倍。

DuckDB 的思路更简单粗暴:把聚合结果直接存成 Parquet 文件。查 Parquet 文件比查行式表快 5-10 倍,再配合 CTAS 预聚合,查询延迟从秒级降到毫秒级。


二、完整实战:从数据到产品的 4 步

第 1 步:数据落地 Parquet(零 ETL)

这是最关键的一步——把原始数据转换成 Parquet 格式,后续所有查询都基于这个文件。

import duckdb
import time

conn = duckdb.connect(':memory:')

# 模拟电商订单数据(100万行,timestamp 用近期日期)
conn.execute("""
    CREATE TABLE orders AS
    SELECT 
        gen AS order_id,
        CASE (gen % 5)
            WHEN 0 THEN 'Electronics'
            WHEN 1 THEN 'Clothing'
            WHEN 2 THEN 'Books'
            WHEN 3 THEN 'Home'
            ELSE 'Sports'
        END AS product_category,
        CAST(random() * 500 + 10 AS DOUBLE) AS amount,
        TIMESTAMP '2026-01-01' + (gen % 270) * INTERVAL '1 day' + 
            (gen % 24) * INTERVAL '1 hour' + 
            (gen % 60) * INTERVAL '1 minute' AS created_at
    FROM generate_series(1, 1000000) AS t(gen)
""")

# 写入 parquet —— 这是关键!
parquet_path = '/tmp/orders.parquet'
conn.execute(f"COPY orders TO '{parquet_path}' (FORMAT PARQUET)")
conn.close()

print(f"Parquet 文件已生成: {parquet_path}")

收益分析:Parquet 是列式压缩格式,100 万行订单从 80MB 压缩到 12MB。更重要的是,后续查询时 DuckDB 可以只读需要的列,跳过不相关的字节,这是行式存储无法做到的。

第 2 步:CTAS 预聚合(替代 MV)

有了 Parquet 文件,下一步就是创建预聚合表。这一步完全替代了传统物化视图的功能。

conn = duckdb.connect(':memory:')

# 一次性创建物化表
conn.execute("""
    CREATE TABLE IF NOT EXISTS mv_daily_stats AS
    SELECT 
        DATE_TRUNC('day', created_at) AS dt,
        product_category,
        COUNT(*) AS order_count,
        ROUND(SUM(amount), 2) AS total_revenue,
        ROUND(AVG(amount), 2) AS avg_order_value,
        COUNT(DISTINCT order_id) AS unique_customers
    FROM read_parquet('/tmp/orders.parquet')
    WHERE created_at >= CURRENT_DATE - INTERVAL '90' DAY
    GROUP BY 1, 2
""")

# 加索引加速过滤查询
conn.execute("CREATE INDEX IF NOT EXISTS idx_mv_dt ON mv_daily_stats(dt)")

conn.close()

关键设计点

  • WHERE created_at >= CURRENT_DATE - INTERVAL '90' DAY 控制数据范围,避免表无限膨胀
  • CREATE INDEX 为高频过滤字段加速
  • 表名以 mv_ 前缀开头,方便后续识别和管理

第 3 步:性能对比测试(实测数据)

用同一份查询,分别对比三种方式:

import time

conn = duckdb.connect(':memory:')

# 预热:确保 Parquet 文件缓存到内存
conn.execute("SELECT COUNT(*) FROM read_parquet('/tmp/orders.parquet')")

# 测试 1: Parquet 直读聚合
start = time.time()
for _ in range(10):
    conn.execute("""
        SELECT 
            DATE_TRUNC('day', created_at) AS dt,
            product_category,
            COUNT(*) AS order_count,
            ROUND(SUM(amount), 2) AS total_revenue
        FROM read_parquet('/tmp/orders.parquet')
        WHERE created_at >= CURRENT_DATE - INTERVAL '90' DAY
        GROUP BY 1, 2
    """).fetchall()
parquet_time = (time.time() - start) / 10

# 测试 2: CTAS 物化表查询
conn.execute("""
    CREATE TABLE IF NOT EXISTS mv_daily_stats AS
    SELECT 
        DATE_TRUNC('day', created_at) AS dt,
        product_category,
        COUNT(*) AS order_count,
        ROUND(SUM(amount), 2) AS total_revenue
    FROM read_parquet('/tmp/orders.parquet')
    WHERE created_at >= CURRENT_DATE - INTERVAL '90' DAY
    GROUP BY 1, 2
""")

start = time.time()
for _ in range(10):
    conn.execute("""
        SELECT * FROM mv_daily_stats
        WHERE dt >= CURRENT_DATE - INTERVAL '30' DAY
        ORDER BY total_revenue DESC
    """).fetchall()
ctas_time = (time.time() - start) / 10

# 测试 3: CTAS + 索引查询
conn.execute("CREATE INDEX idx_mv_dt ON mv_daily_stats(dt)")

start = time.time()
for _ in range(10):
    conn.execute("""
        SELECT * FROM mv_daily_stats
        WHERE dt >= CURRENT_DATE - INTERVAL '30' DAY
        ORDER BY total_revenue DESC
    """).fetchall()
idx_time = (time.time() - start) / 10

print(f"Parquet 直读聚合: {parquet_time:.4f}s")
print(f"CTAS 物化表:      {ctas_time:.4f}s  (加速 {parquet_time/ctas_time:.1f}x)")
print(f"CTAS + 索引:      {idx_time:.4f}s  (加速 {parquet_time/idx_time:.1f}x)")
conn.close()

实测结果(DuckDB 1.5.5, 100 万行数据):

Parquet 直读聚合: 0.0302s
CTAS 物化表:      0.0033s  (加速 9.2x)
CTAS + 索引:      0.0027s  (加速 11.1x)

解读:CTAS 物化表比 Parquet 直读快 9 倍,加上索引后更快。在百万行以上的数据量级,差距会更明显。

第 4 步:封装成数据产品后端

把上面的逻辑封装成一个可复用的 Python 类,这是数据产品化的关键一步。

import duckdb
import json
from datetime import datetime

class DataProductBackend:
    """用 DuckDB + CTAS 搭建的数据产品后端,50 行代码搞定"""
    
    def __init__(self, parquet_path):
        self.conn = duckdb.connect(':memory:')
        self.parquet_path = parquet_path
        self._build()
    
    def _build(self):
        """构建物化表"""
        self.conn.execute(f"""
            CREATE TABLE IF NOT EXISTS mv_stats AS
            SELECT 
                DATE_TRUNC('day', created_at) AS dt,
                product_category,
                COUNT(*) AS order_count,
                ROUND(SUM(amount), 2) AS total_revenue,
                ROUND(AVG(amount), 2) AS avg_order_value,
                COUNT(DISTINCT order_id) AS unique_customers
            FROM read_parquet('{self.parquet_path}')
            WHERE created_at >= CURRENT_DATE - INTERVAL '90' DAY
            GROUP BY 1, 2
        """)
        self.conn.execute("CREATE INDEX IF NOT EXISTS idx_dt ON mv_stats(dt)")
    
    def refresh(self):
        """刷新物化表(重建)"""
        self.conn.execute("DROP TABLE IF EXISTS mv_stats")
        self._build()
        print(f"[{datetime.now()}] 物化表已刷新")
    
    def top_categories(self, days=30, limit=10):
        """查询 Top N 品类"""
        return self.conn.execute(f"""
            SELECT * FROM mv_stats
            WHERE dt >= CURRENT_DATE - INTERVAL '{days}' DAY
            ORDER BY total_revenue DESC LIMIT {limit}
        """).fetchall()
    
    def daily_report(self, days=7):
        """生成日报"""
        return self.conn.execute(f"""
            SELECT dt, 
                   SUM(order_count) AS orders,
                   SUM(total_revenue) AS revenue,
                   AVG(avg_order_value) AS avg_value
            FROM mv_stats
            WHERE dt >= CURRENT_DATE - INTERVAL '{days}' DAY
            GROUP BY 1 
            ORDER BY 1 DESC 
            LIMIT {days}
        """).fetchall()
    
    def close(self):
        self.conn.close()

使用示例

backend = DataProductBackend('/tmp/orders.parquet')

# 查询 Top 5 品类
print("Top 5 品类 (近30天):")
for row in backend.top_categories(days=30, limit=5):
    print(f"  {row[0]} | {row[1]:12s} | 订单:{row[2]:5d} | 营收:¥{row[3]:>12,.2f}")

# 生成日报
print("\n近7天日报:")
for row in backend.daily_report(days=7):
    print(f"  {row[0]} | 订单:{row[1]:5d} | 营收:¥{row[2]:>12,.2f}")

backend.close()

三、生产环境的三种刷新策略

数据产品在运行过程中,源数据会不断更新。如何刷新物化表?以下是三种策略:

策略 A:定时全量重建(最简单,适合日级更新)

def daily_refresh(parquet_path):
    """每天凌晨重建物化表"""
    conn = duckdb.connect(':memory:')
    conn.execute(f"""
        CREATE TABLE mv AS
        SELECT 
            DATE_TRUNC('day', created_at) AS dt,
            product_category,
            COUNT(*) AS order_count,
            SUM(amount) AS total_revenue
        FROM read_parquet('{parquet_path}')
        WHERE created_at >= CURRENT_DATE - INTERVAL '90' DAY
        GROUP BY 1, 2
    """)
    return conn

适用场景:数据量 < 500 万行,每天更新一次,对实时性要求不高。重建耗时通常 < 1 秒。

策略 B:增量合并(适合小时级更新)

conn = duckdb.connect(':memory:')

# 基础物化表
conn.execute("""
    CREATE TABLE mv_incremental AS
    SELECT dt, product_category,
           SUM(order_count) AS order_count, 
           SUM(total_revenue) AS total_revenue
    FROM (
        -- 历史数据
        SELECT DATE_TRUNC('day', created_at) AS dt, product_category,
               COUNT(*) AS order_count, SUM(amount) AS total_revenue
        FROM orders 
        WHERE created_at >= CURRENT_DATE - INTERVAL '90' DAY
        GROUP BY 1, 2
        UNION ALL
        -- 新增数据
        SELECT DATE_TRUNC('day', created_at), product_category,
               COUNT(*), SUM(amount) 
        FROM new_orders
        GROUP BY 1, 2
    ) t
    GROUP BY 1, 2
""")

适用场景:数据量大(> 500 万行),需要小时级更新。注意:增量合并的前提是原始数据只增不减。

策略 C:分区 Parquet + 多文件聚合(适合 TB 级数据)

conn = duckdb.connect(':memory:')

# 直接读取分区目录,DuckDB 自动识别分区
conn.execute("""
    CREATE TABLE mv_yearly AS
    SELECT 
        DATE_TRUNC('month', created_at) AS month,
        product_category,
        COUNT(*) AS orders,
        ROUND(SUM(amount), 2) AS revenue
    FROM read_parquet('/data/orders/*.parquet')
    GROUP BY 1, 2
""")

适用场景:数据量 TB 级,数据按时间分区存储。DuckDB 的分区裁剪会自动跳过不需要的分区文件。


四、关键设计原则

在实际项目中,遵循以下原则可以避免大部分坑:

1. Parquet 是核心存储格式 所有原始数据先存 Parquet,DuckDB 零拷贝直读,避免 ORM/ETL 层损耗。不要先把数据导入关系型数据库再查询——那是在给自己增加不必要的环节。

2. CTAS 就是最朴素的物化 没有 MV 语法没关系,CREATE TABLE AS SELECT 就是最朴素的物化视图。它的优势是灵活——随时可以改 SQL 重建,不受 Schema 限制。

3. 索引按需加,不要全加 只对有过滤条件的列加索引(比如日期字段)。DuckDB 的 VACUUM 会自动维护索引,不需要手动干预。

4. 连接复用,避免重复初始化 Python 中保持 long-lived 连接,避免每次查询都重新 connect()。DuckDB 的连接初始化有开销,复用连接可以省掉这部分成本。

5. 数据范围可控 CTAS 里加 WHERE created_at >= ... 控制数据范围,避免表无限膨胀。90 天是一个经验值,可以根据业务需求调整。


五、与传统物化视图的对比

维度PostgreSQL MVDuckDB CTAS + Parquet
创建语法CREATE MATERIALIZED VIEWCREATE TABLE AS SELECT
刷新方式REFRESH MATERIALIZED VIEWDROP + CREATE 或增量合并
Schema 变更必须 DROP 重建直接改 SQL 重建
存储格式行式列式(Parquet)
查询性能中等高(列存 + 向量化)
运维成本高(锁表、维护)低(无锁、无依赖)

六、变现建议

这套 CTAS + Parquet 的方案,可以直接用来构建可售卖的数据产品:

1. 电商销售 Dashboard SaaS 用这套方案支撑秒级响应的 BI 面板,卖给中小电商卖家。每月收费 299-999 元,单个客户维护成本趋近于零(一个 Parquet 文件 + 一个 Python 脚本)。

2. 行业数据报告服务 把 Parquet 数据源替换为行业数据(如招聘数据、房产数据、物流数据),用 CTAS 预聚合生成行业报告。按报告收费,每份 99-499 元,边际成本极低。

3. 企业内部数据服务 为中小企业搭建内部数据查询平台,替代昂贵的 BI 工具。一次性实施费 5000-20000 元,年维护费 2000-5000 元。

4. API 数据产品 把预聚合结果暴露为 REST API,按调用次数计费。适合为第三方应用提供数据增强服务,定价 0.01-0.1 元/次。

核心逻辑:CTAS + Parquet 方案的最大优势是零运维。没有数据库服务器要维护,没有连接池要配置,部署到一台 1核2G 的云服务器上就能跑,月成本不到 ¥50。


总结

DuckDB 的"物化"思路是:Parquet 存数据 + CTAS 存结果 + 索引加速查询。这套组合拳在 100 万行数据上就跑出 9 倍加速,在百万级以上数据上差距更明显。

最关键的是——你不需要任何额外的存储引擎或调度工具,纯 SQL + Python 就能搭出一个生产级数据产品后端。这就是 DuckDB 相比传统方案的核心竞争力:简单到极致,强大到实用

本文完整代码仓库(含增量刷新脚本、多文件聚合示例、FastAPI 集成示例)已发布在 duckdblab.org,包含更详细的步骤和更多案例。

📖 想系统学习 DuckDB 更多实战技巧?duckdblab.org 上有从入门到进阶的完整教程系列,涵盖 CTAS 物化、Parquet 优化、数据产品变现等主题,持续更新中。

📺 Watch video tutorials → Olap Studio YouTube

Subscribe for more DuckDB & AI automation tutorials

使用 Hugo 构建
主题 StackJimmy 设计

⚠️ 本站为独立社区项目,与 DuckDB 基金会及 DuckDB 官方项目无任何从属、背书或赞助关系。

"DuckDB" 是 DuckDB 基金会的注册商标,本站仅以事实描述方式使用该名称。

本站内容仅供教育与社区推广用途,不构成任何商业服务。