经典 DAG 节点通常只经历一次就绪、一次执行和一次完成。但媒体分片、流式聚合等任务希望在上游部分结果到达后立即工作,而不是等待全部输入完成。
状态语义的变化
多次唤醒后,节点需要区分:
- 收到新的输入批次;
- 本批次处理完成;
- 上游已经封口,不会再有新输入;
- 节点最终完成,可以释放后继节点。
如果仍用一个布尔 completed 表示全部状态,重复通知和晚到数据会很快制造竞态。
批次协议
每次输入携带节点、批次和序列标识。执行端按批次幂等处理,调度端记录已确认水位。上游封口后,节点处理完剩余批次才进入最终完成态。
输入批次 1 ─┐
输入批次 2 ─┼→ 节点多次执行 → 最终封口 → 后继节点
输入批次 N ─┘
控制调度放大
逐条唤醒会把数据量直接转化为调度请求量。应设置批量大小、等待窗口和最大并发,在延迟与吞吐之间选择,而不是固定使用单条派发。
扩展 DAG 时,最重要的是保持“什么时候算最终完成”仍然唯一且可证明。