UW-Madison Logo

The ADvanced Systems Laboratory (ADSL)
Publication abstract

Scale and Performance in Memory-Centric Storage

Chenhao Ye

Department of Computer Sciences
University of Wisconsin-Madison

Abstract:

Modern storage systems increasingly center their designs around fast memory. Across diverse workloads and deployment settings, this shift creates new challenges for achieving scale and performance. On a single multi-tenant node, the system must allocate cache space alongside auxiliary resources (e.g., network and I/O), whose demands correlate with the cache miss ratio. Beyond a single node, systems must scale out through replication, with each replica exploiting multi-core parallelism while preserving consistency. Emerging large language model (LLM) workloads operate on massive data in GPU high-bandwidth memory, requiring a rethink of system abstractions. In this dissertation, we ask: how can we improve the scale and performance of memory-centric storage systems under fairness, consistency, and other workload-specific constraints?

In the first part of this dissertation, we present HARE, a cache-centric multi-resource allocation algorithm that exploits the demand correlations between cache and auxiliary resources. HARE first harvests resources by searching for a cache partition that preserves every tenant's fair-share throughput while reducing the aggregate consumption of auxiliary resources. It then redistributes the harvested resources to improve the throughput of every tenant. We demonstrate HARE's generality with two systems: HopperKV, a cloud-native key-value cache that modifies Redis to cache data from DynamoDB, and BunnyFS, a microkernel-style local filesystem for NVMe SSDs. Evaluations show that HopperKV achieves up to 1.9x higher throughput, and BunnyFS up to 1.4x.

In the second part, we present MultiLog, a highly scalable replication architecture for memory-centric storage. MultiLog scales replication by sharding the log into multiple sublogs across disjoint keyspaces, each of which can be appended and replayed independently. To keep replica reads prefix-consistent across sublogs, MultiLog introduces PrefixLoom, a lightweight timestamp-based protocol that lazily restores consistency. Implemented in Garnet, MultiLog improves append throughput by up to 19x and replay throughput by up to 36x over a single-log baseline. Under a write-heavy replication workload, its replica lags the primary by only ~10 ms at p99.9, two orders of magnitude lower than in the baselines.

In the final part, we present Reference-Oriented Storage (ROS), a new approach to weight transfer in LLM reinforcement learning. ROS presents the illusion that certain versions of the model weights are stored and can be fetched on demand. Underneath, ROS stores no copies of the weights; it instead tracks the workers that hold each version for inference, and serves reads directly from their GPUs. We build TensorHub, a production-quality system that instantiates the ROS idea with topology-aware transfer, model-parallel consistency, and fault tolerance. Across three distinct rollout workloads, TensorHub reduces total GPU stall time by up to 6.7x for standalone rollouts, accelerates weight updates for elastic rollouts by 4.8x, and cuts cross-datacenter rollout stall time by 19x. TensorHub has been deployed at ByteDance for reinforcement-learning training.

Full Paper: PDF   BibTeX

Publications