| Total: 44
Cloud local storage has been a popular service among many vendors thanks to its near-physical performance and affordable price. In this paper, we revisit the evolution of the cloud local storage at Alibaba Cloud. We systematically analyze and evaluate the motivations, architectures and pros/cons of different methodologies from user space stack to hardware offloading. We also explore the future of local storage including a hybrid solution by integrating Elastic Block Storage to achieve better performance, availability and cost efficiency.
TapeOBS is an archive storage service offered by Huawei Cloud, which delivers high cost-efficiency by leveraging tape to store large volumes of archived data. Although tape boasts a low total cost of ownership, its inherent characteristics (e.g., a limited number of drives within a tape library) pose unique challenges when developing a large-scale distributed storage system. To address these challenges, we take a holistic approach in designing TapeOBS. At the high level, we introduce a fully asynchronous tape pool, which supports data scheduling and erasure coding in a batched manner, aligning with the features of tape hardware. Within a tape library, we design a tape-tailored local storage engine and incorporate techniques such as dedicated drives to optimize performance. TapeOBS began its gradual rollout at the end of 2022 and officially started serving customers in 2024. As of this writing, TapeOBS has stored hundreds of petabytes of raw user data.
Read-only compressed file systems have become increasingly popular in space-sensitive scenarios, such as IoT and Docker containers. To construct condensed images, they divide the data into blocks (e.g., 1 MB) and compress blocks separately. However, we observe that block-based compression cannot fully utilize the compression benefits due to the data mixture problem, while its performance issues hinder practical usage. We propose RubikFS, a sort-enhanced read-only file system. Our key idea is to solve data mixture by sorting and clustering similar data chunks in a file system-favored block granularity. This is achieved by similarity sorter, which builds a similarity graph to measure the similarity of data chunks and clusters similar chunks by subgraph partitioning. Moreover, sorting can also group data with the same hotness to minimize read amplification. We then introduce an array of techniques, including data grouper, data chunker, and hotness grouper, to implement condensed and efficient RubikFS. Experiments suggest that, compared to existing read-only compressed file systems, RubikFS increases the compression ratio by up to 42.60% and reduces unnecessary reads by up to 70.70%.
Over the last two decades, with the advent of mobile computing and Internet streaming, Apple has expanded its user base and services significantly. With this growth, we have seen an increased volume and diversity of data storage, ranging from backups, personal photos and videos, to music libraries, TV shows, and live streaming. In this paper, we present ACOS, Apple's object store designed to meet the specific requirements of user-facing and internal services, by accommodating a wide range of content and access patterns. ACOS has been in production for over a decade, storing several exabytes of objects and serving billions of requests per day. With its geo-replicated architecture using both local and regional replication mechanisms, ACOS is cost-efficient and highly scalable, available, and durable. The evaluation of our production deployment shows its throughput and latency performance as well as its resilience to hardware and data center failures. This paper presents the design and evolution of ACOS, evaluates its performance in production, and demonstrates its capacity to scale and support Apple's current storage needs and future growth. Note: This paper's title and abstract have been updated per the errata slip published in the conference proceedings.
AI personal computers (AIPCs) enable the local deployment of large language model (LLM) inference, offering enhanced privacy guarantees and customizable serving. However, such deployments are constrained by limited memory capacity, primarily due to the substantial key-value (KV) cache overhead. This paper introduces SolidAttention, an LLM inference engine which addresses these limitations through a tight co-design of dynamic attention sparsity algorithms and SSD-based storage management. Specifically, to maximize SSD bandwidth utilization, SolidAttention consolidates multiple KV pairs into coarse-grained blocks and implements speculative prefetching mechanisms that exploit temporal locality in sparse attention. By fine-grained orchestration of computation and I/O operations while reusing synchronization points, SolidAttention further minimizes SSD-induced blocking latency. With a 128k-token context, SolidAttention improves the inference speed by up to 3.1× and reduces the KV cache memory footprint by up to 98% without compromising inference accuracy.
Large Language Models (LLMs) are increasingly deployed in agent-based applications with complex prompt structures comprising both invariant and dynamic segments. Existing KV cache reuse strategies—PositionDependent Caching (PDC) and Position-Independent Caching (PIC)—inadequately address these scenarios, imposing either strict positional constraints or introducing significant computational overhead due to Positionally Misaligned KV Drift (PMKD) and window padding problems. We identify a distinct pattern in agent workflows termed Relative-Position-Dependent Caching (RPDC), where reusable segments maintain consistent relative ordering despite absolute position shifts. To address this pattern, we propose CacheSlide, a novel KV cache management system that enhances positional-encoding similarity for fixed segments, computes attention for only a minimal subset of tokens, combines new and cached KVs using learned weights, and implements layer-wise and spill-aware KV-cache optimizations. Our implementation extends vLLM’s KV cache management with Chunked Contextual Position Encoding and Weighted Correction Attention. Experimental evaluation across multiple LLMs and agent benchmarks demonstrates that CacheSlide significantly outperforms state-of-the-art baselines, achieving 3.11-4.3× reduction in latency and 3.5-5.8× improvement in throughput, establishing a new efficiency frontier for agent-based LLM applications.
In interactive LLM serving, historical key–value tensors (KVs) of multi-round conversations are often cached in a two-tier storage system consisting of host memory and SSDs, which provides large capacity at low cost. However, loading KVs from two-tier storage in existing approaches increases serving latency by up to 3.8× and decreases throughput by up to 2.0× compared to an ideal large-memory setting on our interactive conversation workload. This inefficiency arises from poor coordination between compute engine and two-tier storage. This paper proposes Bidaw, an efficient KV caching approach with two-tier storage that enables bidirectional awareness between compute and storage. Bidaw introduces two key mechanisms. First, the compute engine schedules requests with KV-loading latency awareness by separating requests whose KVs reside in different storage layers and reordering them by KV size to reduce blocking. Second, the storage system improves host memory hit rates by leveraging LLM-generated responses to predict user access patterns during KV eviction. For further optimization, Bidaw balances storage footprint against computational savings by selectively caching storage-efficient history tensors. Experiments on our interactive conversation workload and a public multi-round conversation workload of interactive LLM serving show that Bidaw reduces response latency by up to 3.58× and improves throughput by up to 1.83× over state-of-the-art approaches, approaching the theoretical upper bound achieved when all KVs reside entirely in host memory.
This paper examines the model loading bottleneck during the LLM inference startup. Existing solutions often optimize model loading performance at the expense of compatibility. However, compatibility is a crucial factor determining whether a technology can be widely applied in real-world scenarios. This work achieves both high performance and strong compatibility by optimizing the cache policy of the kernel file system. We design PPC, a programmable page cache framework that allows users to customize page cache policies in a non-intrusive, flexible, and lightweight manner. Furthermore, we design MAIO, a cache policy implemented based on PPC, to optimize model loading. MAIO introduces an I/O template-based mechanism to fully utilize SSD bandwidth, XPU affinity, and data locality to enhance the efficiency of prefetching and eviction. Our evaluation shows that MAIO reduces the model loading latency by up to 79% compared to existing optimizations. In a real-world application, MAIO achieves up to 36% improvement in inference throughput over other tested solutions in the elastic deployment scenario.
Approximate Nearest Neighbor Search (ANNS) is widely used in various scenarios. For billion-scale ANNS, on-disk graph-based indexes, which organize the vectors as a graph and store them on disk, are favored for their performance and cost-efficiency. However, existing indexes can not maintain a stable search performance while inserting new vectors. In this paper, we propose to use direct insert, which directly inserts vectors into the on-disk index, rather than buffering them in memory and merging them to disk in batches like existing systems. This approach can even out the interference of insert with frontend search, thus stabilizing the performance. We evaluate direct insert by integrating it into a billion-scale graph-based ANNS index named OdinANN. With a fixed insert rate, OdinANN outperforms state-of-the-art ANNS indexes in search latency and throughput, and it consistently shows stable performance in billion-scale vector datasets.
Disaggregated memory (DM) separates computing and memory resources into distinct resource pools, enhancing resource utilization and scalability. However, this new architecture presents fundamental design challenges on range indexes. Existing works fail to achieve high performance: they either suffer from the network bandwidth bottleneck or are fragile due to high RDMA IOPS demands. The key reason is that they all follow a typical design paradigm that uses private compute-side caching, where each compute server holds a private cache space and aggressively consumes the bandwidth and IOPS between compute servers and memory servers. We propose a new compute-side collaborative design. It offloads data locating and locking operations from memory servers to compute servers and thus fully utilizes unsaturated RDMA resources between compute servers to mitigate bottlenecks on memory servers. We implement a prototype called DMTree. Experiments show that DMTree outperforms existing state-of-the-art range indexes on DM for both point operations (i.e., searches, inserts, and updates) and range operations (i.e., scans) under various workloads and parameter settings.
Cloud-based performance monitoring timeseries systems are emerging due to their flexibility and pay-as-you-go capabilities. However, these systems encounter a major bottleneck in query performance, mainly attributed to the prolonged access latency of cloud storage and metadata redundancy of large number of timeseries. Thus, it is critical to optimize query performance within cloud environment and reduce metadata redundancy. In this paper, we propose CloudTS, which is a novel timeseries data storage model with query optimization for cloud storage. CloudTS separately manages metadata and data, and introduces an efficient global metadata management for both space saving and query speedup. CloudTS also transparently supports the time-partitioned tag-based query model in performance monitoring timeseries systems. For metadata, a global tag dictionary is built to reduce metadata redundancy and a novel timeseries-tag mapping technique with a two-dimension bitmap is designed so the mapping of timeseries and tags can be efficiently accomplished to support tag-based queries. For data, the compressed data chunks are put into objects by timeseries group. We have implemented a fully functional prototype of CloudTS and evaluated it with production timeseries data and synthetic workloads based on Amazon S3. In comparison, Cortex, a cloud-based timeseries system widely adopted by industries, and Apache Parquet and JSON Time Series, two representative cloud storage formats, are utilized in the evaluation. Experimental results show that CloudTS can improve query performance by 1.37x on average compared with Cortex, and outperforms Apache Parquet and JSON Time Series as well.
In cloud block store, indexing is on the critical path of I/O operations and typically resides in memory. With the scaling of users and the emergence of denser storage media, the index has become a primary memory consumer, causing memory strain. Our extensive analysis of production traces reveals that write requests exhibit a strong tendency to target continuous block ranges in cloud storage systems. Thus, compared to current per-block indexing, our insight is that we should directly index block ranges (i.e., range-as-a-key) to save memory. In this paper, we propose RASK, a memory-efficient and high-performance tree-structured index that natively indexes ranges. While range-as-a-key offers the potential to save memory and improve performance, realizing this idea is challenging due to the range overlap and range fragmentation issues. To handle range overlap efficiently, RASK introduces the log-structured leaf, combined with range-tailored search and garbage collection. To reduce range fragmentation, RASK employs range-aware split and merge mechanisms. Our evaluations on four production traces show that RASK reduces memory footprint by up to 98.9% and increases throughput by up to 31.0× compared to ten state-of-the-art indexes.
Mitigating latency fluctuations for distributed key-value (KV) stores is critical, yet it is often hindered by the tight coupling of foreground and background tasks related to data distribution and storage management. Using Cassandra, a widely deployed distributed LSM-tree-based KV store, as a case study, we observe that foreground read tasks are often interfered with by background compaction tasks, yet compaction tasks are critical for achieving high read performance. We propose HATS, a holistic and automated task scheduling framework that judiciously co-schedules read and compaction tasks, so as to mitigate latency fluctuations and achieve load balancing. HATS features coarse-grained and fine-grained replica selection for reads as well as adaptive rate control for compaction. We implement HATS atop Cassandra and demonstrate its improved latency and throughput performance over state-of-the-art distributed LSM-tree-based KV stores.
Input data preprocessing is a common bottleneck when concurrently training multimedia machine learning (ML) models in modern systems. To alleviate these bottlenecks and reduce the training time for concurrent jobs, we present Seneca, a data loading system that optimizes cache partitioning and data sampling for the data storage and ingestion (DSI) pipeline. The design of Seneca contains two key techniques. First, Seneca uses a performance model for the data pipeline to optimally partition the cache for three different forms of data (encoded, decoded, and augmented). Second, Seneca opportunistically serves cached data over uncached ones during random batch sampling so that concurrent jobs benefit from each other. We implement Seneca by modifying PyTorch and demonstrate its effectiveness by comparing it against several state-of-the-art caching systems for DNN training. Seneca reduces the makespan by 45.23% compared to PyTorch and increases data processing throughput by up to 3.45× compared to the next best dataloader.
System-level GPU checkpoint/restore (C/R) enables several critical features such as elastic scaling, task switching, and fault tolerance, for modern GPU workloads in a unified and application-transparent manner. However, existing approaches present fundamental limitations: they fail to simultaneously achieve low C/R latency and low overhead imposed on normal GPU execution, while also lacking efficient support for incremental checkpointing. We propose GCR, a GPU checkpoint/restore system that addresses all these limitations simultaneously. GCR employs a hybrid C/R scheme through control/data separation to deliver low C/R latency and negligible overhead imposed on normal GPU execution. To efficiently support incremental checkpointing, GCR introduces shadow execution on the CPU to reduce the overhead of dirty buffer identification, utilizing dirty templates for both lightweight CPU shadow execution and identification at a fine-grained instruction level. Our evaluations demonstrate that GCR reduces GPU checkpointing latency by 72.1% and 63.6% compared to cuda-ckpt (NVIDIA’s official solution) and PhOS (the current state-of-the-art), respectively, and restoration latency by 54.2% and 87.1%, while imposing negligible overhead (less than 1%). GCR also supports efficient incremental checkpointing, which reduces checkpoint sizes by 86.6% and latency by 43.8%.
The emergence of AI workloads has placed rigorous bandwidth requirements on cloud storage, which are challenging to meet due to inherent hardware restrictions in cost-efficient disaggregated storage architectures, as well as the non-triviality of implementing application-tailored optimizations. This paper presents AITURBO, a cloud storage system for AI jobs with high bandwidth demands. AITURBO first utilizes the high-bandwidth compute fabric between accelerators to meet AI applications’ bandwidth demands without incurring additional storage cost. AITURBO further introduces a simple yet powerful grouped I/O API that allows AITURBO to automatically derive optimized read and write plans at the storage layer. These plans enable optimizations that are comparable or better than application-level ones, because they capture common I/O patterns in AI workloads and have a holistic view from the storage layer’s perspective. Under common AI workloads such as checkpoint reads and writes and KV-cache reads, AITURBO achieves comparable or better performance than state-of-the-art systems, with and without application-level optimizations, including systems such as Megatron, Gemini, and Mooncake, typically with minimal application-level code changes. AITURBO has been deployed in training jobs in HUAWEI’s production cloud to support efficient training workloads.
The development of large language models (LLMs) relies on sophisticated parallel training techniques, involving prolonged training runs with thousands of workers. Checkpointing systems are essential for handling failures in large-scale training. However, existing checkpointing systems are almost offline solutions tailored to specific parallelisms or model architectures. They lack adaptability to diverse parallel strategies and fail to recognize that most model states can be excluded from checkpoints, missing optimization opportunities. In this paper, we present AdaCheck, an adaptive checkpointing system that achieves minimized checkpoint size by characterizing and exploiting state redundancy across various parallelisms, model architectures, and training iterations. We model the state redundancy induced by parallelisms and model architectures using the abstraction tensor redundancy, and propose an offline redundancy utilization method to create checkpoints with a reduced set of states. To fully identify tensor redundancy, we design an efficient redundancy detector, which employs a hash-based data consistency check method and a ring-based communication algorithm. Besides, we introduce a novel online redundancy utilization method, which further reduces checkpoint size by exploiting the state redundancy across training iterations. Experimental results demonstrate that AdaCheck is adaptable to various parallelisms, including irregular parallelisms generated by automatic planners, as well as diverse model architectures, encompassing both dense and sparse architectures. Compared with state-of-the-art checkpointing approaches, AdaCheck can reduce checkpoint size by 6.00–896×, increase the checkpointing frequency by 1.46–111×, and incur almost no overhead on training throughput for LLM training.
File systems are critical OS components that require constant evolution to support new hardware and emerging application needs. However, the traditional paradigm of developing features, fixing bugs, and maintaining the system incurs significant overhead, especially as systems grow in complexity. This paper proposes a new paradigm, generative file systems, which leverages Large Language Models (LLMs) to generate and evolve a file system from prompts, effectively addressing the need for robust evolution. Despite the widespread success of LLMs in code generation, attempts to create a functional file system have thus far been unsuccessful, mainly due to the ambiguity of natural language prompts. This paper introduces SYSSPEC, a framework for developing generative file systems. Its key insight is to replace ambiguous natural language with principles adapted from formal methods. Instead of imprecise prompts, SYSSPEC employs a multi-part specification that accurately describes a file system's functionality, modularity, and concurrency. The specification acts as an unambiguous blueprint, guiding LLMs to generate expected code flexibly. To manage evolution, we develop a DAG-structured patch that operates on the specification itself, enabling new features to be added without violating existing invariants. Moreover, the SYSSPEC toolchain features a set of LLM-based agents with mechanisms to mitigate hallucination during construction and evolution. We demonstrate our approach by generating SPECFS, a concurrent file system. SPECFS demonstrates equivalent level of correctness to that of a manually-coded baseline across hundreds of regression tests. We further confirm its evolvability by seamlessly integrating 10 real-world features from Ext4. Our work shows that a specification-guided approach makes generating and evolving complex systems not only feasible but also highly effective.
We present Cylon, a fast and extensible full-system emulator for CXL-SSDs built on FEMU. Cylon bridges the gap between closed hardware prototypes and slow software simulators by faithfully reproducing sub-μs cache hits and tens-of-μs misses that fall to NAND through a hybrid execution path that mitigates hypervisor trap overheads. Cylon supports configurable caching policies and provides an application-level interface for hardware-software co-design. Validated against a real CXL-SSD prototype, Cylon accurately models performance across a wide range of applications, from microbenchmarks to full-scale workloads. Our evaluation shows that Cylon reproduces realistic latency distributions, executes unmodified applications at near bare-metal speed, and scales to system-level studies. By combining speed, fidelity, and extensibility, Cylon fills a critical gap for evaluating today’s CXL-SSDs and exploring next-generation architectures that blend CXL-enabled memory and storage semantics.
Compute Express Link (CXL) is an emerging industry standard that offers high-performance cache-coherent interconnects to heterogeneous devices, including host CPUs, computation accelerators, and memory devices. It aims to support high system scalability, peer-to-peer communication, and high-speed data transmission. To this end, the latest version of the CXL protocol introduces several new features, including port-based routing, device-managed coherence, and PCIe 6.0 support. However, the absence of CXL hardware and the methodological limitations of existing simulators hinder the exploration of these new architectures. To bridge this gap, we propose Xerxes, a novel simulation framework designed from the ground up to faithfully model the emerging features in the latest CXL protocol. It employs a dedicated interconnect layer to support interconnection within diverse system topologies. It also implements important components to conduct specific functions required by these features. Utilizing Xerxes, we comprehensively explore multiple aspects of CXL systems, including system topologies, device-managed coherences, and impacts of PCIe characteristics, and derive key observations that can inspire new designs of high-performance CXL systems. The codes of Xerxes are open-sourced and available at https://github.com/ChaseLab-PKU/Xerxes.
Flexible Data Placement (FDP) promises to reduce write amplification by steering writes across reclaim unit handles (RUHs), yet outcomes vary widely across devices. This paper presents WARP, the first open emulator and comprehensive study of FDP SSDs. Our cross-device, cross-workload characterization shows that FDP sustains near-1 WAF when RUH isolation aligns with object lifetimes, but fails under misclassification, RUH interference, or adversarial invalidations. WARP reproduces hardware WAF trends while exposing per-RUH dynamics and configurable policies hidden in real devices. With WARP, we explore the firmware design space for FDP and demonstrate policies that reduce WAF beyond current hardware. By combining empirical characterization with a transparent emulator, this work advances FDP research from anecdotal reports to principled understanding and provides a platform for future FDP-aware system design.
This paper presents a scalable OS swap system, ScaleSwap, designed to enhance core and SSD scalability on all-flash swap arrays. Specifically, ScaleSwap first enables a one-to-one swap model where each core exclusively manages its own swap resources, enabling core-centric swap in/out operations. Second, ScaleSwap devises opportunistic inter-core swap assistance, allowing each core to delegate swap metadata access to other cores as needed. Finally, ScaleSwap adopts core-affinity page and LRU management to mitigate LRU lock contention during swap in/out operations. We implement ScaleSwap in the Linux kernel and evaluate its performance on a 128-core machine with an all-flash swap array comprising eight NVMe SSDs. Our evaluation shows that ScaleSwap achieves up to 3.4× higher throughput and up to 11.5× lower average latency than the Linux swap system. Furthermore, ScaleSwap outperforms two prior systems, TMO and ExtMEM, by up to 64% and 5×, respectively.
Modern SSDs demand faster I/O completion methods. While polling is a potential alternative to interrupts, it suffers under CPU contention. Hybrid polling mitigates this by sleeping early and polling later, yet it cannot keep up with rapidly varying I/O latencies and incurs context-switch overheads. We introduce PAS, an accurate latency tracking method for hybrid polling that adjusts sleep duration using the two most recent I/Os, and DPAS, which dynamically switches among polling, interrupts, and PAS to overcome the inherent drawbacks of hybrid polling. Experiments show that PAS reduces CPU usage by 21 percentage points compared to Linux hybrid polling for 4 KB random reads, and DPAS improves YCSB performance by 9% on a 3D XPoint SSD and 5% on a TLC NAND SSD, even under simultaneous CPU contention and I/O interference.
Efficient cloud computing relies on high-performance Elastic Block Storage (EBS) services, where virtual disk (VD) image loading significantly affects user experience. While the commonly used "lazy loading" approach reduces cold-start time from minutes to sub-seconds, our trace analysis of around 160,000 real-world image loading events in Alibaba EBS reveals that slow I/Os during initial block access constitute the primary performance bottleneck, accounting for 40% of all slow I/Os. We propose ThinkAhead, a data-driven image preloading system for VDs. ThinkAhead comprises various techniques to predict efficient block preloading sequences for images with historical traces based on runtime conditions and address corner cases with limited or no historical traces. Trace-driven simulation and cluster experiments show that ThinkAhead improves the data block hit rate by up to 7.27× and reduces tail waiting time by up to 98.7% across various image types compared to lazy loading and various baselines.
The running of applications in containers has emerged as a popular trend in the industry. The cold start of a container involves a sequential time-consuming process of image downloading and image unpacking. The high cold-start latency significantly prolongs the startup time of containerized applications and could potentially violate responsiveness SLAs in serverless computing or during service automatic scaling to handle burst requests. To accelerate container startup, state-of-the-art systems pull the container image on demand. Unfortunately, they suffer from userspace I/O interposition overhead, maintainability, and/or performance fluctuation. This paper presents CoFS, a novel filesystem based on extended FUSE for fast container startup. The insight is that the container image is built only once with a fixed read-only filesystem tree from the perspective of containers. This motivates CoFS to construct a minimal perfect hash function (MPHF) at image building time and to store metadata of files in a container image in a dense array indexed by their hash value. MPHF is collision-free and space-optimal. Leveraging the excellent properties of MPHF, CoFS accomplishes lookup request through less than one single I/O operation in most cases (unless the filename is excessively long) from kernel space, effectively avoiding the costly userspace lookup process in FUSE. Furthermore, CoFS constructs another MPHF that enables parallel lookup based on full path hashing, so as to further accelerate the path resolution. For data access, CoFS leverages sparse files provided by the in-kernel host filesystem to implement fine-grained data caching, and accesses cached data from kernel space. The evaluation shows that CoFS outperforms state-of-the-art systems that achieve fast container startup, and compared to fuse-loopback, a FUSE-based loopback filesystem, the lookup performance improves by up to 86%.