EDBT 2026 Demo / reviewers in the wild / expert
Chen Chen 0067
dblp:65/4423-67
· DBLP profile ↗
52ranked-venue papers
13as first author
42since 2021 · last 2026
0000-0001-9480-5632ORCID · conflict
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 32 · 8 first-author · 26 since 2021Computer networks · 12 · 5 first-author · 8 since 2021Software engineering, systems software and programming languages · 5 · 5 since 2021Artificial intelligence and machine learning · 3 · 3 since 2021Databases, data management, data science and information retrieval · 2 · 2 since 2021
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | Towards Efficient Serving of Network-intensive LLM InferencesabstractPrefix caching has become a key technique for LLM serving, and nowadays the reusable KVCache contents are often hosted on distributed servers. For long-context LLM inferences with high cache hit ratio, cross-server KVCache transmission has become an emerging performance bottleneck; such network-intensive LLM inferences are increasingly prevalent in the coming era of agentic AI. However, existing LLM inference engines are essentially compute-centric; we find that they are highly inefficient when serving such workloads due to compute-stage service blocking and ignorance of KVCache-transfer cost. Chen Chen 0067, Junxue Zhang 0001, Zhusheng Wang, Zixuan Guan, Qizhen Weng 0001, Minyi Guo |
APNet | 2 |
| 2026 | Suika: Efficient and High-quality Re-scheduling of 3D-parallelized LLM Training Jobs in Shared ClustersabstractLarge Language Models (LLMs) are usually trained with 3D (data, tensor, and pipeline) parallelism—in shared GPU clusters where the available resources are highly dynamic. Rescheduling the idle resources to ongoing jobs can help improve cluster utilization, but doing so for 3D-parallelized training jobs suffers large overheads in performance modeling, decision making, and redeployment. We present Suika, a cluster training system that supports efficient and high-quality resource rescheduling for 3D-parallelized LLM training jobs. Suika holistically addresses the complexity challenges by exploiting the incremental nature of rescheduling. For performance modeling, it builds an accurate performance estimator with non-disruptive online profiling. For decision-making, it employs topology-aware sorting and an expand-and-balance algorithm to reduce the complexity of resource allocation and job parallelization, without compromising decision quality. Suika further integrates a device-to-device redeployment method to leverage the overlapping nature of incremental reconfiguration for overhead reduction. Experiments on 64-GPU physical cluster and 1024-GPU simulated cluster show that, Suika achieves 1.29 ~ 1.31× reduction in average JCT compared to state-of-the-art schedulers. Chen Chen 0067, Chunyu Xue, Qizhen Weng 0001, Zeren Li, Xuqi Zhu, Yongqiang Yang, Quan Chen 0002, Minyi Guo |
EuroSys | 3 |
| 2026 | Arena: Efficiently Training Large Models via Dynamic Scheduling and Adaptive Parallelism Co-DesignabstractEfficiently training large-scale models (LMs) in GPU clusters involves two separate avenues: inter-job dynamic scheduling and intra-job adaptive parallelism (AP). However, existing dynamic schedulers struggle with large-model scheduling due to the mismatch between static parallelism (SP)-aware scheduling and AP-based execution, leading to cluster inefficiencies such as degraded throughput and prolonged job queuing. This paper presents Arena, a large-model training system that co-designs dynamic scheduling and adaptive parallelism to achieve high cluster efficiency. To reduce scheduling costs while improving decision quality, Arena designs low-cost, disaggregated profiling and AP-tailored, load-aware performance estimation, while unifying them by sharding the joint scheduling-parallelism optimization space via a grid abstraction. Building on this, Arena dynamically schedules profiled jobs in elasticity and heterogeneity dimensions, and executes them using efficient AP with pruned search space. Evaluated on heterogeneous testbeds and production workloads, Arena reduces job completion time by up to 49.3% and improves cluster throughput by up to 1.6X. Chunyu Xue, Weihao Cui, Quan Chen 0002, Chen Chen 0067, Han Zhao 0005, Shulai Zhang, Linmei Wang, Limin Xiao 0001, Weifeng Zhang 0003, Jing Yang 0017, Bingsheng He, Minyi Guo |
EuroSys | 4 |
| 2026 | FluxZK: Scalable and Efficient Zero-Knowledge Proof Computation via GPU AccelerationabstractZero-knowledge succinct non-interactive arguments of knowledge (zkSNARKs) are a key technology to privacy-preserving applications today. The complexity of proof generation, however, heavily constrains throughput in latency-sensitive environments. The computational burden primarily stems from two fundamental algorithms: Multi-Scalar Multiplication (MSM) and the Number Theoretic Transform (NTT). We propose a series of optimizations for these two kernels, including computation-transfer pipelining, load balancing, and memory access fusion, achieving 1.97 × to 2.16 × proof generation speedup over a state-of-the-art open source GPU acceleration library. Our design also supports out-of-core computation, enabling the generation of large-scale ZKP proofs. Xinwei Qiang, Liukun Yu, Zhengyi Li 0002, Shixuan Sun, Jingwen Leng, Chen Chen 0067, Jiaping Gui, Zhenzhe Zheng 0001, Jin Dong 0004, Minyi Guo |
HPDC | 7 |
| 2026 | RetroInfer: A Vector Storage Engine for Scalable Long-Context LLM Inference
Yaoqi Chen, Jinkai Zhang, Baotong Lu, Qianxi Zhang, Chengruidong Zhang, Jingjia Luo, Huiqiang Jiang, Qi Chen 0009, Bailu Ding, Xiao Yan 0002, Jiawei Jiang 0001, Chen Chen 0067, Cheng Li 0001, Yuqing Yang 0001, Fan Yang 0024, Mao Yang 0004 |
Proc. VLDB Endow. | 14 |
| 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. | 5 |
| 2026 | NPUMeter: Automatic Operator Optimization for Ascend NPU with Accurate Analytical Performance ModelsabstractWith the rapid development of AI and deep learning, computational demands are increasing significantly. While GPUs excel in parallel computing, they fall short in terms of energy efficiency, specialization, and processing latency. In contrast, Neural Processing Units (NPUs), such as the Ascend NPUs, designed specifically for deep learning tasks, demonstrate superior performance. However, the architecture specialization makes operator development more challenging, leading to a reliance on manual tuning and optimization, which incurs significant time cost and developing effort. To address this issue, we propose NPUMeter, an automatic operator optimization framework for Ascend NPUs built upon accurate and comprehensive analytical performance models. NPUMeter comprises two components: (1) an analytical performance model that accurately estimates operator latency on NPU given different configurations of optimization parameters; (2) an efficient design space exploration (DSE) algorithm that automatically searches for the optimal parameter configuration in a large design space within minutes. Experimental results demonstrate that NPUMeter achieves high estimation accuracy, with an average error below 5%. It effectively generates near-optimal configurations for various operators, achieving up to a 1.46× performance speedup compared to the configuration generated by the Ascend C compiler while reducing the DSE time from hours to minutes. Weichuang Zhang, Yufei Shangguan, Yuting Mai, Qiuliang Wang, Chen Chen 0067, Quan Chen 0002, Wenchao Ding 0001, Jieru Zhao, Minyi Guo |
ACM Trans. Archit. Code Optim. | 6 |
| 2026 | Castor: Optimizing Deep Learning Job Scheduling in Multi-Tenant GPU Clusters via Intelligent ColocationabstractDeep learning (DL) has achieved significant success across a wide range of domains, prompting the widespread deployment of GPU clusters equipped with specialized accelerators to support high-performance training workloads. To minimize operational costs while maximizing resource utilization, efficient job scheduling in these clusters is essential. Although recent schedulers have improved cluster efficiency through periodic reallocation or selection of GPU resources, they still face challenges such as preemption and migration overheads, along with the risk of degrading model accuracy. Despite these limitations, the potential of GPUsharing remains largely underexplored. Few existing studies have systematically examined GPU sharing as a strategy to enhance resource utilization and reduce job queuing delays in multi-tenant DL clusters. Motivated by these insights, we propose a job scheduling model that enables multiple jobs to share the same set of GPUs without modifying their original training configurations. We introduce Castor, a simple yet efficient scheduling system, to achieve intelligent GPU colocation for multiple DL jobs. Castor intelligently selects job pairs for GPU sharing and determines runtime parameters (sub-batch size and scheduling time point) to optimize overall system performance while preserving the accuracy of DL convergence through gradient accumulation. Through a combination of physical DL workloads and trace-driven simulations across various configurations, we demonstrate that Castor reduces average job completion time by 26–52% compared to state-of-the-art preemptive DL schedulers, despite operating under a preemption-free policy. Furthermore, Castor effectively identifies optimal resource-sharing configurations, outperforming the baseline first-fit sharing policy (SJF-FFS) by up to 20% on large-scale workload traces. Yizhou Luo, Jiaxin Lai, Shaohuai Shi, Chen Chen 0067, Shuhan Qi, Jiajia Zhang 0001, Qiang Wang 0022 |
IEEE Trans. Cloud Comput. | 4 |
| 2026 | Boosting Gradient-Based Training Diagnosis for Efficient and Accurate Federated LearningabstractFederated Learning (FL) allows edge clients to collaborate in model training with data privacy preserved, yet it is known to suffer low training efficiency and model accuracy. Given that efficiency and accuracy are usually conflicting objectives, existing practices increasingly employ an adaptive scheme that changes the FL configurations (e.g., quantization or sparsification level) based on runtime training status, for which accurate training diagnosis—used for guiding the optimization actions—is crucial. However, while training diagnosis is a common task shared by different optimization schemes, existing works propose their diagnosis methods in an ad-hoc manner, which yield multiple limitations. First, the diagnosis metric in an optimization scheme may sometimes be less accurate than others; second, existing schemes fail to fully exploit the diagnosis result by applying it for only one optimization action; third, existing methods usually do not perceive cross-client data heterogeneity, failing to simultaneously enhance FL accuracy. To tackle those limitations, we make a systematical study on the training diagnosis methods of multiple optimization schemes, and propose metric grafting—replacing a scheme's diagnosis metric with a better one to improve the training performance. Moreover, to fully exploit the potential of training diagnosis, we build a system platform that supports flexible combinations of training diagnosis and optimization actions (i.e., single-diagnosis-multiple actions and multiple-diagnosis-multiple-actions). Evaluation on testbeds show that, with metric grafting and advanced diagnosis action combinations, we can substantially improve the efficiency and accuracy performance of FL. Jiayi Zhang 0006, Zuo Gan, Chen Chen 0067, Zhifeng Jiang 0001, Hao Wang 0022, Yifei Zhu 0001, Quan Chen 0002, Minyi Guo |
IEEE Trans. Mob. Comput. | 3 |
| 2026 | Mitigating Server-Side Communication Bottlenecks in Distributed Learning With Round-Robin Participant CoordinationabstractDeep neural networks are increasingly trained in a distributed manner—in either clusters or with federated devices, where the participants jointly refine the global model with their gradients calculated locally. More often than not, those gradients are collected to a central server in a synchronous manner to avoid the negative impact of stale updates. However, when all the participants communicate their gradients to the server in such a uniform pace, the network on the server side—under intense contention—often becomes a performance bottleneck. To address this problem, for the cluster environment we propose theRound-Robin Synchronous Parallel(R2SP) scheme, which coordinates the participants to make updates in anevenly-gapped,round-robinmanner. This way, we can minimize the network contention with a minimum cost of the update quality; we also propose to incorporate adaptive batch sizing in R2SP to address the hardware heterogeneity among workers. Moreover, for the federated learning (FL) scenarios, we note that it is necessary yet challenging to apply the insight of R2SP to mitigate the network bottleneck in the FL server, given that there are a huge number of participants with unstable resources and inconsistent data distributions. To tackle those challenges, we further propose FL-R2SP, which extents the coordination units from individual participants to participantgroups—with the resource instability and data heterogeneity tackled within each group. We have implemented R2SP and FL-R2SP respectively with TensorFlow and PyTorch, and extensive EC2 experiments show that R2SP and FL-R2SP can respectively speed up model convergence for clustered and federated scenarios by over 20%. Jiayi Zhang 0006, Chen Chen 0067, Zuo Gan, Wei Wang 0030, Bo Li 0001, Minyi Guo |
IEEE Trans. Netw. | 2 |
| 2026 | Flexible Synchronization Control for Accurate and Efficient Federated LearningabstractFederated Learning (FL) is a distributed paradigm that supports collaborated model training while preserving data privacy, where clients periodically synchronize their local gradients once after multiple local iterations. Due to non-uniform data distribution and poor network condition, FL processes often suffer degraded training accuracy and efficiency. In this work, we analyze the microscopic parameter variation behaviors in FL, and find that an effective method to improve FL accuracy is to switch to more frequent synchronization at proper moments. In particular, such frequency-tuning moments—which can be detected from gradient characteristics—areheterogeneousacross different parameters. Motivated by such observations, we propose Parameter-Adaptive Synchronization (PAS), a FL scheme that adaptively tunes the synchronization period for each scalar parameter. The benefits of PAS are two-fold: By switching to more frequent synchronization when necessary, we can improve the FL training accuracy; by synchronizing different parameters independently, we can enable communication-computation overlapping and enhance the network utilization. We have theoretically demonstrated the convergence validity of PAS, and have further extended it with adaptive sparsification capability to jointly reduce the overall communication volume. We implemented PAS atop PyTorch, and extensive experiments show that it can substantially improve FL performance in both accuracy and communication efficiency. Zuo Gan, Chen Chen 0067, Jiayi Zhang 0006, Yifei Zhu 0001, Jieru Zhao, Quan Chen 0002, Minyi Guo |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2026 | gPooling: An Elastic GPU Resource Management Framework for On-Demand Virtualization in Shared Accelerator Clusters
Kaicheng Guo, Chen Chen 0067, Yun Wang 0039, Pengwei Du, Zhengwei Qi, Haibing Guan |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2026 | Enabling Client-Autonomous Training Optimizations for Efficient Federated LearningabstractFederated Learning (FL) enables collaborate model training without privacy violation, where clients periodically report their updates to the server in communication rounds. Due to heterogeneous resource and limited bandwidth, FL processes often suffer from low efficiency. Existing works in that regard are oblivious to the intra-round execution status on clients, failing to tackle runtime stragglers or hide the communication overheads for some early-converged layers. In this paper, we propose FedCA, a novel mechanism that allows clients to autonomously exploit intra-round training status for higher efficiency while preserving accuracy performance. We first devise a metric to help quantify the statistical contribution of different iterations in a round, which can be efficiently profiled at runtime with the periodical sampling strategy. With the instantaneous system and statistical status, to improve computation efficiency, clients under FedCA can adaptively determine the intra-round workloads based on a utility function depicting the marginal computation benefit. Besides, to mitigate the communication bottleneck, for some parameters attaining fast local convergence, clients under FedCA can eagerly transmit their updates to the FL server prior to round completion. We also extend FedCA to FedCA+, integrating speculative sparsification to futher reduce the cumulative communication amount. We implemented FedCA and FedCA+ atop PyTorch, and large-scale experiments show that they can improve the FL efficiency by up to 45.3%. Na Lyu, Jiayi Zhang 0006, Zhi Shen, Chen Chen 0067, Zhifeng Jiang 0001, Quan Chen 0002, Minyi Guo |
IEEE Trans. Parallel Distributed Syst. | 4 |
| 2025 | SemanticPrefetcher: Accelerate Data Lake Access with Semantics-Aware File PrefetchingabstractStorage-compute disaggregation has become a mainstream paradigm in cloud computing, yet data lake workloads introduce distinct prefetching challenges: massive numbers of small files, interleaved multi-tenant streams, and frequent one-time accesses. Under these conditions, traditional sequential, correlation-based, and semantic prefetching methods become ineffective, leading to cache inefficiency and high latency. We propose SemanticPrefetcher, a lightweight, semantic-aware prefetching mechanism that incrementally constructs meaningful access streams at runtime. The system operates in three stages: it tokenizes file paths or object names into a unified semantic representation, clusters related requests into coherent streams despite multi-tenant interleaving, and detects naming regularities to predict future accesses. This design transforms mixed and seemingly disordered requests into predictable access flows without application modifications. Implemented on JuiceFS, SemanticPrefetcher reduces end-to-end execution time by up to 39.6% and read latency by 79.3% compared to state-of-the-art baselines. These results demonstrate that implicit semantics in file paths can be effectively leveraged for robust and efficient prefetching in cloud-scale data lakes. Tianze Wang, Guanjie Wang, Mingyan Yang, Manqi Luo, Mingchuan Zou, Chen Chen 0067, Minyi Guo |
CloudCom | 6 |
| 2025 | DevTrace: Lightweight Plug-In Design for PCIe Transaction Tracing in Edge Intelligence WorkloadsabstractThe complexity of host-peripheral interactions during high-load tasks poses significant challenges for system optimization, with existing tracing tools degrading performance by up to 5.39×. We introduce DevTrace, a novel low-overhead tracing framework for peripheral interactions. Its modular architecture separates data collection from kernel-level operations, enabling lightweight tracing with minimal driver modifications across entire classes of devices. By eliminating heavy kernel tracing interrupts, DevTrace reduces overhead to negligible levels while maintaining data accuracy. In edge-based intelligence deployments, DevTrace achieves a 128× reduction in memory usage and approximately 10× lower CPU overhead compared to page-fault-based solutions. It significantly reduces data loss and performance degradation under high-load conditions, establishing it as a reliable tool for analyzing host-peripheral interactions and optimizing performance in resource-constrained environments. We also discuss potential extensions to eBPF to further decouple tracing from driver frameworks. Zhibai Huang, Kailiang Xu, Zhixiang Wei, Yinghao Deng, Chen Chen 0067, Yun Wang 0039, Fangxin Liu, Mingyuan Xia 0001, Zhengwei Qi |
ICCAD | 5 |
| 2025 | FedSU: Communication-efficient Federated Learning with Speculative UpdatingabstractFederated learning enables mobile devices to collaboratively learn a global model in iterative communication rounds. Many sparsification methods have been proposed for communication compression of FL, working by not synchronizing insignificant updates. However, we find there still exist unexploited sparsification opportunities: given the update similarity across different rounds, parameters often exhibit a linear updating pattern; motivated by speculative execution in computer architecture domain, it is promising to use the predicted gradients to refine the linearly-updating parameters without synchronization. To that end, we propose Federated Learning with Speculative Updating, or FedSU, to attain larger sparsification ratio without compromising model accuracy. In particular, to identify the linearly-updating parameters efficiently at runtime, we devise a regression-free method that diagnoses parameter linearity based on whether the second-order parameter difference is oscillating around 0. Meanwhile, to ensure convergence validity, FedSU leverages the prediction error as a feedback signal—so as to timely return to regular updating if the parameter no longer follows the linear pattern in reality. We have implemented FedSU as a Python module, and large-scale experiments in an emulated FL setup confirm that FedSU can remarkably improve the communication efficiency of FL, with a convergence speedup of over 40%. Chen Chen 0067, Qinbin Li, Jieru Zhao, Shixuan Sun, Bo Li 0001, Minyi Guo |
ICDCS | 2 |
| 2025 | LLMSched: Uncertainty-Aware Workload Scheduling for Compound LLM ApplicationsabstractDeveloping compound Large Language Model (LLM) applications is becoming an increasingly prevalent approach to solving real-world problems. In these applications, an LLM collaborates with various external modules, including APIs and even other LLMs, to realize complex intelligent services. However, we reveal that the intrinsic duration and structural uncertainty in compound LLM applications pose great challenges for LLM service providers in serving and scheduling them efficiently. In this paper, we propose LLMSched, an uncertainty-aware scheduling framework for emerging compound LLM applications. In LLMSched, we first design a novel DAG-based model to describe the uncertain compound LLM applications. Then, we adopt the Bayesian network to comprehensively profile compound LLM applications and identify uncertainty-reducing stages, along with an entropy-based mechanism to quantify their uncertainty reduction. Combining an uncertainty reduction strategy and a job completion time (JCT)-efficient scheme, we further propose an efficient scheduler to reduce the average JCT. Evaluation of both simulation and testbed experiments on various representative compound LLM applications shows that compared to existing state-of-the-art scheduling schemes, LLMSched can reduce the average JCT by 14 ~ 79%. Botao Zhu, Chen Chen 0067, Xiaoyi Fan 0001, Yifei Zhu 0001 |
ICDCS | 2 |
| 2025 | Reducing the End-to-End Latency of DNN-Based Recommendation Systems in GPU PoolsabstractWhile intelligent applications (e.g., recommendation systems) prefer different CPU-GPU ratios, GPU pooling technique that decouples the GPU and CPU resources yields substantial flexibility when serving diverse applications. With such architecture, DNN-based recommendation services often offload the compute-intensive neural network layers to the remote GPU pool for high resource utilization. However, such a paradigm results in the long end-to-end latency due to two causes: 1) the intermediate data is copied for multiple times during the entire process in current GPU pooling practices, incurring heavy overheads; 2) the content transferred to the GPU pool involves multiple small tensors, suffering from poor bandwidth efficiency. To solve these problems, we design Zero, a runtime system that incorporates a zero-copy transmission mechanism as well as a dynamic tensor merging policy. The zero-copy transmission mechanism unifies memory management across the inference framework and the RPC framework, accompanied by an elaborated serialization protocol to fully eliminate redundant data copying. Meanwhile, the tensor merging policy deliberately organizes small tensors into larger data blocks, so as to transfer them with higher efficiency. Experimental results show that, compared with prior work, Zero reduces the latency of typical recommendation models by up to 15.1% (10.1% on average). Guangqiang Luan, Pu Pang, Quan Chen 0002, Chen Chen 0067, Guoyao Xu, Chi Zhang 0005, Yanyi Zi, Yinghao Yu, Liping Zhang 0013, Minyi Guo |
IPDPS | 4 |
| 2025 | Lumina: Real-Time Neural Rendering by Exploiting Computational Redundancyabstract3D Gaussian Splatting (3DGS) has vastly advanced the pace of neural rendering, but it remains computationally demanding on today's mobile SoCs.To address this challenge, we propose Lumina, a hardware-algorithm co-designed system, which integrates two principal optimizations: a novel algorithm, S 2 , and a radiance caching mechanism, RC, to improve the efficiency of neural rendering.S 2 algorithm exploits temporal coherence in rendering to reduce the computational overhead, while RC leverages the color integration process of 3DGS to decrease the frequency of intensive rasterization computations.Coupled with these techniques, we propose an accelerator architecture, LuminCore, to further accelerate cache lookup and address the fundamental inefficiencies in Rasterization.We show that Lumina achieves 4.5× speedup and 5.3× energy reduction against a mobile Volta GPU, with a marginal quality loss (< 0.2 dB peak signal-to-noise ratio reduction) across synthetic and real-world datasets. Yu Feng 0007, Weikai Lin, Yuge Cheng, Zihan Liu 0002, Jingwen Leng, Minyi Guo, Chen Chen 0067, Shixuan Sun, Yuhao Zhu 0001 |
ISCA | 7 |
| 2025 | ARMing x86 Games: Accelerating Binary Translation Using Software-Only Validated Flag Speculation
James Yen, Zhibai Huang, Zhixiang Wei, Chen Chen 0067, Senhao Yu, Yun Wang 0039, Hao Wang 0022, Zhengwei Qi |
MobiSys | 6 |
| 2025 | RetrievalAttention: Accelerating Long-Context LLM Inference via Vector RetrievalabstractTransformer-based Large Language Models (LLMs) have become increasingly important. However, scaling LLMs to longer contexts incurs slow inference speed and high GPU memory consumption for caching key-value (KV) vectors. This paper presents RetrievalAttention, a training-free approach to both accelerate the decoding phase and reduce GPU memory consumption by pre-building KV vector indexes for fixed contexts and maintaining them in CPU memory for efficient retrieval. Unlike conventional KV cache methods, RetrievalAttention integrate approximate nearest neighbor search (ANNS) indexes into attention computation. We observe that off-the-shelf ANNS techniques often fail due to the out-of-distribution (OOD) nature of query and key vectors in attention mechanisms. RetrievalAttention overcomes this with an attention-aware vector index. Our evaluation shows RetrievalAttention achieves near full attention accuracy while accessing only 1-3\% of the data, significantly reducing inference costs. Remarkably, RetrievalAttention enables LLMs with 8B parameters to handle 128K tokens on a single NVIDIA RTX4090 (24GB), achieving a decoding speed of 0.107 seconds per token. Baotong Lu, Huiqiang Jiang, Zhenhua Han, Qianxi Zhang, Qi Chen 0009, Chengruidong Zhang, Bailu Ding, Chen Chen 0067, Fan Yang 0024, Yuqing Yang 0001, Lili Qiu |
NeurIPS | 11 |
| 2025 | RapidStore: An Efficient Dynamic Graph Storage System for Concurrent QueriesabstractDynamic graph storage systems are essential for real-time applications such as social networks and recommendation, where the graph continuously evolves. However, they face significant challenges in efficiently handling concurrent read and write operations. We find that existing methods suffer from write queries interfering with read efficiency, substantial time and space overhead due to per-edge versioning, and an inability to balance performance, such as slow searches. To address these issues, we propose RapidStore, a holistic approach for efficient in-memory dynamic graph storage designed for read-intensive workloads. Our key idea is to exploit the characteristics of graph queries through a decoupled system design that separates the management of read and write queries and decouples version data from graph data. Besides, we design an efficient dynamic graph store to cooperate with the graph concurrency control mechanism. Experiments show that RapidStore enables fast and scalable concurrent graph queries, effectively balancing the performance of inserts, searches, and scans, and significantly improving efficiency in dynamic graph storage systems. Chiyu Hao, Jixian Su, Shixuan Sun, Hao Zhang 0098, Jianwen Zhao, Chenyi Zhang 0002, Jieru Zhao, Chen Chen 0067, Minyi Guo |
Proc. VLDB Endow. | 9 |
| 2025 | Taming Flexible Job Packing in Deep Learning Training ClustersabstractJob packing is an effective technique to harvest the idle resources allocated to the deep learning (DL) training jobs but not fully utilized, especially when clusters may experience low utilization, and users may overestimate their resource needs. However, existing job packing techniques tend to be conservative due to the mismatch in scope and granularity between job packing and cluster scheduling. In particular, tapping the potential of job packing in the training cluster requires a local and fine-grained coordination mechanism. To this end, we propose a novel job-packing middleware named Gimbal , which operates between the cluster scheduler and the hardware resources. As middleware, Gimbal must not only facilitate coordination among the packed jobs but also support various scheduling objectives of different schedulers. Gimbal achieves dual functionality by introducing a set of worker calibration primitives designed to calibrate workers’ execution status in a fine-grained manner. The primitives obscure the complexity of the underlying job and resource management mechanisms, thus offering the generality and extensibility for crafting coordination policies tailored to various scheduling objectives. We implement Gimbal on a real-world GPU cluster and evaluate it with a set of representative DL training jobs. The results show that Gimbal improves different scheduling objectives up to 1.32× compared with the state-of-the-art job packing techniques. Pengyu Yang, Weihao Cui, Chunyu Xue, Han Zhao 0005, Chen Chen 0067, Quan Chen 0002, Jing Yang 0017, Minyi Guo |
ACM Trans. Archit. Code Optim. | 5 |
| 2025 | EDAS: Enabling Fast Data Loading for GPU Serverless ComputingabstractIntegrating GPUs into serverless computing platforms is crucial for improving efficiency. Many GPU functions, such as DNN inferences and scientific services, benefit from GPU usage, which requires only tens to hundreds of milliseconds for pure computation. Under these circumstances, fast data loading is imperative for function performance. However, existing GPU serverless systems face significant data stall issues, leading to extremely low GPU efficiency. Faced with the above problems, we observe opportunities to optimize data loading, such as data preloading and deduplicated data loading. However, these optimizations are impossible in existing GPU serverless systems due to the lack of insights into data information, such as data sizes and read-write attributes of function inputs. To address this, we propose a novel GPU serverless system, EDAS. EDAS first enhances user request specifications, allowing users to annotate data retrieved by GPU functions from the database with additional attributes. Based on this, EDAS takes over data loading from GPU functions and proposes two innovative data loading management schemes: a parallelized data loading scheme and a multi-stage resource exit scheme. Our experimental results show that EDAS reduces function duration by 16.2× and improves system throughput by 1.91× compared with the state-of-the-art serverless platform. Han Zhao 0005, Weihao Cui, Quan Chen 0002, Zijun Li 0001, Zhenhua Han, Yu Feng 0007, Jieru Zhao, Chen Chen 0067, Jingwen Leng, Minyi Guo |
ACM Trans. Archit. Code Optim. | 9 |
| 2025 | Exploring Efficient Hardware Accelerator for Learning-Based Image CompressionabstractRecently, learning-based image compression (LIC) methods have surpassed manually designed approaches in both compression quality and bitrate. However, increasing computational demands and insufficient optimizations in codec performance have hindered the advancement of LIC acceleration. Most researches focus on optimizing specific components, often neglecting the sources of underutilization during the execution of LIC models. Generally, efficient LIC acceleration encounters three primary challenges: 1) extra overheads introduced by individual optimizations; 2) load and computation imbalances in small kernels; and 3) mismatches between hardware configurations and the LIC models. To address these challenges, we propose a framework named extensive accelerator for LIC (X-LIC) for efficiently exploring the design space under constrained resources. First, we quantitatively characterize a representative LIC model, including its latency, computation size, and temporal utilization across various accelerators. We design a hardware-optimized quantization method to compensate for the lack of LIC-oriented research, particularly regarding data precision, distortion, and resource consumption. Additionally, we propose a parameterized LIC accelerator architecture that integrates seamlessly with existing loop optimization models and supports various LIC operators. Two optimization schemes are proposed for redundant computation in transposed convolution and load and computation imbalance in small kernels. Experimental results show that our framework demonstrates significant flexibility across a broad design space, achieving an average of 78%–95% of the theoretical peak performance and up to 688.2/759.1 GOP/s en/de-coder performance with INT8 precision. As a result, the en/de-coder performance can reach up to 33/36 FPS in 720P resolution. An FPGA demo of X-LIC is available athttps://github.com/sjtu-tcloud/X-LIC. Chen Chen 0067, Kaicheng Guo, Xingzi Yu, Weidong Qiu, Zhengwei Qi, Haibing Guan |
IEEE Trans. Comput. Aided Des. Integr. Circuits Syst. | 1 |
| 2025 | Trident: A Provider-Oriented Resource Management Framework for Serverless Computing PlatformsabstractServerless computing has become increasingly popular due to its flexible and hassle-free service, relieving users from traditional resource management burdens. However, the shift in responsibility has led to unprecedented challenges for serverless providers in managing virtual machines (VMs) and serving heterogeneous function instances. Serverless providers need to purchase, provision and manage VM instances from IaaS providers, aiming to minimize VM provisioning costs while ensuring compliance with Service Level Objectives (SLOs). In this paper, we propose Trident, a provider-oriented resource management framework for serverless computing platforms. Trident optimizes three major serverless computing provisioning problems for serverless providers: workload prediction, VM provisioning, and function placement. Specifically, Trident introduces a novel dynamic model selection algorithm for more accurate workload prediction. With the prediction results, Trident then carefully designs a hierarchical reinforcement learning (HRL)-based approach for VM provisioning with a mix of types and configurations. To further improve resource utilization, Trident employs an effective collocation placement strategy for efficient function container scheduling. Evaluations on the Azure Function dataset demonstrate that Trident maintains the lowest probability of violating SLOs while simultaneously achieving substantial cost savings of up to 71.8% in provisioning expense compared to state-of-the-art methods from industry and academia. Botao Zhu, Yifei Zhu 0001, Chen Chen 0067, Linghe Kong |
IEEE Trans. Serv. Comput. | 3 |
| 2024 | An Optimizing Framework on MLIR for Efficient FPGA-based Accelerator GenerationabstractWith the increasing demand for computing capability given limited resource and power budgets, it is prominent to deploy applications to customized accelerators like FPGAs. However, FPGA programming is non-trivial. Although existing high-level synthesis (HLS) tools improve productivity to a certain extent, they are limited in scope and capability to support sufficient FPGA-oriented transformations and optimizations. This paper focuses on FPGA-based accelerators and proposes POM, an end-to-end optimizing framework built on multi-level intermediate representation (MLIR). POM has several features which demonstrate its scope and capability of performance optimization. First, most HLS tools depend exclusively on a single-level IR like LLVM IR to perform all the optimizations, introducing excessive information into the IR and making debugging an arduous task. In contrast, POM explicitly introduces three layers of IR to perform operations at suitable abstraction levels, streamlining the implementation and debugging process and exhibiting better flexibility, extensibility, and systematicness. Second, POM integrates the polyhedral model into MLIR and hence enables advanced dependence analysis and a wide range of FPGA-oriented loop transformations. By representing nested loops with integer sets and maps at suitable IR, loop transformations can be conducted conveniently through a series of manipulations on polyhedral semantics. Finally, to further relieve design effort, POM is equipped with a user-friendly programming interface (DSL) that allows a concise description of computation and includes a rich collection of scheduling primitives. An automatic design space exploration (DSE) engine is also provided to search for high-performance optimization schemes efficiently and generate optimized accelerators automatically. Experimental results show that POM achieves a 6.46× average speedup on typical benchmark suites and a 6.06 ×average speedup on real-world applications compared to the state-of-the-art. Weichuang Zhang, Jieru Zhao, Guan Shen, Quan Chen 0002, Chen Chen 0067, Minyi Guo |
HPCA | 5 |
| 2024 | FedCA: Efficient Federated Learning with Client AutonomyabstractFederated Learning (FL) enables collaborate model training without privacy violation, where clients periodically report their updates to the server in communication rounds. Due to heterogeneous resource and limited bandwidth, FL processes often suffer from low efficiency. Existing works in that regard are oblivious to the intra-round execution status on clients; however, such status information has great potential to support flexible efficiency optimizations. In this paper, we propose FedCA, a novel mechanism that allows clients to autonomously exploit intra-round training status for higher efficiency. We first devise a metric to help quantify the statistical contribution of different iterations in a round, which can be efficiently profiled at runtime with the periodical sampling strategy. With the instantaneous system and statistical status, to improve computation efficiency, clients under FedCA can adaptively determine the intra-round workloads based on a utility function. Besides, to mitigate the communication bottleneck, for some parameters attaining fast local convergence, clients under FedCA can eagerly transmit their updates to the FL server prior to round completion. We implemented FedCA atop PyTorch, and large-scale experiments show that it can improve FL efficiency by over 15%. Zhi Shen, Chen Chen 0067, Zhifeng Jiang 0001, Jiayi Zhang 0006, Quan Chen 0002, Minyi Guo |
ICPP | 3 |
| 2024 | DPBalance: Efficient and Fair Privacy Budget Scheduling for Federated Learning as a ServiceabstractFederated learning (FL) has emerged as a prevalent distributed machine learning scheme that enables collaborative model training without aggregating raw data. Cloud service providers further embrace Federated Learning as a Service (FLaaS), allowing data analysts to execute their FL training pipelines over differentially-protected data. Due to the intrinsic properties of differential privacy, the enforced privacy level on data blocks can be viewed as a privacy budget that requires careful scheduling to cater to diverse training pipelines. Existing privacy budget scheduling studies prioritize either efficiency or fairness individually. In this paper, we propose DPBalance, a novel privacy budget scheduling mechanism that jointly optimizes both efficiency and fairness. We first develop a comprehensive utility function incorporating data analyst-level dominant shares and FL-specific performance metrics. A sequential allocation mechanism is then designed using the Lagrange multiplier method and effective greedy heuristics. We theoretically prove that DPBalance satisfies Pareto Efficiency, Sharing Incentive, Envy-Freeness, and Weak Strategy Proofness. We also theoretically prove the existence of a fairness-efficiency tradeoff in privacy budgeting. Extensive experiments demonstrate that DPBalance outperforms state-of-the-art solutions, achieving an average efficiency improvement of 1.44× ~ 3.49×, and an average fairness improvement of 1.37×~24.32×. Zibo Wang 0001, Yifei Zhu 0001, Chen Chen 0067 |
INFOCOM | 4 |
| 2024 | PAS: Towards Accurate and Efficient Federated Learning with Parameter-Adaptive SynchronizationabstractFederated Learning (FL) is a distributed paradigm that supports collaborated model training while preserving data privacy, where clients periodically synchronize their local gradients once after multiple local iterations. Due to non-uniform data distribution and poor network condition, FL processes often suffer degraded training accuracy and efficiency. In this work, we analyze the microscopic parameter variation behaviors in FL, and find that an effective method to improve FL accuracy is to switch to more frequent synchronization at proper moments. Moreover, such moments can be detected from gradient characteristics, and are heterogeneous across different parameters. Motivated by such observations, we propose Parameter-Adaptive Synchronization (PAS), a FL scheme that adaptively tunes the synchronization period for each scalar parameter. The benefits of PAS are two-fold: By switching to more frequent synchronization when necessary, we can improve the FL training accuracy; by synchronizing different parameters independently, we can enable communication-computation overlapping and enhance the network utilization. We implemented PAS atop PyTorch, and extensive experiments show that it can substantially improve FL performance in both accuracy and communication efficiency. Zuo Gan, Chen Chen 0067, Jiayi Zhang 0006, Gaoxiong Zeng, Yifei Zhu 0001, Jieru Zhao, Quan Chen 0002, Minyi Guo |
IWQoS | 2 |
| 2024 | Towards Efficient Compound Large Language Model System Serving in the WildabstractUtilizing compound Large Language Model (LLM) systems, instead of a monolithic LLM model, is gradually becoming a practical solution to realize a diverse range of industry applications. In compound LLM systems, an LLM collaborates with other external tools, APIs, or LLMs to offer intelligent services. In this poster, we identify the unique challenges, namely temporal and topological uncertainty, brought about by compound LLM systems in system serving. We then propose a priority-based scheduling policy to schedule different stages in DAG-represented compound LLM systems. The preliminary results show promising performance of uncertainty-aware scheduling policies. Yifei Zhu 0001, Botao Zhu, Chen Chen 0067, Xiaoyi Fan 0001 |
IWQoS | 3 |
| 2024 | Parrot: Efficient Serving of LLM-based Applications with Semantic Variable
Chaofan Lin, Zhenhua Han, Chengruidong Zhang, Yuqing Yang 0001, Fan Yang 0024, Chen Chen 0067, Lili Qiu |
OSDI | 6 |
| 2024 | Synchronize Only the Immature Parameters: Communication-Efficient Federated Learning By Freezing Parameters AdaptivelyabstractFederated learning allows edge devices to collaboratively train a global model without sharing their local private data. Yet, with limited network bandwidth at the edge, communication often becomes a severe bottleneck. In this paper, we find that it is unnecessary to always synchronize the full model in the entire training process, because many parameters already become mature (i.e., stable) prior to model convergence, and can thus be excluded from later synchronizations. This allows us to reduce the communication overhead without compromising the model accuracy. However, challenges are that the local parameters excluded from global synchronization may diverge on different clients, and meanwhile some parameters may stabilize only temporally. To address these challenges, we propose a novel scheme called Adaptive Parameter Freezing (APF), which fixes (freezes) the non-synchronized stable parameters in intermittent periods. Specifically, the freezing periods are tentatively adjusted in an additively-increase and multiplicatively-decrease manner—depending on whether the previously-frozen parameters remain stable in subsequent iterations. We also extend APF into APF# and APF++, which freeze parameters in a more aggressive manner to achieve larger performance benefit for large complex models. We implemented APF and its variants as Python modules with PyTorch, and extensive experiments show that APF can reduce data transfer amount by over 60%. Chen Chen 0067, Hong Xu 0001, Wei Wang 0030, Baochun Li, Bo Li 0001, Li Chen 0008, Gong Zhang 0001 |
IEEE Trans. Parallel Distributed Syst. | 1 |
| 2023 | DataFlower: Exploiting the Data-flow Paradigm for Serverless Workflow OrchestrationabstractServerless computing that runs functions with auto-scaling is a popular task execution pattern in the cloud-native era. By connecting serverless functions into workflows, tenants can achieve complex functionality. Prior research adopts the control-flow paradigm to orchestrate a serverless workflow. However, the control-flow paradigm inherently results in long response latency, due to the heavy data persistence overhead, sequential resource usage, and late function triggering. Zijun Li 0001, Chuhao Xu, Quan Chen 0002, Jieru Zhao, Chen Chen 0067, Minyi Guo |
ASPLOS (4) | 5 |
| 2023 | STAG: Enabling Low Latency and Low Staleness of GNN-based Services with Dynamic GraphsabstractMany emerging user-facing services adopt Graph Neural Networks (GNNs) to improve serving accuracy. When the graph used by a GNN model changes, representations (embedding) of nodes in the graph should be updated accordingly. However, the node representation update is too slow, resulting in either long response latency of user queries (inference is performed after update completes) or high staleness problem (inference is performed based on stale data).Our in-depth analysis shows that the slow update is mainly due to neighbor explosion problem in graphs and duplicated computation. Based on such findings, we propose STAG, a GNN serving framework that enables low latency and low staleness of GNN-based services. It comprises a collaborative serving mechanism and an additivity-based incremental propagation strategy. With collaborative serving mechanism, only part of node representations are updated during the update phase, and the final representations are calculated in the inference phase. It alleviates the neighbor explosion problem. The additivity-based incremental propagation strategy reuses intermediate data during update phase, eliminating duplicated computation. Experimental results show that STAG greatly reduces staleness time with a slight increase in response latency, and support 2.7~27x workload compared to existing approaches. Quan Chen 0002, Deze Zeng, Chen Chen 0067, Minyi Guo |
ICCD | 5 |
| 2023 | ASFL: Adaptive Semi-asynchronous Federated Learning for Balancing Model Accuracy and Total Latency in Mobile Edge NetworksabstractFederated learning (FL) is a new paradigm for privacy-preserving learning. This is particularly appealing in the mobile edge network (MEN), in which devices collectively train a global model with their own set of data. It is, however, routinely difficult for FL algorithms to satisfy different training task preferences in terms of the total latency and model accuracy due to a number of factors including the straggler effect, data heterogeneity, communication bottleneck and device mobility. To this end, we propose an Adaptive Semi-asynchronous Federated Learning (ASFL) framework, which adaptively balances the total latency and model accuracy according to the task preferences in MEN. Specifically, ASFL conducts a two-stage operation: i) Device selection stage. Each global round selects a set of devices that can maximize the model accuracy to eliminate data heterogeneity and communication bottlenecks; ii) Training stage. We first define a latency-accuracy objective value to model the balance between the latency and accuracy. Then in each global round, we use a deep reinforcement learning (DRL) algorithm based on soft actor-critic with discrete actions to intelligently derive the number of picked devices (i.e., participants in the current global aggregation) and the lag tolerance at each global round to maximize the latency-accuracy objective value. Extensive experiments show that ASFL can improve the latency-accuracy objective value by up to 94% compared with three state-of-the-art FL frameworks. Jieling Yu, Ruiting Zhou, Chen Chen 0067, Bo Li 0001, Fang Dong 0001 |
ICPP | 3 |
| 2023 | HaeNAS: Hardware-Aware Efficient Neural Architecture Search via Zero-Cost ProxyabstractThe practical use of advanced DNN models is hindered by limited hardware resources and high computation demands.Neural Architecture Search (NAS) is becoming a default technique which automatically discovers architectures that are competitive with handcraft ones.However, existing methods only prioritize accuracy and overlook hardware-related factors.To address this, we introduce HaeNAS (Hardwareaware efficient NAS), which considers both accuracy and computation cost on specific hardware platforms.The search space of HaeNAS consists of several stages, each allowing different convolution kernels, layer numbers, layer widths, and operation types.We use a data-driven approach to predict the latency and energy consumption on target hardware, and we improve the zero-cost proxy based on network pruning research to speed up the NAS process.With these techniques, HaeNAS finds a target network within 160 GPU hours, which achieves 80.7% top-1 accuracy on the ImageNet, with a latency of 10.4ms and energy consumption of 931mJ. Yaodanjun Ren, Chen Chen 0067, Zhengwei Qi |
SEKE | 2 |
| 2023 | GIFT: Toward Accurate and Efficient Federated Learning With Gradient-Instructed Frequency TuningabstractFederated learning (FL) enables distributed clients to collectively train a global model without revealing their private data, and for efficiency clients synchronize their gradients periodically. However, this can lead to the inaccuracy in model convergence due to inconsistent data distributions among clients. In this work, we find that there is a strong correlation between FL accuracy loss and the synchronization frequency, and seek to fine tune the synchronization frequency at training runtime to make FL accurate and also efficient. Specifically, aware that under the FL privacy requirement only gradients can be utilized for making frequency tuning decisions, we propose a novel metric called gradient consistency, which can effectively reflect the training status despite the instability of realistic FL scenarios. We further devise a feedback-driven algorithm called Gradient-Instructed Frequency Tuning (GIFT), which adaptively increases or decreases the synchronization frequency based on the gradient consistency metric. We have implemented GIFT in PyTorch, and large-scale evaluations show that it can improve FL accuracy by up to 10.7% with a time reduction of 58.1%. Chen Chen 0067, Hong Xu 0001, Wei Wang 0030, Baochun Li, Bo Li 0001, Li Chen 0008, Gong Zhang 0001 |
IEEE J. Sel. Areas Commun. | 1 |
| 2023 | Accelerating Distributed Learning in Non-Dedicated EnvironmentsabstractMachine learning (ML) models are increasingly trained with distributed workers possessing heterogeneous resources. In such scenarios, model training efficiency may be negatively affected bystragglers—workers that run much slower than others. Efficient model training requires eliminating such stragglers, yet for modern ML workloads, existing load balancing strategies are inefficient and even infeasible. In this article, we propose a novel strategy, calledsemi-dynamic load balancing, to eliminate stragglers of distributed ML workloads. The key insight is that ML workers shall be load-balanced atiteration boundaries, being non-intrusive to intra-iteration execution. Based on it we further develop LB-BSP, an integrated worker coordination mechanism that adapts workers’ load to their instantaneous processing capabilities—by right-sizing the sample batches at the synchronization barriers. We have designed distinct load tuning algorithms for ML in CPU clusters, in GPU clusters as well as in federated learning setups, based on their respective characteristics. LB-BSP has been implemented as a Python module for ML frameworks like TensorFlow and PyTorch. Our EC2 deployment confirms that LB-BSP is practical, effective and light-weight, and is able to accelerating distributed training by up to 54 percent. Chen Chen 0067, Qizhen Weng 0001, Wei Wang 0030, Baochun Li, Bo Li 0001 |
IEEE Trans. Cloud Comput. | 1 |
| 2022 | Characterizing and orchestrating VM reservation in geo-distributed clouds to improve the resource efficiencyabstractCloud providers often build a geo-distributed cloud from multiple datacenters in different geographic regions, to serve tenants at different locations. The tenants that run large scale applications often reserve resources based on their peak loads in the region close to the end users to handle the ever changing application load, wasting a large amount of resources. We therefore characterize the VM request patterns of the top tenants in our production public geo-distributed cloud, and open-source the VM request traces in four months from the top 20 tenants of our cloud. The characterization shows that the resource usage of large tenants has various temporal and spatial patterns on the dimensions of time series, regions, and VM types, and has the potential of peak shaving between different tenants to further reduce the resource reservation cost. Based on the findings, we propose a resource reservation and VM request scheduling scheme named ROS to minimize the resource reservation cost while satisfying the VM allocation requests. Our experiments show that ROS reduces the overall deployment cost by 75.4% and the reservation resources by 60.1%, compared to the tenant-specified reservation strategy. Jiuchen Shi, Kaihua Fu, Quan Chen 0002, Changpeng Yang, Mosong Zhou, Jieru Zhao, Chen Chen 0067, Minyi Guo |
SoCC | 8 |
| 2021 | Communication-Efficient Federated Learning with Adaptive Parameter FreezingabstractFederated learning allows edge devices to collaboratively train a global model by synchronizing their local updates without sharing private data. Yet, with limited network bandwidth at the edge, communication often becomes a severe bottleneck. In this paper, we find that it is unnecessary to always synchronize the full model in the entire training process, because many parameters gradually stabilize prior to the ultimate model convergence, and can thus be excluded from being synchronized at an early stage. This allows us to reduce the communication overhead without compromising the model accuracy. However, challenges are that the local parameters excluded from global synchronization may diverge on different clients, and meanwhile some parameters may stabilize only temporally. To address these challenges, we propose a novel scheme called Adaptive Parameter Freezing (APF), which fixes (freezes) the non-synchronized stable parameters in intermittent periods. Specifically, the freezing periods are tentatively adjusted in an additively-increase and multiplicatively-decrease manner, depending on if the previously-frozen parameters remain stable in subsequent iterations. We implemented APF as a Python module in PyTorch. Our extensive array of experimental results show that APF can reduce data transfer by over 60%. Chen Chen 0067, Hong Xu 0001, Wei Wang 0030, Baochun Li, Bo Li 0001, Li Chen 0008, Gong Zhang 0001 |
ICDCS | 1 |
| 2021 | Two-Dimensional Learning Rate Decay: Towards Accurate Federated Learning with Non-IID DataabstractIn federated learning a global model is trained with training data geographically distributed over a number of clients. To reduce the communication cost over the expensive wide area network, clients complete multiple local iterations before synchronization. However, since the training data are non-iid, such infrequent synchronization would compromise the accuracy after model convergence. In order to tackle this problem, we propose Two-Dimensional Learning Rate Decay (2D-LRD) in this paper, which aims to improve the model performance by adaptively tuning the learning rate on two dimensions: round-dimension and iteration-dimension during the model training. That is, we gradually decrease the learning rate and decrease the learning rates of local iterations in a synchronization round with different speeds. Based on our experiments and analysis, we find that the sum of the inner product of round updates is a valuable signal for learning rate tuning. We perform evaluation and demonstrate that 2D-LRD can make great progress compared to the baseline scheme. Kaiwei Mo, Chen Chen 0067, Jiamin Li 0002, Hong Xu 0001, Chun Jason Xue |
IJCNN | 2 |
| 2020 | Semi-dynamic load balancing: efficient distributed learning in non-dedicated environmentsabstractMachine learning (ML) models are increasingly trained in clusters with non-dedicated workers possessing heterogeneous resources. In such scenarios, model training efficiency can be negatively affected by stragglers---workers that run much slower than others. Efficient model training requires eliminating such stragglers, yet for modern ML workloads, existing load balancing strategies are inefficient and even infeasible. In this paper, we propose a novel strategy called semi-dynamic load balancing to eliminate stragglers of distributed ML workloads. The key insight is that ML workers shall be load-balanced at iteration boundaries, being non-intrusive to intra-iteration execution. We develop LB-BSP based on such an insight, which is an integrated worker coordination mechanism that adapts workers' load to their instantaneous processing capabilities by right-sizing the sample batches at the synchronization barriers. We have custom-designed the batch sizing algorithm respectively for CPU and GPU clusters based on their own characteristics. LB-BSP has been implemented as a Python module for ML frameworks like TensorFlow and PyTorch. Our EC2 deployment confirms that LB-BSP is practical, effective and light-weight, and is able to accelerating distributed training by up to 54%. Chen Chen 0067, Qizhen Weng 0001, Wei Wang 0030, Baochun Li, Bo Li 0001 |
SoCC | 1 |
| 2020 | Metis: learning to schedule long-running applications in shared container clusters at scaleabstractOnline cloud services are increasingly deployed as long-running applications (LRAs) in containers. Placing LRA containers is known to be difficult as they often have sophisticated resource interferences and I/O dependencies. Existing schedulers rely on operators to manually express the container scheduling requirements as placement constraints and strive to satisfy as many constraints as possible. Such schedulers, however, fall short in performance as placement constraints only provide qualitative scheduling guidelines and minimizing constraint violations does not necessarily result in the optimal performance.In this work, we present Metis, a general-purpose scheduler that learns to optimally place LRA containers using deep reinforcement learning (RL) techniques. This eliminates the complex manual specification of placement constraints and offers, for the first time, concrete quantitative scheduling criteria. As directly training an RL agent does not scale, we develop a novel hierarchical learning technique that decomposes a complex container placement problem into a hierarchy of subproblems with significantly reduced state and action space. We show that many subproblems have similar structures and can hence be solved by training a unified RL agent offline. Large-scale EC2 deployment shows that compared with the traditional constraint-based schedulers, Metis improves the throughput by up to 61%, optimizes various performance metrics, and easily scales to a large cluster where 3K containers run on over 700 machines. Qizhen Weng 0001, Wei Wang 0030, Chen Chen 0067, Bo Li 0001 |
SC | 4 |
| 2019 | Round-Robin Synchronization: Mitigating Communication Bottlenecks in Parameter ServersabstractDeep learning is usually performed in GPU clusters where each worker machine iteratively refines the model parameters by communicating the update with the Parameter Server (PS). More often than not, workers communicate in a synchronous manner, so as to avoid using out-of-dated parameters and make high-quality refinement in each iteration. However, as all workers synchronize with the PS simultaneously, the communication becomes a severe bottleneck. To address this problem, in this paper we propose the Round-Robin Synchronous Parallel (R2SP) scheme, which coordinates workers to make updates in an evenly-gapped, round-robin manner. This way, R2SP can minimize the network contention at a minimum cost of the refinement quality. We further extend R2SP to heterogeneous clusters by adaptively tuning the batch size of each worker based on its processing capability. We have implemented R2SP as a ready-to-use python library for status-quo deep learning frameworks. EC2 deployment in GPU clusters show that R2SP effectively mitigates the communication bottlenecks, accelerating the training of popular image classification models by up to 25%. Chen Chen 0067, Wei Wang 0030, Bo Li 0001 |
INFOCOM | 1 |
| 2018 | Fast Distributed Deep Learning via Worker-adaptive Batch SizingabstractIn heterogeneous or shared clusters, distributed learning processes are slowed down by straggling workers. In this work, we propose LB-BSP, a new synchronization scheme that eliminates stragglers by adapting each worker's training load (batch size) to its processing capability. For training in shared production clusters, a prerequisite for deciding the workers' batch sizes is to know their processing speeds before each iteration starts. To this end, we adopt NARX, an extended recurrent neural network that accounts for both the historical speeds and the driving factors such as CPU and memory in prediction. Chen Chen 0067, Qizhen Weng 0001, Wei Wang 0030, Baochun Li, Bo Li 0001 |
SoCC | 1 |
| 2018 | FlowTime: Dynamic Scheduling of Deadline-Aware Workflows and Ad-Hoc JobsabstractWith rapidly increasing volumes of data to be processed in modern data analytics, it is commonplace to run multiple data processing jobs with inter-job dependencies in a datacenter cluster, typically as recurring data processing workloads. Such a group of inter-dependent data analytic jobs is referred to as a workflow, and may have a deadline due to its mission-critical nature. In contrast, non-recurring ad-hoc jobs are typically best-effort in nature, and rather than meeting deadlines, it is desirable to minimize their average job turnaround time. The state-of-the-art scheduling mechanisms focused on meeting deadlines for individual jobs only, and are oblivious to workflow deadlines. In this paper, we present FlowTime, a new system framework designed to make scheduling decisions for workflows so that their deadlines are met, while simultaneously optimizing the performance of ad-hoc jobs. To achieve this objective, we first adopt a divide-and-conquer strategy to transform the problem of workflow scheduling to a deadline-aware job scheduling problem, and then design an efficient algorithm that tackles the scheduling problem with both deadline-aware jobs and ad-hoc jobs by solving its corresponding optimization problem directly using a linear program solver. Our experimental results have clearly demonstrated that FlowTime achieves the lowest deadline-miss rates for deadline-aware workflows and 2-10 times shorter average job turnaround time, as compared to the state-of-the-art scheduling algorithms. Zhiming Hu 0001, Baochun Li, Chen Chen 0067, Xiaodi Ke |
ICDCS | 3 |
| 2018 | Performance-Aware Fair Scheduling: Exploiting Demand Elasticity of Data Analytics JobsabstractEfficient resource management is of paramount importance in today's production clusters. In this paper, we identify the demand elasticity of data-parallel jobs. Demand elasticity allows jobs to run with a significantly less amount of resources than they ideally need, at the expense of only a modest performance penalty. Our EC2 experiment using popular Spark benchmark suites confirms that running a job using 50% of demanded slots is sufficient to achieve at least 75% of the ideal performance. We show that such an elasticity is an intrinsic property of data-parallel jobs and can be exploited to speed up average job completion. In this regard, we propose Performance-Aware Fair (PAF) scheduler to identify the demand elasticity and use it to improve the average job performance, while still attaining near-optimal isolation guarantee close to fair sharing. PAF starts with a fair allocation and iteratively adjusts it by transferring resources from one job to another, improving the performance of resource-taker without penalizing resource-giver by a noticeable amount. We implemented PAF in Spark and evaluated its effectiveness through both EC2 experiments and large-scale simulations. Evaluation results show that compared with fair allocation, PAF improves the average job performance by 13%, while penalizing resource-givers by no more than 1%. Chen Chen 0067, Wei Wang 0030, Bo Li 0001 |
INFOCOM | 1 |
| 2017 | Speculative Slot Reservation: Enforcing Service Isolation for Dependent Data-Parallel ComputationsabstractPriority scheduling is a fundamental tool to provide service isolation for different jobs in shared clusters. Ideally, the performance of a high-priority job should not be dragged down by another with a lower priority. However, we show in this paper that simply assigning a high priority provides no isolation for jobs with dependent computations. A job, even receiving the highest priority, may give up compute slots to another before proceeding to the downstream computation, which is because of barrier, i.e., that the downstream computation cannot start until all the upstream tasks have completed. Such an interruption of execution inevitably results in a significant delay. In this paper, we propose speculative slot reservation that judiciously reserves slots for downstream computations, so as to retain service isolation for high-priority jobs. To mitigate the utilization loss due to slot reservation, we analyze the trade-off between utilization and isolation, and expose a tunable knob to navigate the trade-off. We also propose a complementary straggler mitigation strategy that uses the reserved slots to run extra copies of slow tasks. We have implemented speculative slot reservation in Spark. Evaluations based on both cluster deployment and trace-driven simulations show that our approach enforces strict service isolation for high-priority jobs, without slowing down the other jobs with a lower priority. Chen Chen 0067, Wei Wang 0030, Bo Li 0001 |
ICDCS | 1 |
| 2017 | Cluster fair queueing: Speeding up data-parallel jobs with delay guaranteesabstractCluster scheduler serves as a critical component to data-parallel systems in datacenters. Ideally, a scheduler should provide predictable performance with guarantees on the maximal job completion delay, while at the same time ensuring the minimal mean response time. Practically however, performance predictability and optimality are often conflicting with each other. The results often are a plethora of scheduling policies that either achieve predictable performance at the expense of long response times (e.g., max-min fairness), or run the risk of starving some jobs to obtain the minimal mean response time (e.g., Shortest Remaining Processing Time First). To address these problems, we develop a new scheduler, Cluster Fair Queueing (CFQ), which preferentially offers resources to jobs that complete the earliest under a fair sharing policy. We show that CFQ is able to minimize the mean response time while at the same time ensuring jobs to finish within a constant time after their completion under fair sharing. Our Spark deployment on a 100-node EC2 cluster demonstrates that compared to the built-in fair scheduler, CFQ can decrease the mean response time by 40%, which speeds up more than 40% of jobs by over 75% on average. Chen Chen 0067, Wei Wang 0030, Shengkai Zhang, Bo Li 0001 |
INFOCOM | 1 |
| 2016 | Software-defined inter-domain routing revisitedabstractThe decoupling of control and data plane in software-defined networking (SDN) has been shown to be promising to improve routing performance in the context of intradomain routing. The applicability of SDN in inter-domain routing, especially with respect to route convergence, has not been properly explored. In this work, we propose a mathematical model to quantify the BGP convergence time for inter-domain routing by capturing only the essential components in BGP convergence process. Based on the model and some practical observations, we study how SDN may help to facilitate the interdomain routing. We further present a greedy algorithm that selects Autonomous Systems (ASes) for incremental SDN deployment with the objectives of minimizing the BGP convergence time. The simulation result based on the real world Internet topology confirms the effectiveness of our proposed algorithm. Chen Chen 0067, Bo Li 0001, Dong Lin, Baochun Li |
ICC | 1 |
| 2010 | A Study of a Software Cache Implementation of the OpenMP Memory Model for Multicore and Manycore Architectures
Chen Chen 0067, Joseph B. Manzano, Ge Gan, Guang R. Gao, Vivek Sarkar |
Euro-Par (2) | 1 |