Dawn: Orchestrating Task Locality via Inter-Job Dependency Awareness

Dependency-Aware Network Adaptive Scheduling of Data-Intensive Parallel Jobs

2018-08-27
Shaoqi Wang, Wei Chen, Xiaobo Zhou, Liqiang Zhang, Yin Wang
Summary
Problem
Method
Results
Takeaways
Abstract

Deeply integrated into Apache Yarn, Dawn is a dependency-aware, network-adaptive scheduler designed for data-intensive parallel jobs. By leveraging inter-job data dependencies and dynamic network status, it achieves SOTA performance, improving cluster throughput by up to 73% compared to the default Fair Scheduler.

TL;DR

Modern datacenters are often strangled by limited cross-rack bisection bandwidth. While current schedulers optimize for single-job locality, they are "blind" to the fact that one job's output is another job's input. Dawn breaks this silo by analyzing the workload's Directed Acyclic Graph (DAG) to proactively aggregate data, resulting in a 73% boost in throughput and a massive reduction in cross-rack traffic.

The "Independence" Trap in Distributed Scheduling

In data-parallel frameworks like MapReduce or Apache Tez, we’ve long relied on Delay Scheduling to ensure map tasks run where the data sits. However, a critical pain point has emerged: Inter-job dependencies.

In common Machine Learning (ML) pipelines and complex SQL queries, Job B cannot start until Job A finishes. If a scheduler places Job A's tasks randomly across the cluster, Job B's input data ends up fragmented across multiple racks. When Job B starts, it is forced to perform "all-to-all" cross-rack reads, saturating the Top-of-Rack (ToR) switches and aggregation layers.

Methodology: Proactive Aggregation & Network Adaptation

Dawn's core philosophy is: Sometimes you must sacrifice current locality to gain significantly more locality in the future.

1. Online Planning (The "Where")

Dawn uses a Dependency Detector to extract the job relationship matrix. Its Online Plan then performs:

  • Proactive Task Aggregation: In the map stage, even if data is slightly scattered, Dawn forces tasks into "preferred racks" to ensure the resulting intermediate data is already bundled together.
  • Dependency-Aware Assignment: For reduce tasks, Dawn doesn't just look at where the current map outputs are; it looks ahead at which future jobs will need this output, ensuring the final "dependent data" is co-located.

Model Architecture

2. Adaptive Scheduling (The "When")

Network load in clusters is bursty. Dawn monitors real-time traffic to decide:

  • Unsaturated Periods: Use available bandwidth to perform "remote reads" for aggregation.
  • Saturated Periods: Schedule tasks whose data has already been aggregated on the local rack, effectively bypassing the congested core network.

Experimental Results: Slaying the Traffic Monster

Evaluated on 25-node physical clusters and 37-node virtual clusters, Dawn shows a clear linear advantage as DAG depth increases (e.g., Matrix Factorization or PageRank).

  • Throughput Gain: Up to 73% vs. Fair Scheduler.
  • Traffic Reduction: Effectively localized 89% of tasks that would have otherwise triggered cross-rack transfers.

Performance Comparison Figure: Detailed breakdown of cross-rack traffic. Note how Dawn significantly reduces traffic in the 'Dependent Data' and 'Shuffle' phases compared to ShuffleWatcher.

Critical Insight: The "Dependency-Aware" Advantage

The most striking takeaway is the Trade-off Logic. As shown in the "Map Task Scheduling" analysis, Dawn might actually increase traffic during the initial input read (if maps are aggregated away from their source). However, because 80% of jobs in a typical DAG have parents, this small initial "tax" pays dividends by making the subsequent 80% of data flow nearly cost-free at the rack level.

Conclusion

Dawn proves that the network bottleneck isn't just a hardware limitation—it's an information problem. By making the scheduler "aware" of the workload's structure, we can transform random, chaotic network traffic into localized, manageable patterns. While the current implementation focuses on Yarn, the principles of proactive aggregation are highly applicable to modern AI training clusters where GPU-to-GPU communication is the ultimate bottleneck.

Limitations: Dawn assumes the input/intermediate data sizes are relatively predictable, which holds for recurrent ML tasks but may require more complex estimation for highly dynamic web workloads.

Find Similar Papers

Try Our Examples

  • Search for recent papers on cross-rack traffic optimization in multi-tenant Kubernetes or Yarn clusters using Directed Acyclic Graph (DAG) scheduling.
  • Which paper first proposed the concept of 'Delay Scheduling' for data locality, and how has this specific paper improved upon it for inter-job dependencies?
  • Explore if these dependency-aware aggregation techniques have been applied to distributed deep learning training frameworks like PyTorch or TensorFlow to optimize collective communication.
Contents
Dawn: Orchestrating Task Locality via Inter-Job Dependency Awareness
1. TL;DR
2. The "Independence" Trap in Distributed Scheduling
3. Methodology: Proactive Aggregation & Network Adaptation
3.1. 1. Online Planning (The "Where")
3.2. 2. Adaptive Scheduling (The "When")
4. Experimental Results: Slaying the Traffic Monster
5. Critical Insight: The "Dependency-Aware" Advantage
6. Conclusion