EDBT 2026 Demo / reviewers in the wild / expert
Xusheng Chen
dblp:38/7674
· DBLP profile ↗
33ranked-venue papers
2as first author
25since 2021 · last 2026
—ORCID · conflict
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 20 · 2 first-author · 15 since 2021Software engineering, systems software and programming languages · 6 · 6 since 2021Security and privacy · 5 · 1 since 2021Databases, data management, data science and information retrieval · 3 · 3 since 2021Computer networks · 2 · 1 since 2021Artificial intelligence and machine learning · 1 · 1 since 2021
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | Towards High-Goodput LLM Serving with Prefill-decode Multiplexing
Yukang Chen, Weihao Cui, Han Zhao 0005, Xiaoze Fan, Xusheng Chen, Yangjie Zhou 0001, Shixuan Sun, Bingsheng He, Quan Chen 0002 |
ASPLOS (2) | 6 |
| 2026 | Lessons Learned from Incorporating Formal Methods in Huawei Cloud ReliabilityabstractFormal methods are increasingly adopted in systems where reliability and correctness are critical, enabled by improvements in tool usability, speed, and automation. This industrial experience report presents three projects at Huawei Cloud showcasing different trade-offs in investment and assurance levels. We applied probabilistic concurrency testing, model checking, and deductive verification to two foundational services in the database and networking domains: the K2 transactional key-value store and the Global Server Load Balancer (GSLB). Claudia Cauli, Timo Lang, Sebti Mouelhi, Subhajit Bandopadhyay, Xusheng Chen, Yazhi Feng, Haoze Song, Linhua Tang, Zhenli Sheng, Ananth Shrinivas Srinath |
EuroSys | 7 |
| 2026 | Hermes: Efficient Serving of LLM Applications with Probabilistic Demand ModelingabstractApplications based on Large Language Models (LLMs) contain a series of tasks to address real-world problems with boosted capability, which have dynamic demand volumes on diverse backends. Existing serving systems treat the resource demands of LLM applications as a blackbox, compromising end-to-end efficiency due to improper queuing order and backend warm up latency. We find that the resource demands of LLM applications can be modeled in a general and accurate manner with Probabilistic Demand Graph (PDGraph). We then propose Hermes, which leverages PDGraph for efficient serving of LLM applications. Confronting probabilistic demand description, Hermes applies the Gittins policy to determine the scheduling order that can minimize the average application completion time. It also uses the PDGraph model to help prewarm cold backends at proper moments. Experiments with diverse LLM applications confirm that Hermes can effectively improve the application serving efficiency, reducing the average completion time by over 70% and the P95 completion time by over 80%. Zuo Gan, Zhenghao Gan, Chen Chen 0067, Yizhou Shan, Xusheng Chen, Zhenhua Han, Yifei Zhu 0001, Shixuan Sun, Minyi Guo |
ACM Trans. Archit. Code Optim. | 7 |
| 2025 | EPIC: Efficient Position-Independent Caching for Serving Large Language ModelsabstractLarge Language Models (LLMs) show great capabilities in a wide range of applications, but serving them efficiently becomes increasingly challenging as requests (prompts) become more complex. Context caching improves serving performance by reusing Key-Value (KV) vectors, the intermediate representations of tokens that are repeated across requests. However, existing context caching requires exact prefix matches across requests, limiting reuse cases in settings such as few-shot learning and retrieval-augmented generation, where immutable content (e.g., documents) remains unchanged across requests but is preceded by varying prefixes. Position-Independent Caching (PIC) addresses this issue by enabling modular reuse of the KV vectors regardless of prefixes. We formalize PIC and advance prior work by introducing EPIC, a serving system incorporating our new LegoLink algorithm, which mitigates the inappropriate “attention sink” effect at every document beginning, to maintain accuracy with minimal computation. Experiments show that EPIC achieves up to 8$\times$ improvements in Time-To-First-Token (TTFT) and 7$\times$ throughput gains over existing systems, with negligible or no accuracy loss. Wenrui Huang, Haoyi Wang, Tiancheng Hu, Xusheng Chen, Yizhou Shan, Tao Xie 0001 |
ICML | 8 |
| 2025 | Toppings: CPU-Assisted, Rank-Aware Adapter Serving for LLM Inference
Suyi Li 0002, Hanfeng Lu, Tianyuan Wu, Minchen Yu, Qizhen Weng 0001, Xusheng Chen, Yizhou Shan, Binhang Yuan, Wei Wang 0030 |
USENIX ATC | 6 |
| 2025 | DEEPSERVE: Serverless Large Language Model Serving at Scale
Zhixia Liu, Yuetao Chen, Baoquan Zhang, Shining Wan, Gengyuan Dan, Zhiyu Dong, Zhihao Ren, Changhong Liu, Tao Xie 0001, Dayun Lin, Xusheng Chen, Yizhou Shan |
USENIX ATC | 20 |
| 2025 | Perseus: Achieving Strong Consistency and High Data Freshness for Scalable Geo-distributed HTAPabstractThe rise of global data-driven applications has made geo-distributed hybrid transactional and analytical processing (HTAP) databases increasingly desirable. Existing distributed HTAP systems provide users with good performance on both transactions and analytical queries, and this good performance is scalable across a large number of data nodes. Unfortunately, these systems either provide weak consistency or incur bad data freshness when deployed geographically. In this paper, we present P erseus , a scalable HTAP database that enforces strong consistency for both transactions and analytical queries. To handle consistency efficiently, P erseus augments the classical dependency graph in concurrency control protocols to explicitly record the versions of data and their complete dependencies, implying which data needs to be read together in a snapshot. To minimize data staleness on analytical queries (another important goal of HTAP), P erseus further introduces a new dynamic snapshot algorithm that chooses updates selectively. Extensive evaluation results show that, compared to the HTAP databases with even weaker consistency, P erseus achieves up to 90% lower visibility delay, a metric of data freshness, capturing the time interval during which transactional updates are committed to the database and can be visible to analytical queries. Besides, Perseus is scalable across many nodes and robust to network instability. Haoze Song, Xusheng Chen, Ruijie Gong, Zekai Sun, Tianxiang Shen, Cheng Li 0001, Sen Wang 0004, Heming Cui |
Proc. ACM Manag. Data | 2 |
| 2025 | K2: On Optimizing Distributed Transactions in a Multi-region Data Store with True-time ClocksabstractTrueTime clocks (TTCs) that offer accurate and reliable time within limited uncertainty bounds have been increasingly implemented in many clouds. Multi-region data stores that seek decentralized synchronization for high performance represent an ideal application of TTC. However, the co-designs between the two often failed to realize their full potential. This paper proposes K2, a multi-region data store that explores the opportunity of using TTC for distributed transactions. Compared to its pioneer, Google Spanner, K2 augments TTC's semantics in three core design pillars. First, K2 carries a new timestamp-generating scheme that is capable of providing a small time uncertainty bound at scale. Second, K2 revitalizes existing multi-version timestamp-ordered concurrency control to realize multi-version properties for read-write transactions. Third, K2 introduces a new TTC-based visibility control protocol that provides efficient reads at replicas. Our evaluation shows that, K2 achieves an order of magnitude higher transaction throughput relative to other geo-distributed transaction protocols while ensuring a lower visibility delay at asynchronous replicas. Haoze Song, Xusheng Chen, Yazhi Feng, Xieyun Fang, Heming Cui, Linghe Kong |
Proc. VLDB Endow. | 3 |
| 2025 | FlatStor: An Efficient Embedded-Index Based Columnar Data Layout for Multimodal Data Workloads
Chi Zhang 0005, Yunfei Gu, Chentao Wu, Jie Li 0002, Xusheng Chen |
Proc. VLDB Endow. | 7 |
| 2025 | ShuffleInfer: Disaggregate LLM Inference for Mixed Downstream WorkloadsabstractTransformer-based large language model (LLM) inference serving is now the backbone of many cloud services. LLM inference consists of a prefill phase and a decode phase. However, existing LLM deployment practices often overlook the distinct characteristics of these phases, leading to significant interference. To mitigate interference, our insight is to carefully schedule and group inference requests based on their characteristics. We realize this idea in ShuffleInfer through three pillars. First, it partitions prompts into fixed-size chunks so that the accelerator always runs close to its computation-saturated limit. Second, it disaggregates prefill and decode instances so each can run independently. Finally, it uses a smart two-level scheduling algorithm augmented with predicted resource usage to avoid decode scheduling hotspots. Results show that ShuffleInfer improves time-to-first-token (TTFT), job completion time (JCT), and inference efficiency in terms of performance per dollar by a large margin, e.g., it uses 38% less resources all the while lowering average TTFT and average JCT by 97% and 47%, respectively. Cunchen Hu, Heyang Huang, Liangliang Xu, Xusheng Chen, Chenxi Wang 0005, Sa Wang, Yungang Bao, Ninghui Sun, Yizhou Shan |
ACM Trans. Archit. Code Optim. | 4 |
| 2025 | PipeMesh: Achieving Memory-Efficient Computation-Communication Overlap for Training Large Language ModelsabstractEfficiently training large language models (LLMs) on commodity cloud resources remains challenging due to limitations in network bandwidth and accelerator memory capacity. Existing training systems can be categorized based on their pipeline schedules. Depth-first scheduling, employed by systems like Megatron, prioritizes memory efficiency but restricts the overlap between communication and computation, causing accelerators to remain idle for over 20% of the training time. Conversely, breadth-first scheduling maximizes communication overlap but generates excessive intermediate activations, exceeding memory capacity and slowing computation by more than 34%. To address these limitations, we propose a novel elastic pipeline schedule that enables fine-grained control over the trade-off between communication overlap and memory consumption. Our approach determines the number of micro-batches scheduled together according to the communication time and the memory available. Furthermore, we introduce a mixed sharding strategy and a pipeline-aware selective recomputation technique to reduce memory usage. Experimental results demonstrate that our system eliminates most of the 28% all-accelerator idle time caused by communication, with recomputation accounting for less than 1.9% of the training time. Compared to existing baselines,PIPEMESHimproves training throughput on commodity clouds by 20.1% to 33.8%. Fanxin Li, Shixiong Zhao, Yuhao Qing, Jianyu Jiang, Xusheng Chen, Heming Cui |
IEEE Trans. Parallel Distributed Syst. | 5 |
| 2025 | Slarm: SLA-Aware, Reliable and Efficient Transaction Dissemination for Permissioned BlockchainsabstractThe blockchain paradigm has attracted diverse applications to be deployed upon. However, no service-level agreement (SLA) mechanism has been proposed to enforce the SLA disseminating deadlines to commit blockchain transactions, although these transactions are often interactively submitted by clients and desire short SLA deadlines (e.g., tens of seconds). Existing peer-to-peer (P2P) multicast protocols for blockchains take the unidirectional approach to disseminate transactions regardless of their SLA deadlines, making transactions easily violate their deadlines. Moreover, these protocols are vulnerable to malicious P2P nodes, and their protocol messages (e.g., SLAstringent transactions) are vulnerable to deferring attacks. We propose SLARM, the first bidirectional P2P multicast protocol for permissioned blockchains, which conservatively adjusts transactions' dissemination speed to satisfy their SLA deadlines according to the trustworthy SLA feedback of previously disseminated transactions. SLARM guarantees transactions' SLAs in a decentralized way and defends against the deferring attacks using TEE. Evaluation of SLARM with five notable P2P multicast protocols and five diverse real-world applications shows that: even with transaction spikes and attacked nodes, SLARM achieves a much higher transaction SLA satisfaction rate with reasonably high commit throughput. Ji Qi 0002, Tianxiang Shen, Jianyu Jiang, Xusheng Chen, Xiapu Luo, Fengwei Zhang, Heming Cui |
IEEE Trans. Serv. Comput. | 4 |
| 2023 | Skadi: Building a Distributed Runtime for Data Systems in Disaggregated Data CentersabstractData-intensive systems are the backbone of today's computing and are responsible for shaping data centers. Over the years, cloud providers have relied on three principles to maintain cost-effective data systems: use disaggregation to decouple scaling, use domain-specific computing to battle waning laws, and use serverless to lower costs. Although they work well individually, they fail to work in harmony: an issue amplified by emerging data system workloads. Cunchen Hu, Chenxi Wang 0005, Sa Wang, Ninghui Sun, Yungang Bao, Jieru Zhao, Sanidhya Kashyap, Pengfei Zuo, Xusheng Chen, Liangliang Xu, Yizhou Shan |
HotOS | 9 |
| 2023 | Coorp: Satisfying Low-Latency and High-Throughput Requirements of Wireless Network for Coordinated Robotic LearningabstractIn coordinated robotic learning, multiple robots share the same wireless channel for communication, and bring together latency-sensitive (LS) network flows for control and bandwidth-hungry (BH) flows for distributed learning. Unfortunately, existing wireless network supporting systems cannot coordinate these two network flows to meet their own requirements: 1) prioritized contention systems (e.g., EDCA) prevent LS messages from timely acquiring the wireless channel because multiple wireless network interface cards (WNICs) with BH messages are contending for the channel 2) global planning systems (e.g., SchedWiFi) have to reserve a notable time window in the shared channel for each LS flow, suffering from severe bandwidth degradation (up to 42%). We present the coordinated preemption method to meet both requirements for LS flows and BH flows. Globally (among multiple robots), coordinated preemption eliminates unnecessary contention of BH flows by making them transmit in a round-robin manner, such that LS flows have the highest chance to win the contention against BH flows, without sacrificing overall bandwidth from the perspective of coordinated robotic learning applications. Locally (within the same robot), coordinated preemption in real time predicts the periodic transmission of LS flows from the upper application and conservatively limits packets of BH flows buffered in the WNIC only before LS packets arriving, reducing the bandwidth devoted to preemption. COORP, our implementation of coordinated preemption, reduced the violation of latency requirements from 53.9% (EDCA) to 8.8% (comparable to SchedWiFi). Regarding learning quality, COORP achieved a comparable (at times the same) learning reward with EDCA, which grew up to 76% faster than SchedWiFi. Shengliang Deng, Xiuxian Guan, Zekai Sun, Shixiong Zhao, Tianxiang Shen, Xusheng Chen, Tianyang Duan, Jia Pan 0001, Libo Zhang 0001, Heming Cui |
IEEE Internet Things J. | 6 |
| 2023 | Fold3D: Rethinking and Parallelizing Computational and Communicational Tasks in the Training of Large DNN ModelsabstractTraining a large DNN (e.g., GPT3) efficiently on commodity clouds is challenging even with the latest 3D parallel training systems (e.g., Megatron v3.0). In particular, along the pipeline parallelism dimension, computational tasks that produce a whole DNN's gradients with multiple input batches should be concurrently activated; along the data parallelism dimension, a set of heavy-weight communications (for aggregating the accumulated outputs of computational tasks) isinevitably serializedafter the pipelined tasks, undermining the training performance (e.g., in Megatron, data parallelism caused all GPUs idle for over 44% of the training time) over commodity cloud networks. To deserialize these communicational and computational tasks, we propose the AIAO scheduling (for 3D parallelism) which slices a DNN into multiple segments, so that the computational tasks processing the same DNN segment can be scheduled together, and the communicational tasks that synchronize this segment can be launched and overlapped (deserialized) with other segments’ computational tasks. We realized this idea in ourFold3Dtraining system. Extensive evaluation showsFold3Deliminated most of the all-GPU 44% idle time in Megatron (caused by data parallelism), leading to 25.2%–42.1% training throughput improvement compared to four notable baselines over various settings;Fold3D's high performance scaled to many GPUs. Fanxin Li, Shixiong Zhao, Yuhao Qing, Xusheng Chen, Xiuxian Guan, Sen Wang 0004, Gong Zhang 0001, Heming Cui |
IEEE Trans. Parallel Distributed Syst. | 4 |
| 2023 | A Geography-Based P2P Overlay Network for Fast and Robust Blockchain SystemsabstractNumerous blockchain systems with various consensus protocols have emerged to achieve high transaction rates (2$\sim$10K tps). However, their underlying P2P network primitives constrain further improvements due to two problems (i) high message redundancy and (ii) long broadcast convergence time. The first problem is caused by the excessive robustness of the dominant broadcast approach Gossip. All state-of-the-art blockchain systems only tolerate 20-50% node failure while Gossip can withstand up to 90%. The reason for (ii) is that existing broadcast topologies ignore geographical distances among nodes and incur paths with unnecessarily high latency. We presentFRing, a geography-based P2P overlay network for fast and robust broadcast in blockchain systems.FRinghas three main features: sufficient robustness, low message redundancy, and fast convergence. To reduce convergence time,FRingforms the network topology by considering geographical proximity. A novel broadcast algorithm based onFRingtopology is proposed to lower message redundancy while maintaining sufficient robustness. One major challenge is to eliminate the risk of topology inference by traffic pattern analysis.FRingleverages Intel SGX to guarantee nodes’ behavior integrity and incorporates pattern obfuscation to prevent traffic pattern analysis. The evaluation shows thatFRingimproved the throughput of EOS by 220% and Hyperledger Fabric by 210%. Haoran Qiu, Shixiong Zhao, Xusheng Chen, Ji Qi 0002, Heming Cui, Sen Wang 0004 |
IEEE Trans. Serv. Comput. | 4 |
| 2022 | NASPipe: high performance and reproducible pipeline parallel supernet training via causal synchronous parallelismabstractSupernet training, a prevalent and important paradigm in Neural Architecture Search, embeds the whole DNN architecture search space into one monolithic supernet, iteratively activates a subset of the supernet (i.e., a subnet) for fitting each batch of data, and searches a high-quality subnet which meets specific requirements. Although training subnets in parallel on multiple GPUs is desirable for acceleration, there inherently exists a race hazard that concurrent subnets may access the same DNN layers. Existing systems support neither efficiently parallelizing subnets’ training executions, nor resolving the race hazard deterministically, leading to unreproducible training procedures and potentiallly non-trivial accuracy loss. Shixiong Zhao, Fanxin Li, Xusheng Chen, Tianxiang Shen, Li Chen 0008, Sen Wang 0004, Nicholas Zhang, Cheng Li 0001, Heming Cui |
ASPLOS | 3 |
| 2022 | ROG: A High Performance and Robust Distributed Training System for Robotic IoTabstractCritical robotic tasks such as rescue and disaster response are more prevalently leveraging ML (Machine Learning) models deployed on a team of wireless robots, on which data parallel (DP) training over Internet of Things of these robots (robotic IoT) can harness the distributed hardware resources to adapt their models to changing environments as soon as possible. Unfortunately, due to the need for DP synchronization across all robots, the instability in wireless networks (i.e., fluctuating bandwidth due to occlusion and varying communication distance) often leads to severe stall of robots, which affects the training accuracy within a tight time budget and wastes energy stalling. Existing methods to cope with the instability of datacenter networks are incapable of handling such straggler effect. That is because they are conducting model-granulated transmission scheduling, which is much more coarse-grained than the granularity of transient network instability in real-world robotic IoT networks, making a previously reached schedule mismatch with the varying bandwidth during transmission. We present ROG, the first ROw-Granulated distributed training system optimized for ML training over unstable wireless networks. ROG confines the granularity of transmission and synchronization to each row of a layer’s parameters and schedules the transmission of each row adaptively to the fluctuating bandwidth. In this way the ML training process can update partial and the most important gradients of a stale robot to avoid triggering stalls, while provably guaranteeing convergence. The evaluation shows that, given the same training time, ROG achieved about 4.9%~6.5% training accuracy gain compared with the baselines and saved 20.4%~50.7% of the energy to achieve the same training accuracy. Xiuxian Guan, Zekai Sun, Shengliang Deng, Xusheng Chen, Shixiong Zhao, Zongyuan Zhang, Tianyang Duan, Chenshu Wu, Yong Cui 0001, Libo Zhang 0001, Rui Wang 0007, Heming Cui |
MICRO | 4 |
| 2022 | CRONUS: Fault-isolated, Secure and High-performance Heterogeneous Computing for Trusted Execution EnvironmentabstractWith the trend of processing a large volume of sensitive data on PaaS services (e.g., DNN training), a TEE architecture that supports general heterogeneous accelerators, enables spatial sharing on one accelerator, and enforces strong isolation across accelerators is highly desirable. However, none of the existing TEE solutions meet all three requirements. In this paper, we propose CRONUS, the first TEE architecture that achieves the three crucial requirements. The key idea of CRONUS is to partition heterogeneous computation into isolated TEE enclaves, where each enclave encapsulates only one kind of computation (e.g., GPU computation), and multiple enclaves can spatially share an accelerator. Then, CRONUS constructs heterogeneous computing using remote procedure calls (RPCs) among enclaves. With CRONUS, each accelerator’s hardware and its software stack are strongly isolated from others’, and each enclave trusts only its own hardware. To tackle the security challenge caused by inter-enclave interactions, we design a new streaming remote procedure call abstraction to enable secure RPCs with high performance. CRONUS is software-based, making it general to diverse accelerators. We implemented CRONUS on ARM TrustZone. Evaluation on diverse workloads with CPUs, GPUs and NPUs shows that, CRONUS achieves less than 7.1% extra computation time compared to native (unprotected) executions. Jianyu Jiang, Ji Qi 0002, Tianxiang Shen, Xusheng Chen, Shixiong Zhao, Sen Wang 0004, Li Chen 0008, Gong Zhang 0001, Xiapu Luo, Heming Cui |
MICRO | 4 |
| 2022 | SOTER: Guarding Black-box Inference for General Neural Networks at the Edge
Tianxiang Shen, Ji Qi 0002, Jianyu Jiang, Siyuan Wen, Xusheng Chen, Shixiong Zhao, Sen Wang 0004, Li Chen 0008, Xiapu Luo, Fengwei Zhang, Heming Cui |
USENIX ATC | 6 |
| 2022 | Efficient and DoS-resistant Consensus for Permissioned Blockchains
Xusheng Chen, Shixiong Zhao, Ji Qi 0002, Jianyu Jiang, Haoze Song, Cheng Wang 0021, Tsz On Li, T.-H. Hubert Chan, Fengwei Zhang, Xiapu Luo, Sen Wang 0004, Gong Zhang 0001, Heming Cui |
Perform. Evaluation | 1 |
| 2022 | DAENet: Making Strong Anonymity Scale in a Fully Decentralized NetworkabstractTraditional anonymous networks (e.g., Tor) are vulnerable to traffic analysis attacks that monitor the whole network traffic to determine which users are communicating. To preserve user anonymity against traffic analysis attacks, the emerging mix networks mess up the order of packets through a set of centralized and explicit shuffling nodes. However, this centralized design of mix networks is insecure against targeted DoS attacks that can completely block these shuffling nodes. In this article, we presentDAENet, an efficient mix network that resists both targeted DoS attacks and traffic analysis attacks with a new abstraction calledStealthy Peer-to-Peer (P2P) Network. Thestealthy P2P networkeffectively hides the shuffling nodes used in a routing path into the whole network, such that adversaries cannot distinguish specific shuffling nodes and conduct targeted DoS attacks to block these nodes. In addition, to handle traffic analysis attacks, we leverage the confidentiality and integrity protection of Intel SGX to ensure trustworthy packet shuffles at each distributed host and use multiple routing paths to prevent adversaries from tracking and revealing user identities. We show that our system is scalable with moderate latency (2.2s) when running in a cluster of 10,000 participants and is robust in the case of machine failures, making it an attractive new design for decentralized anonymous communication. DAENet ’s code is released onhttps://github.com/hku-systems/DAENet. Tianxiang Shen, Jianyu Jiang, Yunpeng Jiang, Xusheng Chen, Ji Qi 0002, Shixiong Zhao, Fengwei Zhang, Xiapu Luo, Heming Cui |
IEEE Trans. Dependable Secur. Comput. | 4 |
| 2022 | vPipe: A Virtualized Acceleration System for Achieving Efficient and Scalable Pipeline Parallel DNN TrainingabstractThe increasing computational complexity of DNNs achieved unprecedented successes in various areas such as machine vision and natural language processing (NLP), e.g., the recent advanced Transformer has billions of parameters. However, as large-scale DNNs significantly exceed GPU's physical memory limit, they cannot be trained by conventional methods such as data parallelism. Pipeline parallelism that partitions a large DNN into small subnets and trains them on different GPUs is a plausible solution. Unfortunately, the layer partitioning and memory management in existing pipeline parallel systems are fixed during training, making them easily impeded by out-of-memory errors and the GPU under-utilization. These drawbacks amplify when performing neural architecture search (NAS) such as the evolved Transformer, where different network architectures of Transformer needed to be trained repeatedly. vPipe is the first system that transparently provides dynamic layer partitioning and memory management for pipeline parallelism. vPipe has two unique contributions, including (1) an online algorithm for searching a near-optimal layer partitioning and memory management plan, and (2) a live layer migration protocol for re-balancing the layer distribution across a training pipeline. vPipe improved the training throughput of two notable baselines (Pipedream and GPipe) by 61.4-463.4 percent and 24.8-291.3 percent on various large DNNs and training settings. Shixiong Zhao, Fanxin Li, Xusheng Chen, Xiuxian Guan, Jianyu Jiang, Dong Huang 0005, Yuhao Qing, Sen Wang 0004, Peng Wang 0037, Gong Zhang 0001, Cheng Li 0001, Ping Luo 0002, Heming Cui |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2021 | Achieving low tail-latency and high scalability for serializable transactions in edge computingabstractA distributed database utilizing the wide-spread edge computing servers to provide low-latency data access with the serializability guarantee is highly desirable for emerging edge computing applications. In an edge database, nodes are divided into regions, and a transaction can be categorized as intra-region (IRT) or cross-region (CRT) based on whether it accesses data in different regions. In addition to serializability, we insist that a practical edge database should provide low tail latency for both IRTs and CRTs, and such low latency must be scalable to a large number of regions. Unfortunately, none of existing geo-replicated serializable databases or edge databases can meet such requirements. Xusheng Chen, Haoze Song, Jianyu Jiang, Chaoyi Ruan, Cheng Li 0001, Sen Wang 0004, Gong Zhang 0001, Reynold Cheng, Heming Cui |
EuroSys | 1 |
| 2021 | Bidl: A High-throughput, Low-latency Permissioned Blockchain Framework for Datacenter NetworksabstractA permissioned blockchain framework typically runs an efficient Byzantine consensus protocol and is attractive to deploy fast trading applications among a large number of mutually untrusted participants (e.g., companies). Unfortunately, all existing permissioned blockchain frameworks adopt sequential workflows for invoking the consensus protocol and executing applications' transactions, making the performance of these applications much lower than deploying them in traditional systems (e.g., in-datacenter stock exchange). Ji Qi 0002, Xusheng Chen, Yunpeng Jiang, Jianyu Jiang, Tianxiang Shen, Shixiong Zhao, Sen Wang 0004, Gong Zhang 0001, Li Chen 0008, Man Ho Au, Heming Cui |
SOSP | 2 |
| 2020 | Uranus: Simple, Efficient SGX Programming and its ApplicationsabstractApplications written in Java have strengths to tackle diverse threats in public clouds, but these applications are still prone to privileged attacks when processing plaintext data. Intel SGX is powerful to tackle these attacks, and traditional SGX systems rewrite a Java application's sensitive functions, which process plaintext data, using C/C++ SGX API. Although this code-rewrite approach achieves good efficiency and a small TCB, it requires SGX expert knowledge and can be tedious and error-prone. To tackle the limitations of rewriting Java to C/C++, recent SGX systems propose a code-reuse approach, which runs a default JVM in an SGX enclave to execute the sensitive Java functions. However, both recent study and this paper find that running a default JVM in enclaves incurs two major vulnerabilities, Iago attacks, and control flow leakage of sensitive functions, due to the usage of OS features in JVM. In this paper, Uranus creates easy-to-use Java programming abstractions for application developers to annotate sensitive functions, and Uranus automatically runs these functions in SGX at runtime. Uranus effectively tackles the two major vulnerabilities in the code-reuse approach by presenting two new protocols: 1) a Java bytecode attestation protocol for dynamically loaded functions; and 2) an OS-decoupled, efficient GC protocol optimized for data-handling applications running in enclaves. We implemented Uranus in Linux and applied it to two diverse data-handling applications: Spark and ZooKeeper. Evaluation shows that: 1) Uranus achieves the same security guarantees as two relevant SGX systems for these two applications with only a few annotations; 2) Uranus has reasonable performance overhead compared to the native, insecure applications; and 3) Uranus defends against privileged attacks. Uranus source code and evaluation results are released on https://github.com/hku-systems/uranus. Jianyu Jiang, Xusheng Chen, Tsz On Li, Cheng Wang 0021, Tianxiang Shen, Shixiong Zhao, Heming Cui, Cho-Li Wang, Fengwei Zhang |
AsiaCCS | 2 |
| 2020 | UPA: An Automated, Accurate and Efficient Differentially Private Big-Data Mining SystemabstractIn the era of big-data, individuals and institutions store their sensitive data on clouds, and these data are often analyzed and computed by MapReduce frameworks (e.g., Spark). However, releasing the computation result on these data may leak privacy. Differential Privacy (DP) is a powerful method to preserve the privacy of an individual data record from a computation result. Given an input dataset and a query, DP typically perturbs an output value with noise proportional to sensitivity, the greatest change on an output value when a record is added to or removed from the input dataset. Unfortunately, directly computing the sensitivity value for a query and an input dataset is computationally infeasible, because it requires adding or removing every record from the dataset and repeatedly running the same query on the dataset: a dataset of one million input records requires running the same query for more than one million times. This paper presents UPA, the first automated, accurate, and efficient sensitivity inferring approach for big-data mining applications. Our key observation is that MapReduce operators often have commutative and associative properties in order to enable parallelism and fault tolerance among computers. Therefore, UPA can greatly reduce the repeated computations at runtime while computing a precise sensitivity value automatically for general big-data queries. We compared UPA with FLEX, the most relevant work that does static analysis on queries to infer sensitivity values. Based on an extensive evaluation on nine diverse Spark queries, UPA supports all the nine evaluated queries, while FLEX supports only five of the nine queries. For the five queries which both UPA and FLEX can support, UPA enforces DP with five orders of magnitude more accurate sensitivity values than FLEX. UPA has reasonable performance overhead compared to native Spark. UPA's source code is available on https://github.com/hku-systems/UPA. Tsz On Li, Jianyu Jiang, Ji Qi 0002, Chi Chiu So, Jiacheng Ma 0003, Xusheng Chen, Tianxiang Shen, Heming Cui, Peng Wang 0070 |
DSN | 6 |
| 2020 | HAMS: High Availability for Distributed Machine Learning Service GraphsabstractMission-critical services often deploy multiple Machine Learning (ML) models in a distributed graph manner, where each model can be deployed on a distinct physical host. Practical fault tolerance for such ML service graphs should meet three crucial requirements: high availability (fast failover), low normal case performance overhead, and global consistency under non-determinism (e.g., threads in a GPU can do floating point additions in random order). Unfortunately, despite much effort, existing fault tolerance systems, including those taking the primary-backup approach or the checkpoint-replay approach, cannot meet all these three requirements. To tackle this problem, we present HAMS, which starts from the primary-backup approach to replicate each stateful ML model, and we leverage the causal logging technique from the checkpoint-replay approach to eliminate the notorious stop-and-buffer delay in the primary-backup approach. Extensive evaluation on 25 ML models and six ML services shows that: (1) in normal case, HAMS achieved 0.5%-3.7% overhead on latency compared with bare metal; (2) HAMS took 116.12ms-254.19ms to recover one stateful model in all services, 155.1X-1067.9X faster than a relevant system Lineage Stash (LS); and (3) HAMS recovered these services with global consistency even when the GPU non-determinism exists, not supported by LS. HAMS's code is released ongithub.com/hku-systems/hams. Shixiong Zhao, Xusheng Chen, Cheng Wang 0021, Fanxin Li, Heming Cui, Cheng Li 0001, Sen Wang 0004 |
DSN | 2 |
| 2019 | Fulva: Efficient Live Migration for In-Memory Key-Value Stores with Zero DowntimeabstractA key-value store live migration approach migrates key-value tuples and their client requests from an overloaded machine (source) to an idle machine (destination), while still serving client requests. Existing migration approaches fall into two categories. First, a source-driven approach (e.g., DrTM-B) executes all client requests on the source and incrementally propagates the updated key-value tuples to the destination. This approach has an inevitable downtime to completely propagate the updated tuples at the end of a migration. Second, a destination-driven approach (e.g., RockSteady) executes all read and write requests on the destination, and pulls tuples from source for read requests on-demand. This approach has zero downtime, but incurs extra network round-trips due to the on-demand pull, greatly increasing request latency. Overall, a live migration approach that has zero downtime and no performance degradation during the migration is highly desirable but missing. The key observation of our Fulva system is that the source and destination can cooperatively drive the migration and serve requests, and we need only to design an efficient protocol to ensure linearizability (i.e., reads see the updates from the latest writes). To this end, when a migration starts, Fulva works by three steps. First, all write requests are redirected to the destination. Second, each client program uses a Fulva's RPC library to track the migration progress. For read requests accessing the already-migrated tuples, Fulva RPC library sends the requests to the destination. Third, for read requests accessing not-yet-migrated tuples, Fulva sends to both machines. The first step avoids downtime because all updated tuples are already on the destination. The second and third steps avoid the on-demand pull and ensure linearizability, making Fulva efficient. We implemented Fulva using DPDK and integrated it with RAMCloud, a popular in-memory key-value store. We compared Fulva with two notable systems, RockSteady (destination-driven approach) and RAMCloud's default source-driven approach. Extensive evaluation shows that Fulva had much higher throughput and lower latency than the two systems, and Fulva's network bandwidth usage is comparable with RockSteady. All Fulva's source code and raw evaluation results are released on github.com/hku-systems/fulva. Jiewen Hai, Cheng Wang 0021, Xusheng Chen, Tsz On Li, Heming Cui, Sen Wang 0004 |
SRDS | 3 |
| 2018 | PLOVER: Fast, Multi-core Scalable Virtual Machine Fault-tolerance
Cheng Wang 0021, Xusheng Chen, Weiwei Jia 0001, Boxuan Li, Haoran Qiu, Shixiong Zhao, Heming Cui |
NSDI | 2 |
| 2018 | Effectively Mitigating I/O Inactivity in vCPU Scheduling
Weiwei Jia 0001, Cheng Wang 0021, Xusheng Chen, Jianchen Shan, Xiaowei Shang, Heming Cui, Xiaoning Ding, Luwei Cheng, Francis C. M. Lau 0001, Yuangang Wang |
USENIX ATC | 3 |
| 2017 | APUS: fast and scalable paxos on RDMAabstractState machine replication (SMR) uses Paxos to enforce the same inputs for a program (e.g., Redis) replicated on a number of hosts, tolerating various types of failures. Unfortunately, traditional Paxos protocols incur prohibitive performance overhead on server programs due to their high consensus latency on TCP/IP. Worse, the consensus latency of extant Paxos protocols increases drastically when more concurrent client connections or hosts are added. This paper presents APUS, the first RDMA-based Paxos protocol that aims to be fast and scalable to client connections and hosts. APUS intercepts inbound socket calls of an unmodified server program, assigns a total order for all input requests, and uses fast RDMA primitives to replicate these requests concurrently. Cheng Wang 0021, Jianyu Jiang, Xusheng Chen, Ning Yi, Heming Cui |
SoCC | 3 |
| 2017 | A Fast, General Storage Replication Protocol for Active-Active Virtual Machine Fault ToleranceabstractCloud computing enables more and more online services deployed in virtual machines (VMs), making fast VM fault tolerance particularly crucial. Unfortunately, despite much effort, achieving fast VM fault tolerance remains an open problem. A traditional way to provide VM fault tolerance is the active-passive approach, which frequently transfers tremendous updated states, including memory and storage, of a primary VM to a suspended secondary VM. The other emerging approach, namely the active-active approach, runs the secondary VM concurrently with the primary. Compared to active-passive, active-active is faster because it only performs the transfer when the externally visible states (e.g., network outputs) of the primary and secondary diverge. However, active-active aggravates the performance issue on I/O intensive workloads. In existing active-active systems, storage replication protocols hold updated storage states from both the primary and secondary on the secondary, incurring excessive I/O contention. For instance, both our evaluation and prior study show that a well-engineered active-active system, COLO, degrades the throughput of I/O intensive services by up to 61.6%. To tackle this open problem, this paper presents GANNET, a fast and general storage replication protocol for active-active VM fault tolerance systems. It greatly alleviates the I/O contention on the secondary's storage by efficiently buffering the updated disk states from both the primary and secondary VM in memory. GANNET carries a lightweight storage checkpoint algorithm to avoid consuming too much memory. GANNET is proved to be as reliable as existing storage replication protocols. We integrated GANNET into two popular active-active systems. Evaluation on six widely used services shows that GANNET incurred 15.9% overhead compared to the native executions and outperformed COLO's storage replication protocol by 1.2X~2.6X. GANNET's source code is available at github.com/hku-systems/gannet. Cheng Wang 0021, Xusheng Chen, Youwei Zhu, Heming Cui |
ICPADS | 2 |