WRITING / 2018.08.12

工程日志:DAG 节点的多次唤醒与批量下发

当一个节点需要边接收上游结果边执行时,如何扩展 DAG 的单次完成语义,同时控制调度放大。

经典 DAG 节点通常只经历一次就绪、一次执行和一次完成。但媒体分片、流式聚合等任务希望在上游部分结果到达后立即工作,而不是等待全部输入完成。

状态语义的变化

多次唤醒后,节点需要区分:

  • 收到新的输入批次;
  • 本批次处理完成;
  • 上游已经封口,不会再有新输入;
  • 节点最终完成,可以释放后继节点。

如果仍用一个布尔 completed 表示全部状态,重复通知和晚到数据会很快制造竞态。

批次协议

每次输入携带节点、批次和序列标识。执行端按批次幂等处理,调度端记录已确认水位。上游封口后,节点处理完剩余批次才进入最终完成态。

输入批次 1 ─┐
输入批次 2 ─┼→ 节点多次执行 → 最终封口 → 后继节点
输入批次 N ─┘

控制调度放大

逐条唤醒会把数据量直接转化为调度请求量。应设置批量大小、等待窗口和最大并发,在延迟与吞吐之间选择,而不是固定使用单条派发。

扩展 DAG 时,最重要的是保持“什么时候算最终完成”仍然唯一且可证明。