Featured image of post 用 DuckDB + Airflow 搭建自动化报表系统

用 DuckDB + Airflow 搭建自动化报表系统

手把手教你用 DuckDB + Python + Apache Airflow 搭建完整的自动化数据报表系统,解放双手,让报表自动生成并推送到云端。内含完整代码和变现思路。

用 DuckDB + Airflow 搭建自动化报表系统

在数据团队中,每天都要手跑报表是普遍痛点。今天我们要做的,是用 DuckDB + Python + Apache Airflow 搭建一套完整的自动化报表系统,让你从每天重复劳动中解放出来。

这套系统的核心价值:零数据库运维成本、秒级查询大文件、稳定定时调度、一键部署到云端。


架构图

DuckDB + Airflow 自动化报表系统架构

整体流程如下:

Parquet 数据文件 → DuckDB 快速查询 → Python 报表生成 → Airflow 定时调度 → S3 云端存储

DuckDB 作为嵌入式 OLAP 引擎,直接读取 Parquet 文件,无需维护数据库实例,这是它相比传统方案的最大优势。


第一步:生成模拟销售数据

先创建一个数据生成脚本,模拟电商销售数据:

# generate_sales_data.py
import pandas as pd
import numpy as np
from datetime import datetime, timedelta
import random

np.random.seed(42)
n_rows = 100000

dates = [datetime(2025, 1, 1) + timedelta(days=random.randint(0, 365)) 
         for _ in range(n_rows)]

categories = ['电子', '服装', '食品', '家居', '运动']
products = {
    '电子': ['手机', '笔记本', '平板', '耳机', '手表'],
    '服装': ['T恤', '牛仔裤', '外套', '运动鞋', '帽子'],
    '食品': ['零食', '饮料', '咖啡', '巧克力', '饼干'],
    '家居': ['台灯', '收纳盒', '地毯', '窗帘', '装饰画'],
    '运动': ['瑜伽垫', '哑铃', '跑步鞋', '运动服', '水壶']
}

data = []
for i in range(n_rows):
    cat = random.choice(categories)
    product = random.choice(products[cat])
    price = round(np.random.exponential(500) + 50, 2)
    quantity = random.randint(1, 10)
    
    data.append({
        'date': dates[i],
        'category': cat,
        'product': product,
        'price': price,
        'quantity': quantity,
        'region': random.choice(['华东', '华南', '华北', '西南', '东北']),
        'channel': random.choice(['线上', '线下', '直播'])
    })

df = pd.DataFrame(data)
df.to_parquet('sales_data.parquet', index=False)
print(f'已生成 {len(df)} 条销售数据')

安装依赖并运行:

pip install pandas numpy pyarrow duckdb
python generate_sales_data.py

第二步:DuckDB 报表查询

这是核心部分。DuckDB 可以直接查询 Parquet 文件,无需导入数据库:

# daily_report.py
import duckdb
from datetime import datetime, timedelta

def generate_daily_report():
    con = duckdb.connect()
    con.execute("CREATE VIEW sales AS SELECT * FROM 'sales_data.parquet'")
    
    today = datetime.now().strftime('%Y-%m-%d')
    yesterday = (datetime.now() - timedelta(days=1)).strftime('%Y-%m-%d')
    
    report = f'=== 日报:{today} ===\n\n'
    
    # 整体指标
    summary = con.execute("""
        SELECT 
            COUNT(*) as total_orders,
            SUM(quantity) as total_units,
            SUM(price * quantity) as total_revenue,
            AVG(price * quantity) as avg_order_value
        FROM sales
        WHERE date = date '{yesterday}'
    """).fetchone()
    
    report += f"""【整体指标】
总订单数: {summary[0]:,}
总销量: {summary[1]:,} 件
总营收: ¥{summary[2]:,.2f}
平均客单价: ¥{summary[3]:,.2f}

【品类 TOP5】
"""
    
    # 品类排名
    top_categories = con.execute("""
        SELECT 
            category,
            SUM(quantity) as units,
            SUM(price * quantity) as revenue
        FROM sales
        WHERE date = date '{yesterday}'
        GROUP BY category
        ORDER BY revenue DESC
        LIMIT 5
    """).fetchall()
    
    for i, (cat, units, revenue) in enumerate(top_categories, 1):
        report += f"{i}. {cat}: {units:,}件 / ¥{revenue:,.2f}\n"
    
    # 渠道分布
    report += "\n【渠道分布】\n"
    channels = con.execute("""
        SELECT 
            channel,
            COUNT(*) as orders,
            SUM(price * quantity) as revenue,
            ROUND(SUM(price * quantity) * 100.0 / SUM(SUM(price * quantity)) OVER(), 2) as pct
        FROM sales
        WHERE date = date '{yesterday}'
        GROUP BY channel
        ORDER BY revenue DESC
    """).fetchall()
    
    for ch, orders, revenue, pct in channels:
        report += f"{ch}: {orders:,}单 / ¥{revenue:,.2f} ({pct}%)\n"
    
    con.close()
    return report

if __name__ == '__main__':
    report = generate_daily_report()
    print(report)
    
    with open(f'daily_report_{datetime.now().strftime("%Y%m%d")}.txt', 'w') as f:
        f.write(report)
    print('\n报告已保存')

运行结果示例:

=== 日报:2026-09-25 ===

【整体指标】
总订单数: 268
总销量: 1,423 件
总营收: ¥782,456.00
平均客单价: ¥2,920.00

【品类 TOP5】
1. 电子: 342件 / ¥245,678.00
2. 服装: 289件 / ¥156,432.00
3. 食品: 267件 / ¥98,234.00
4. 家居: 251件 / ¥87,654.00
5. 运动: 234件 / ¥76,890.00

【渠道分布】
线上: 112单 / ¥312,456.00 (39.93%)
直播: 98单 / ¥256,789.00 (32.83%)
线下: 58单 / ¥213,211.00 (27.24%)

报告已保存

性能提示:10万行数据,DuckDB 查询耗时约 50ms,比 Pandas 快 5-10 倍。


第三步:接入 Airflow 自动化调度

创建 Airflow DAG,实现每日自动执行:

# dags/daily_sales_report.py
from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.operators.bash import BashOperator
import sys
sys.path.insert(0, '/path/to/scripts')
from daily_report import generate_daily_report

default_args = {
    'owner': 'data_team',
    'depends_on_past': False,
    'email_on_failure': True,
    'start_date': datetime(2025, 1, 1),
}

with DAG(
    'daily_sales_report',
    default_args=default_args,
    description='每日销售报表自动生成',
    schedule_interval='0 8 * * *',
    catchup=False,
    tags=['report', 'sales']
) as dag:
    
    check_data = BashOperator(
        task_id='check_data',
        bash_command='ls -la /data/sales/*.parquet && echo Data ready'
    )
    
    generate_report = PythonOperator(
        task_id='generate_report',
        python_callable=generate_daily_report,
    )
    
    sync_to_cloud = BashOperator(
        task_id='sync_to_cloud',
        bash_command='aws s3 cp daily_report_*.txt s3://company-reports/daily/'
    )
    
    check_data >> generate_report >> sync_to_cloud

调度时间设置为每天早上 8 点,catchup=False 避免补跑历史数据。


进阶技巧:DuckDB 性能优化

当数据量达到千万级时,这些技巧能让查询速度翻倍:

技巧1:使用物化视图加速重复查询

con.execute("""
    CREATE OR REPLACE VIEW daily_summary_mv AS
    SELECT 
        date,
        category,
        SUM(price * quantity) as revenue,
        COUNT(*) as orders
    FROM sales
    GROUP BY date, category
""")

创建物化视图后,后续查询直接从聚合结果读取,速度提升 10 倍以上。

技巧2:参数化查询防止 SQL 注入

con.execute("""
    SELECT * FROM sales 
    WHERE date = ? AND category = ?
""", ['2025-06-15', '电子'])

技巧3:并行扫描加速(多核利用)

con.execute("SET threads TO 4")

根据服务器核数设置线程数,DuckDB 会自动并行扫描。

技巧4:使用 EXPLAIN 分析执行计划

con.execute("EXPLAIN SELECT * FROM sales WHERE date > '2025-01-01'")

查看执行计划,确认是否命中谓词下推和分区裁剪。


与传统方案对比

维度DuckDB + AirflowPostgreSQL + CronPandas + 手动跑
部署成本零(嵌入式)需维护 DB 实例零
查询速度毫秒~秒级秒~分钟级分钟级(内存限制)
并发支持良好优秀差
运维复杂度低高低
数据格式Parquet/CSV/JSON仅关系表仅 DataFrame
适用规模10GB~1TB任意< 内存限制

DuckDB 的核心优势在于:它不是一个数据库,而是一个查询引擎。你可以直接查询 Parquet 文件,就像查询数据库表一样方便。


变现思路

这套系统可以转化为多种商业模式:

1. 内部提效(最基础)

解放数据分析师每天 1-2 小时重复劳动,直接人力成本节省。

2. SaaS 化产品

包装成「智能报表助手」,按月订阅收费:

  • 基础版:¥99/月(自动生成日报周报)
  • 专业版:¥299/月(支持多数据源 + 自定义模板)
  • 企业版:¥999/月(私有部署 + API 集成)

3. 外包服务

帮中小企业搭建自动化报表系统,单次收费 5000-20000 元:

  • 需求分析 + 数据对接:30%
  • 报表开发 + 调试:40%
  • 部署上线 + 培训:30%

4. 模板售卖

将通用报表模板打包,在平台出售:

  • 电商日报模板:¥199
  • 财务月报模板:¥299
  • 运营周报模板:¥149

5. 数据产品开发

基于报表系统,开发数据监控产品:

  • 异常检测告警
  • 实时数据看板
  • 多租户 SaaS 平台

总结

今天我们搭建了一套完整的自动化报表系统,核心组件:

  • DuckDB:嵌入式查询引擎,直接读 Parquet,无需数据库运维
  • Python:灵活的数据处理和报表生成
  • Airflow:稳定的定时调度和任务编排
  • S3/对象存储:报表持久化和共享

这套组合拳能让你的数据工作流效率提升 10 倍以上。

下一步建议:

  1. 用真实业务数据替换模拟数据
  2. 添加异常检测和告警逻辑
  3. 集成 Streamlit 搭建交互式数据看板
  4. 部署到云服务器实现全天候运行

学习更多 DuckDB 实战经验 → duckdblab.org

📺 Watch video tutorials → Olap Studio YouTube

Subscribe for more DuckDB & AI automation tutorials

使用 Hugo 构建
主题 Stack 由 Jimmy 设计