Dawn: Orchestrating Task Locality via Inter-Job Dependency Awareness
Dependency-Aware Network Adaptive Scheduling of Data-Intensive Parallel Jobs
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.

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.
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.
