DOI: 10.1145/3837106 ISSN: 2836-6573

DQSim: Does the Task Assignment Slow Down Your Distributed Database System?

Maximilian Rieger, Thomas Neumann

In distributed OLAP databases, assigning data and tasks to compute nodes is a non-trivial problem with large performance impact. Common strategies include simple round-robin schedulers or minimizing network transfers. However, in some cases these techniques cannot deliver optimal query latency. For example, in scale-out scenarios, existing techniques fail to balance the trade-off between utilization of new nodes and increased network cost from cold caches, leaving substantial room for improvement.

Addressing these problems in existing systems is difficult, because schedulers are often tightly integrated and cannot be easily adapted for experiments. Evaluating a new scheduler in a system is equally difficult, requiring time-consuming and costly experiments for many scenarios.

To tackle this problem, we formalize the distributed task scheduling problem (DTS). Further, we present DQSim, a fine-grained simulator for distributed query execution that estimates the execution time of DTS schedules. This allows us to accurately model the effects of new scheduling algorithms without having to rebuild large parts of a system. DQSim can run thousands of experiments with user-defined system and cluster configurations and queries in seconds to evaluate scheduling techniques quickly.

Finally, we propose two new scheduling algorithms, NetHEFT and cCEFT, that yield good results across a wide range of scenarios. They are on par with round-robin and network-minimizing baselines on simple scenarios where no improvement is expected. In more challenging scenarios, they can reduce query latency by over 2x.