dijiekstra commented on issue #18626:
URL:
https://github.com/apache/dolphinscheduler/issues/18626#issuecomment-5615151289
## 📋 核心改造点速查表
### Master侧改造(< 500行)
#### 改造点1:DependencyResolver集成
**位置**:
`dolphinscheduler-master/src/main/java/.../service/DependencyResolverService.java`(新增)
```java
public interface DependencyResolver {
DependencyStatus resolveDependency(Long workflowCode, Long taskCode,
String dependencyGroup);
void registerAssetStateChangeListener(AssetStateChangeListener listener);
}
```
#### 改造点2:ASSET_SENSOR任务类型
**位置**:
`dolphinscheduler-task-plugin/dolphinscheduler-task-asset-sensor/`(新增模块)
- AssetSensorParameters extends AbstractParameters
- AssetSensorLogicTask extends AbstractLogicTask
- AssetSensorTracker implements TaskTracker
#### 改造点3:TaskProcessor中的依赖检查
**位置**:
`dolphinscheduler-master/src/main/java/.../processor/TaskProcessor.java`(改动 ~5行)
```java
// 在submitTask()方法中,发送给Worker前添加
if (task.getType() == TaskType.ASSET_SENSOR) {
if (!dependencyResolver.resolveDependency(...).isReady()) {
return; // 保持WAITING_DEPENDENCY状态
}
}
```
#### 改造点4:DAGExecutionEngine中的工作流启动检查
**位置**: `dolphinscheduler-master/.../engine/DAGExecutionEngine.java`(改动 ~10行)
工作流实例创建后、DAG生成前,可选地提前验证工作流级资产依赖。
#### 改造点5:Master配置和补偿扫描器
**位置**:
`dolphinscheduler-master/.../scanner/AssetEventCompensationScanner.java`(新增)
周期:5~15分钟
职责:
- 补写漏采集的Paimon快照事件
- 重试卡顿的trigger_history记录
---
## 🗂️ 数据库表清单
| 表名 | 主要字段 | 唯一约束 | 索引 |
|-----|--------|--------|------|
| t_ds_asset | assetKey, catalogName, databaseName, tableName | uk_asset_key
| idx_asset_key |
| t_ds_asset_event | eventId, sourceType, assetKey, snapshotId |
uk_dedup(source, asset, event_type, snapshot) | idx_asset_snapshot,
idx_event_type_time |
| t_ds_asset_state | assetKey, latestSnapshotId, qualityStatus |
uk_asset_key | idx_asset_key |
| t_ds_asset_dependency | workflowCode, taskCode, assetKey, dependencyGroup
| - | idx_asset_key, idx_workflow |
| t_ds_asset_trigger_history | triggerKey, workflowCode, taskCode |
uk_trigger_key | idx_status_time |
| t_ds_asset_event_consumer_offset | sourceType, assetKey, consumerGroup |
uk_offset | - |
---
## 🔍 关键设计要点
### 1. trigger_key组成规则
```
workflow_{workflowCode}_group_{dependencyGroup}_snapshot_{id1}_{id2}_...
```
- 标识一次完整的触发决策
- 唯一约束保证幂等性
- 无需分布式锁
### 2. 三层依赖判定
### 3. 幂等性保证
- 同一trigger_key最多成功插入一次
- 数据库唯一约束天然处理并发竞态
- 补偿重试时自动去重
---
## ⚠️ 8项已识别风险详解
| # | 风险 | 影响 | 缓解措施 |
|----|-----|------|--------|
| 1 | 补数(backfill)语义 | 历史快照是否触发? | 设计backfill_mode标记 |
| 2 | 多Master竞态 | trigger_history插入冲突 | uk_trigger_key约束 + 异常捕获 |
| 3 | 任务轮询压力 | DB查询压力 | 缓存 + 事件驱动替代轮询 |
| 4 | CommandType兼容性 | 现有逻辑失效 | 全仓库搜索使用点适配 |
| 5 | 跨项目权限 | 权限模型不清 | 后续单独设计 |
| 6 | Push/Poll去重 | 时钟不同步 | 版本号比较 + 唯一约束 |
| 7 | exactly-once范围 | 过度承诺 | 明确边界:trigger_history层面 |
| 8 | 分区混合依赖 | 表达不清 | 待后续补充规范 |
---
## 📦 测试覆盖清单
### 单元测试
- [ ] 幂等性:同一trigger_key不会被创建两次
- [ ] 去重:重复事件被数据库约束拒绝
- [ ] 乱序:旧快照不会覆盖新快照
- [ ] AND/OR:多资产依赖的组合判定
### 集成测试
- [ ] 快照推进流程:event → state → dependency ready → trigger → dag execute
- [ ] 多Master场景:并发条件下的幂等性
- [ ] 补偿流程:漏采集事件的补写、卡顿触发的重试
### 端到端测试
- [ ] GMV场景:完整的ODS→DWD→DWS→ADS链路
- [ ] 性能测试:并发snapshot推进下的响应时间
- [ ] 故障测试:Master宕机、网络抖动等场景恢复
---
## 🚀 建议的社区讨论重点
1. **Paradigm Shift适用范围**
- 确认是否所有湖仓场景都应该采用纯snapshot模型?
- 是否允许混合时间+快照触发?
2. **Master改造风险评估**
- 5个改造点是否充分?是否遗漏隐藏的依赖?
- TaskProcessor#submitTask()的1行检查是否影响性能?
3. **性能与可观测性**
- CompensationScanner的5~15分钟周期是否合理?
- 22+指标是否覆盖了关键决策路径?
4. **生态与扩展**
- MVP只支持Paimon,Iceberg/Hudi的适配成本如何估算?
- Push事件接收的认证和限流策略如何设计?
---
💬 **欢迎评论与讨论!**
--
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]