Skip to content

SOS: A High-Performance Distributed Key-Value Store for Large-Scale Online Services

Aug 2026 · Proceedings of the VLDB Endowment · 0 citations · 44 references

Abstract

Large-scale online services—including web search, recommendation, and LLM inference workloads such as Retrieval-Augmented Generation (RAG) and KV-cache offloading—demand storage that handles petabyte-scale data under millisecond tail-latency SLAs. In-memory stores are cost-prohibitive at scale; disk-based systems sacrifice latency for capacity. We present SOS , a distributed key-value store that bridges this gap by guaranteeing at most one physical disk I/O per operation , enforced by three co-designed components: a flat in-memory hash index with O (1) lookup; a disk storage engine with a hierarchical variable-size layout that limits internal fragmentation to ≈7.93%; and a write cache with asynchronous flush that fully decouples client-visible latency from device I/O. Distributed strong consistency is provided by Primary-Based Chain Replication with Two-Phase Commit and an in-chain retry mechanism that reduces write failures to near-zero under transient hardware jitter. Single-node benchmarks show SOS achieves 704 K QPS under uniform YCSB-B (256 client threads) at P99 = 2 ms, and outperforms RocksDB by 10.0 × in throughput and 24.6 × in P99 latency on write-heavy YCSB-A, with write amplification of 1.00× (vs. 14–16×). Deployed at Tencent for several years, SOS manages over 50 PB and 7 trillion records—peaking at 550 M QPS—powering Yuanbao RAG, Hunyuan KV-cache, user portrait and other services with P99.9 read latencies below 5 ms, at a fraction of in-memory cost.

View source

We use cookies to run the site and, with your consent, for analytics and to show ads. See our Cookie Policy.