进阶 dagster.io 2026-10-10 01:48:56 · 6 阅读
第16章 现代 Lakehouse 编排实战:Dagster 与 Airflow 深度对比
# 第16章:现代 Lakehouse 编排指南:Dagster 与 Airflow 对比
从一个出色的 Lakehouse 教程出发,用现代数据编排进一步升华
灵感来源
我最近读到了一篇非常棒的文章:《Data Engineer Things 团队:用 dbt、Airflow、Trino、Iceberg 和 MinIO 在笔记本电脑上搭建 Lakehouse》。这是一份优秀的教程,展示了如何利用现代开源工具在笔记本上构建完整的 Lakehouse 技术栈。
读完他们的实现方案后,我心想:“这已经很棒了,但如果用 Dagster 来优化一下,岂不是更好?”
于是,我决定用 Dagster 替代 Airflow,重新实现同样的 Lakehouse 架构,结果令人印象深刻。我强烈建议您先阅读原文,体验他们的实现方案,再与我这里构建的版本进行对比。
为什么这个对比很重要
两种实现方案都使用了相同的 Lakehouse 核心技术:
- 使用 MinIO 作为 S3 兼容的对象存储
- 使用 Trino 作为分布式 SQL 查询引擎
- 使用 Iceberg 作为 ACID 表格式
- 使用 Nessie 作为数据目录和版本管理
- 使用 dbt 进行分析转换
关键区别在哪里?在于编排层。这次对比凸显了选择合适的编排工具如何显著提升数据平台的能力。
我优化了哪些方面
1. 基于时间窗口的智能分区
*原文内容:* 基础的数据日处理,缺乏精细的分区机制。
*我的改进:* 使用 Dagster 的 TimeWindowPartitionsDefinition 实现规范的分区键管理 ```python from dagster import TimeWindowPartitionsDefinition daily_partitions = TimeWindowPartitionsDefinition( cron_schedule="0 0 * * *", # 每日零点作为边界 start="2024-01-01", end="2027-01-01", # 为未来预留分区空间 fmt="%Y-%m-%d", ) ``` 这带来了以下优势: - 支持回填(Backfill)——高效处理历史数据 - 支持增量处理——仅处理新分区 - 支持分区感知传感器——仅在新数据到达时触发 - 支持选择性重跑——重新处理特定日期范围的数据 2. 基于事件的架构与传感器 *原文内容:* 仅基于时间表的调度。
我的增强:智能传感器,监控 MinIO 中的新数据 @sensor( minimum_interval_seconds=60, job_name="daily_pipeline", required_resource_keys={"minio"} ) def sales_sensor(context: SensorEvaluationContext): """Monitor MinIO for new sales files and trigger processing.""" minio = context.resources.minio # Check for new files in the last 7 days for days_back in range(1, 8): date = datetime.now() - timedelta(days=days_back) date_str = date.strftime("%Y-%m-%d") prefix = f"raw/sales/dt={date_str}/" file_count = minio.count_objects(prefix) if file_count > previous_count: yield RunRequest( partition_key=date_str, tags={"trigger": "new_data_detected"} ) 优势: ✅ 响应数据到达——无需等待定时调度 ✅ 基于 Cursor 的追踪——记住哪些已处理 ✅ 资源高效利用——仅在需要时运行 ✅ 支持多个传感器——不同数据类型可配置不同逻辑 3. 纯 SQL 的 Lakehouse 模式 原文章:使用 Python 加载数据。
我的增强:纯 Trino SQL,采用 Hive 到 Iceberg 的 CTAS 模式 @asset( required_resource_keys={"trino"}, kinds={"trino", "hive", "iceberg"} ) def product_category_data(context: AssetExecutionContext): """Create Iceberg table from MinIO using pure SQL.""" trino = context.resources.trino # Step 1: Create Hive external table trino.execute_statement(""" CREATE TABLE hive.raw.product_category ( category_id BIGINT, category_name VARCHAR, created_date TIMESTAMP ) WITH ( external_location = 's3a://lake/raw/product_category/', format = 'PARQUET' ) """) # Step 2: CTAS to Iceberg managed table trino.execute_statement(""" CREATE TABLE iceberg.raw.product_category WITH (format = 'PARQUET') AS SELECT * FROM hive.raw.product_category """) 优势: ✅ 可扩展——充分利用 Trino 的分布式处理能力 ✅ 没有 Python 瓶颈——纯 SQL 操作 ✅ 性能更佳——catalog 之间直接传输数据 ✅ 符合行业标准——遵循 lakehouse 最佳实践 4. 全面的数据质量框架 原文章:基础的 dbt 测试。
我的增强:基于 Dagster 资产检查的多层次数据质量 @asset_check(asset="product_category_data", required_resource_keys={"trino"}) def product_category_completeness(context: AssetExecutionContext): """确保商品分类数据符合质量标准。""" trino = context.resources.trino result = trino.execute_query("SELECT COUNT(*) FROM iceberg.raw.product_category") count = result[0][0] if result else 0 passed = count >= 5 # 应至少包含 5 个分类 return AssetCheckResult( passed=passed, metadata={"category_count": count}, description=f"商品分类数量:{count}(要求 >= 5)" ) 额外优势: ✅ 新鲜度策略 — 监控数据新鲜度的 SLA ✅ 资产检查 — 程序化的数据质量验证 ✅ dbt 测试集成 — 复用现有的 dbt 测试 ✅ 血缘感知质量 — 质量检查遵循数据依赖关系 5. 高级编排功能 原文:基础 DAG 调度。
我的增强:具备资产选择能力的复杂作业定义 # 灵活的作业定义 setup_job = define_asset_job( name="setup_job", selection=AssetSelection.keys( ["raw", "product_category_data"], ["raw", "product_subcategory_data"], ["raw", "product_data"], ["raw", "territory_data"] ), description="从 MinIO 加载维度数据至 Iceberg 表") daily_pipeline = define_asset_job( name="daily_pipeline", selection=AssetSelection.keys(["raw", "daily_sales_data"]) | AssetSelection.keys(["curated", "product_dim_simple"]) | AssetSelection.keys(["marts", "sales_summary"]), partitions_def=daily_partitions, description="每日流程:销售数据 → dbt 整理层 → dbt 分析层") 功能特性: ✅ 资产选择 — 运行管道中的特定子集 ✅ 动态作业 — 根据上游变更自动适配的作业 ✅ 分区感知作业 — 优雅处理时序数据 ✅ 资源优化 — 按作业精细化调配资源 6. 现代开发体验 原文:传统配置。
我的增强方案:现代 Dagster 工具与 dg 命令行界面 # 现代项目结构 dg dev # 启动开发服务器 # - 代码变更时自动热重载 # - 功能丰富的 Web UI,支持血缘关系可视化 # - 集成日志与监控 # - 带元数据的资产目录 # 现代项目结构dg dev # 启动开发服务器# - 代码变更时自动热重载# - 功能丰富的 Web UI,支持血缘关系可视化# - 集成日志与监控# - 带元数据的资产目录 开发者的益处: ✅ 自动发现——自动检测资产 ✅ 热重载——即时查看变更 ✅ 丰富界面——可视化开发流水线 ✅ 集成测试——对资产进行独立测试 架构对比 原始方案(Airflow) 原始数据 → Python 处理 → Iceberg 表 → dbt → 分析 ↑ 定时调度 ↑ 手动 ETL ↑ 基础 DAG ↑ 基于时间
我的增强方案(Dagster) 原始数据 → Hive 外部表 → Iceberg 托管表 → dbt → 分析 ↑ 传感器
事件驱动
纯 SQL CTAS
分布式
处理
资产校验
数据质量
框架
智能
任务
SLA
*我的改进:* 使用 Dagster 的 TimeWindowPartitionsDefinition 实现规范的分区键管理 ```python from dagster import TimeWindowPartitionsDefinition daily_partitions = TimeWindowPartitionsDefinition( cron_schedule="0 0 * * *", # 每日零点作为边界 start="2024-01-01", end="2027-01-01", # 为未来预留分区空间 fmt="%Y-%m-%d", ) ``` 这带来了以下优势: - 支持回填(Backfill)——高效处理历史数据 - 支持增量处理——仅处理新分区 - 支持分区感知传感器——仅在新数据到达时触发 - 支持选择性重跑——重新处理特定日期范围的数据 2. 基于事件的架构与传感器 *原文内容:* 仅基于时间表的调度。
我的增强:智能传感器,监控 MinIO 中的新数据 @sensor( minimum_interval_seconds=60, job_name="daily_pipeline", required_resource_keys={"minio"} ) def sales_sensor(context: SensorEvaluationContext): """Monitor MinIO for new sales files and trigger processing.""" minio = context.resources.minio # Check for new files in the last 7 days for days_back in range(1, 8): date = datetime.now() - timedelta(days=days_back) date_str = date.strftime("%Y-%m-%d") prefix = f"raw/sales/dt={date_str}/" file_count = minio.count_objects(prefix) if file_count > previous_count: yield RunRequest( partition_key=date_str, tags={"trigger": "new_data_detected"} ) 优势: ✅ 响应数据到达——无需等待定时调度 ✅ 基于 Cursor 的追踪——记住哪些已处理 ✅ 资源高效利用——仅在需要时运行 ✅ 支持多个传感器——不同数据类型可配置不同逻辑 3. 纯 SQL 的 Lakehouse 模式 原文章:使用 Python 加载数据。
我的增强:纯 Trino SQL,采用 Hive 到 Iceberg 的 CTAS 模式 @asset( required_resource_keys={"trino"}, kinds={"trino", "hive", "iceberg"} ) def product_category_data(context: AssetExecutionContext): """Create Iceberg table from MinIO using pure SQL.""" trino = context.resources.trino # Step 1: Create Hive external table trino.execute_statement(""" CREATE TABLE hive.raw.product_category ( category_id BIGINT, category_name VARCHAR, created_date TIMESTAMP ) WITH ( external_location = 's3a://lake/raw/product_category/', format = 'PARQUET' ) """) # Step 2: CTAS to Iceberg managed table trino.execute_statement(""" CREATE TABLE iceberg.raw.product_category WITH (format = 'PARQUET') AS SELECT * FROM hive.raw.product_category """) 优势: ✅ 可扩展——充分利用 Trino 的分布式处理能力 ✅ 没有 Python 瓶颈——纯 SQL 操作 ✅ 性能更佳——catalog 之间直接传输数据 ✅ 符合行业标准——遵循 lakehouse 最佳实践 4. 全面的数据质量框架 原文章:基础的 dbt 测试。
我的增强:基于 Dagster 资产检查的多层次数据质量 @asset_check(asset="product_category_data", required_resource_keys={"trino"}) def product_category_completeness(context: AssetExecutionContext): """确保商品分类数据符合质量标准。""" trino = context.resources.trino result = trino.execute_query("SELECT COUNT(*) FROM iceberg.raw.product_category") count = result[0][0] if result else 0 passed = count >= 5 # 应至少包含 5 个分类 return AssetCheckResult( passed=passed, metadata={"category_count": count}, description=f"商品分类数量:{count}(要求 >= 5)" ) 额外优势: ✅ 新鲜度策略 — 监控数据新鲜度的 SLA ✅ 资产检查 — 程序化的数据质量验证 ✅ dbt 测试集成 — 复用现有的 dbt 测试 ✅ 血缘感知质量 — 质量检查遵循数据依赖关系 5. 高级编排功能 原文:基础 DAG 调度。
我的增强:具备资产选择能力的复杂作业定义 # 灵活的作业定义 setup_job = define_asset_job( name="setup_job", selection=AssetSelection.keys( ["raw", "product_category_data"], ["raw", "product_subcategory_data"], ["raw", "product_data"], ["raw", "territory_data"] ), description="从 MinIO 加载维度数据至 Iceberg 表") daily_pipeline = define_asset_job( name="daily_pipeline", selection=AssetSelection.keys(["raw", "daily_sales_data"]) | AssetSelection.keys(["curated", "product_dim_simple"]) | AssetSelection.keys(["marts", "sales_summary"]), partitions_def=daily_partitions, description="每日流程:销售数据 → dbt 整理层 → dbt 分析层") 功能特性: ✅ 资产选择 — 运行管道中的特定子集 ✅ 动态作业 — 根据上游变更自动适配的作业 ✅ 分区感知作业 — 优雅处理时序数据 ✅ 资源优化 — 按作业精细化调配资源 6. 现代开发体验 原文:传统配置。
我的增强方案:现代 Dagster 工具与 dg 命令行界面 # 现代项目结构 dg dev # 启动开发服务器 # - 代码变更时自动热重载 # - 功能丰富的 Web UI,支持血缘关系可视化 # - 集成日志与监控 # - 带元数据的资产目录 # 现代项目结构dg dev # 启动开发服务器# - 代码变更时自动热重载# - 功能丰富的 Web UI,支持血缘关系可视化# - 集成日志与监控# - 带元数据的资产目录 开发者的益处: ✅ 自动发现——自动检测资产 ✅ 热重载——即时查看变更 ✅ 丰富界面——可视化开发流水线 ✅ 集成测试——对资产进行独立测试 架构对比 原始方案(Airflow) 原始数据 → Python 处理 → Iceberg 表 → dbt → 分析 ↑ 定时调度 ↑ 手动 ETL ↑ 基础 DAG ↑ 基于时间
我的增强方案(Dagster) 原始数据 → Hive 外部表 → Iceberg 托管表 → dbt → 分析 ↑ 传感器
事件驱动
纯 SQL CTAS
分布式
处理
资产校验
数据质量
框架
智能
任务
SLA
监控
技术栈
两套实现共享同一套坚实的基础设施:
基础设施层
- Docker Compose —— 本地开发环境
- MinIO —— S3 兼容对象存储
- Trino —— 分布式 SQL 查询引擎
- Nessie —— 支持 Git 式版本管理的数据目录
- Iceberg —— 具备 ACID 特性的开放表格式
数据层
- Raw Zone —— MinIO 中的 Parquet 文件(s3://lake/raw/)
- Curated Zone —— 带规范 schema 的 Iceberg 表
- Analytics Zone —— 经过 dbt 转换的维度模型
编排层
- 原版:Apache Airflow
- 增强版:Dagster,采用现代化特性
为什么这对你的组织很重要
卓越运维
- 可靠性:事件驱动处理减少数据遗漏
- 效率:智能分区最小化算力浪费
- 可见性:丰富的监控与血缘追踪
- 可维护性:清晰的资产定义与依赖关系
开发者生产力
- 更快迭代:热重载与集成测试
- 更好调试:详细日志与执行上下文
- 轻松上手:自文档化的资产目录
- 现代工具链:业界标准的开发体验
业务价值
- 更快获得洞察:数据一到达即可用
- 更高数据质量:全面的校验框架
- 更低运营成本:高效的资源利用
- 面向未来的架构:基于现代数据栈构建
快速上手
我建议你:
- 阅读原文 —— 它是理解 lakehouse 概念的绝佳入门
- 试用他们的实现 —— 熟悉核心技术
- 对比 Dagster 版本 —— 看看编排层的差异
- 两者都动手试试 —— 理解各自的取舍
快速启动我的实现
# 克隆并启动基础设施
git clone https://github.com/eric-thomas-dagster/dagster-lakehouse
cd dagster-lakehouse
./start_lakehouse.sh
# 生成示例数据
python generate_adventureworks_data.py
# 启动 Dagster
cd dagster_project
dg dev
然后访问:
- Dagster UI:http://localhost:3000
- Trino:http://localhost:8080
- MinIO 控制台:http://localhost:9000
结语
原版 lakehouse 文章为理解现代数据架构打下了很好的基础。通过引入 Dagster 的高级编排能力,我打造了一个更健壮、更可扩展、更易维护的数据平台。
关键洞察是什么?技术选型在每一层都很重要。核心 lakehouse 技术(MinIO、Trino、Iceberg、dbt)提供了基础,但编排层决定了你能多高效地发挥它们的价值。
Dagster 以资产为中心的理念,结合其 sensors、分区和数据质量功能,构建出的数据平台不仅能用,而且是真正生产可用的。
两种方案都试试
我强烈建议你体验两套实现:
- 原版展示了出色的 lakehouse 基础实践
- 我的增强版演示了现代编排能力
通过对比,你会深刻理解编排选型如何影响数据平台的能力、可维护性和开发者体验。
想构建自己的增强版 lakehouse?欢迎查看我的实现:https://github.com/eric-thomas-dagster/dagster-lakehouse,也欢迎告诉我你的想法!
更多资源
- 原版 Lakehouse 文章
- Dagster 文档
- Apache Iceberg
- Trino
- dbt
- MinIO
有反馈或问题?欢迎在 Slack 或 GitHub 上参与讨论。想和我们一起工作?查看我们的招聘职位。想看更多这类内容?在 LinkedIn 上关注我们。