dijiekstra opened a new issue, #18625:
URL: https://github.com/apache/dolphinscheduler/issues/18625

   # RFC: 湖仓资产版本事件驱动调度设计
   
   ## 概述
   
   本RFC提议为Apache 
DolphinScheduler引入**资产事件驱动调度能力**,以支持现代CDC湖仓架构(Paimon/Iceberg/Hudi)下的数据版本驱动工作流调度。
   
   ## 核心目标
   
   将DolphinScheduler从"时间DAG调度系统"升级为"数据资产版本驱动的调度系统",使得:
   
   ```text
   湖仓表 snapshot 推进
       → 事件被感知并持久化为 AssetEvent
       → 更新资产的当前状态 AssetState
       → 按依赖条件进行匹配评估
       → 触发/释放 DolphinScheduler 工作流或任务实例
       → 执行、观测、失败重试与补偿
   ```
   
   ## 设计概要
   
   ### 适用场景
   
   - ✅ **新一代数仓架构**:ODS→DWD→DWS→ADS全链路基于Paimon/Iceberg/Hudi
   - ✅ **纯快照驱动**:整个DAG完全由上游资产快照版本推进决定,不涉及时间或cron
   - ✅ **多资产协调**:DWD等中间层依赖多个ODS资产快照同时就绪
   - ❌ **不适用**:仍在使用传统时间触发的数仓
   
   ### 核心设计决策
   
   1. **不改Master内核** - 复用现有CommandType/Command/WorkflowInstance/TaskInstance状态机
   2. **插件化集成** - 新增ASSET_SENSOR任务类型,DependencyResolverService作为Service层
   3. **幂等优先** - trigger_key + 数据库唯一约束保证单次触发
   4. **补偿可靠** - Push事件 + Poll补偿 + CompensationScanner兜底
   5. **分阶段实现** - MVP(Paimon单资产)→ 多资产AND/OR → Iceberg/Hudi
   
   ## 关键架构要素
   
   ### 新增数据库表(6张)
   
   ```sql
   t_ds_asset              -- 资产元数据
   t_ds_asset_event        -- 原始事件(证据)
   t_ds_asset_state        -- 资产当前状态(调度真相来源)
   t_ds_asset_dependency   -- 工作流/任务对资产的依赖声明
   t_ds_asset_trigger_history  -- 幂等触发历史
   t_ds_asset_event_consumer_offset  -- 事件消费位点
   ```
   
   ### 三层依赖判定
   
   | 层级 | 判定方式 | 时机 |
   |-----|--------|------|
   | 工作流级 | DependencyResolver 评估 t_ds_asset_dependency | 创建 WorkflowInstance 
前(Strategy C,二期) |
   | 任务级 | DAG 前后置依赖 + ASSET_SENSOR 轮询 + DEPENDENT | TaskProcessor 提交给 Worker 前 
|
   | DAG 内部 | 正常拓扑执行 | 任务完成后自动 |
   
   ### Master改造(< 500行代码)
   
   5个改造点,全部插件化:
   
   1. **DependencyResolver集成** - 新Service评估资产依赖
   2. **ASSET_SENSOR任务** - 新任务类型 + AssetSensorTracker
   3. **工作流启动检查** - WorkflowInstance创建后验证(可选)
   4. **补偿扫描器** - AssetEventCompensationScanner(5~15分钟周期)
   5. **配置项** - dolphinscheduler.asset-event-scheduling.* 
   
   ## 实施路线图
   
   ### Phase 1 (MVP,2~3周)
   - [ ] 6张表DDL + Migration脚本
   - [ ] DependencyResolverService(仅snapshot判定)
   - [ ] ASSET_SENSOR任务插件 + AssetSensorTracker
   - [ ] Paimon AssetEventScanner(轮询模式)
   - [ ] 端到端测试:快照推进 → 任务释放 → DAG执行
   
   ### Phase 2 (扩展,3~4周)
   - [ ] 多资产AND/OR依赖
   - [ ] Watermark/Quality/Schema条件
   - [ ] CompensationScanner补偿逻辑
   - [ ] 工作流级ASSET_EVENT_TRIGGER
   
   ### Phase 3 (生态)
   - [ ] Iceberg/Hudi事件源
   - [ ] Push事件接收HTTP端点
   - [ ] UI:资产大盘 + 依赖可视化
   - [ ] 告警规则和可观测性仪表盘
   
   ## 端到端示例:GMV场景
   
   ```text
   MySQL CDC → Paimon CDC写入
              ↓
   ODS_order_latest, ODS_payment_latest snapshot推进
              ↓
   DWD_trade_event_fact 等待两个ODS快照都READY(AND依赖)
              ↓
   DWD_trade_event_fact 执行并产出快照 snapshot_id=2001
              ↓
   DWS_gmv_daily 等待DWD_trade_event_fact快照推进
              ↓
   DWS_gmv_daily 执行并产出快照 snapshot_id=3001
              ↓
   ADS_gmv_dashboard 等待DWS_gmv_daily快照推进
              ↓
   整个DAG完全由快照版本推进驱动,无需cron或时间表达式
   ```
   
   ## 并发与安全性
   
   ### 多Master场景
   
   - **Command创建去重**:依赖PK或sequence
   - **trigger_history唯一约束**:uk_trigger_key保证单次创建
   - **无需分布式锁**:数据库约束天然处理竞态
   
   ### 幂等性
   
   ```
   trigger_key = workflow_code + dependency_group + snapshot_id_combo
   同一trigger_key最多成功创建一次trigger_history记录
   => 同一资产快照组合最多触发一次工作流
   => 幂等性天然保证
   ```
   
   ## 观测性设计
   
   ### 22+ Prometheus指标
   
   覆盖:事件ingestion、状态更新、依赖解析、触发、补偿各环节
   
   ### 20+结构化日志字段
   
   eventId, assetKey, snapshotId, triggerKey, workflowCode, taskCode, status, 
reason等
   
   ### 告警规则
   
   8条告警:状态滞后、瀑布失败、补偿积压、触发失败等
   
   ## 对现有架构的影响
   
   ### ✅ 不改动(零风险)
   
   - WorkflowInstance/TaskInstance核心状态机
   - DAG前后置依赖表达和执行
   - DEPENDENT任务现有逻辑
   - Worker任务执行生命周期
   - Alert/Monitor/Backfill等周边功能
   
   ### ⚠️ 需改动(插件化,低风险)
   
   - 新增ASSET_SENSOR任务类型(如DependentTask一样)
   - 新增DependencyResolverService(Service层,不改内核)
   - 在TaskProcessor#submitTask()前加1行依赖检查
   - 新增CommandType.ASSET_EVENT_TRIGGER(可选)
   
   ### 代码改动规模
   
   **预估 < 500行**(相对于DolphinScheduler数万行代码量极小)
   
   ## 风险与开放问题
   
   ### 8项已识别风险
   
   1. 补数(backfill)与事件驱动的语义定义
   2. 多Master并发触发的竞态窗口
   3. 任务轮询频率与数据库压力
   4. CommandType语义扩展的兼容性
   5. 跨catalog/跨项目资产依赖的权限模型
   6. Push/Poll双路径去重的时钟/顺序依赖
   7. exactly-once声明范围
   8. 分区与无分区表的混合依赖
   
   所有风险都已分析,大部分可通过设计规范规避。
   
   ## 设计文档
   
   **完整RFC文档**(1178行)已准备在:
   - `/docs/docs/zh/architecture-design/lakehouse-asset-event-scheduling.md`
   
   包含:
   - ✅ 背景与目标(paradigm shift说明)
   - ✅ 现有基础设施评估(Command/DEPENDENT等)
   - ✅ 总体架构(领域模型、事件流、数据库表、触发链路)
   - ✅ 集成方案(Strategy A/C、API、UI)
   - ✅ **Master改造详解**(4.5章节,346行新增)
   - ✅ 端到端示例(GMV场景完整演绎)
   - ✅ 观测性设计(指标、日志、告警)
   - ✅ 资产标识规范(AssetKey命名、注册发现)
   - ✅ 实施路线图(Phase 1/2/3)
   - ✅ 架构升级指南
   - ✅ FAQ(7个常见问题)
   - ✅ 风险分析(8项)
   
   ## 讨论与反馈
   
   欢迎社区的反馈和讨论:
   
   - **功能完整性**:是否遗漏必要的元数据、事件类型或依赖条件?
   - **Master改造**:5个改造点的代码位置和改动规模是否合理?
   - **性能影响**:CompensationScanner的周期、BatchSize是否合理?
   - **权限模型**:跨项目资产依赖的权限设计如何规范?
   - **补偿策略**:补数场景下对历史快照是否触发?
   - **生态支持**:Iceberg/Hudi的事件模型是否适配?
   
   ## 下一步
   
   1. 社区评审本RFC
   2. 确认paradigm shift的适用范围和限制
   3. 讨论MVP的工作量分配
   4. 启动Phase 1开发(如果获得批准)
   
   ---
   
   **相关标签**:design, enhancement, lakehouse, event-driven, scheduling, paimon, 
iceberg, hudi
   
   **优先级**:Medium(新功能,不影响现有功能)
   
   **类型**:RFC / Discussion / Design
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: 
[email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to