External Merge Sort for Top-K Queries: Eager input filtering guided by histograms
Yannis Chronis, Thanh Do, Goetz Graefe, Keith Peters
Abstract
Business intelligence and web log analysis workloads often use queries with top-k clauses to produce the most relevant results. Values ofk range from small to rather large and sometimes the requested output exceeds the capacity of the available main memory. When the requested output fits in the available memory existing top-k algorithms are efficient, as they can eliminate almost all but the topk results before sorting them. When the requested output exceeds the main memory capacity, existing algorithms externally sort the entire input, which can be very expensive. Furthermore, the drastic difference in execution cost when the memory capacity is exceeded results in an unpleasant user experience. Every day, tens of thousands of production top-k queries executed on F1 Query resort to an external sort of the input. To address these challenges, we introduce a new top-k algorithm that is able to eliminate parts of the input before sorting or writing them to secondary storage, regardless of whether the requested output fits in the available memory. To achieve this, at execution time our algorithm creates a concise model of the input using histograms. The proposed algorithm is implemented as part of F1 Query and is used in production, where significantly accelerates top-k queries with outputs larger than the available memory. We evaluate our algorithm against existing top-k algorithms and show that it reduces I/O traffic and can be up to 11 times faster.
Ask about this paper
Ask your agent about it.
Lune has read the top-tier papers around this one, so every answer names the papers it rests on.
Your agent calls
Lunesearch_papers
Free to start. No credit card required.
Terminal
Install the CLIlune papers get f3f97e54-3a86-4d35-8f1b-951516c7fdf2Cited by top-tier papers3
- Relational Algorithms for Top-k Query EvaluationQichen Wang, Qiyao Luo, Yilei WangSIGMOD 2024 · 5 citations
- Cache-Efficient Top-k Aggregation over High Cardinality Large DatasetsTarique Siddiqui, Vivek R. Narasayya, Marius Dumitru, Surajit ChaudhuriVLDB 2024
- CrocSort: Resource-Efficient, Skew-Resilient Parallel External Merge SortRiki Otaki, Charles Benello, Fuheng Zhao, Aaron J. Elmore et al.VLDB 2026
Related papers
- Approximating Opaque Top-k QueriesJiwon Chang, Fatemeh NargesianSIGMOD 2025 · 2 citations
- Scalable top-k retrieval with SpartaGali Sheffi, Dmitry Basin, Edward Bortnikov, David Carmel et al.PPoPP 2020
- FaScalSQL: A Fast and Scalable GPU-Accelerated SQL Query Engine for Out-of-Memory TablesChaemin Lim, Suhyun Lee, Jinwoo Choi, Kwanghyun Park et al.ICDE 2026
- Parallel Top-K Algorithms on GPU: A Comprehensive Study and New MethodsJingrong Zhang, Akira Naruse, Xipeng Li, Yong WangSC 2023 · 17 citations
- BCCE: Block-Centric GPU Co-Design for Real-Time Range-Top-K Query at ScaleChengying Huan, Ziheng Meng, Zhengyi Yang, Yongchao Liu et al.HPDC 2026
