EDBT 2026 Demo / reviewers in the wild / expert
Yuan Yuan 0014
dblp:64/5845-14
· DBLP profile ↗
28ranked-venue papers
8as first author
16since 2021 · last 2026
—ORCID · conflict
Domains — the database's venue-derived domains; a paper can count in several
Databases, data management, data science and information retrieval · 10 · 3 first-author · 2 since 2021Systems, architecture and hardware · 9 · 3 first-author · 6 since 2021Artificial intelligence and machine learning · 3 · 1 first-author · 2 since 2021Applied, interdisciplinary, general and emerging computing · 3 · 2 first-author · 1 since 2021Computer networks · 2 · 1 first-author · 2 since 2021Theory of computation · 2 · 2 since 2021Security and privacy · 1 · 1 since 2021
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | PARS: Optimizing High-Dimensional Vector Search with Dimensional Reduction and Ray Tracing Core
Yiling Ma, Mengbai Xiao, Zhixiong Xiao, Yuan Yuan 0014, Dongxiao Yu, Feng Li 0002 |
ICDCS | 4 |
| 2025 | PDUDT: Provable Decentralized Unlearning under Dynamic TopologiesabstractThis paper investigates decentralized unlearning, aiming to eliminate the impact of a specific client on the whole decentralized system. However, decentralized communication characterizations pose new challenges for effective unlearning: the indirect connections make it difficult to trace the specific client's impact, while the dynamic topology limits the scalability of retraining-based unlearning methods.
In this paper, we propose the first **P**rovable **D**ecentralized **U**nlearning algorithm under **D**ynamic **T**opologies called PDUDT. It allows clients to eliminate the influence of a specific client without additional communication or retraining. We provide rigorous theoretical guarantees for PDUDT, showing it is statistically indistinguishable from perturbed retraining. Additionally, it achieves an efficient convergence rate of $\mathcal{O}(\frac{1}{T})$ in subsequent learning, where $T$ is the total communication rounds. This rate matches state-of-the-art results. Experimental results show that compared with the Retrain method, PDUDT saves more than 99\% of unlearning time while achieving comparable unlearning performance. Jing Qiao, Yu Liu 0085, Zengzhe Chen, Yuan Yuan 0014, Xiao Zhang 0015, Dongxiao Yu |
ICML | 5 |
| 2025 | Memory-Enhanced Invariant Prompt Learning for Urban Flow Prediction Under Distribution Shifts
Haiyang Jiang 0017, Tong Chen 0005, Wentao Zhang 0001, Nguyen Quoc Viet Hung, Yuan Yuan 0014, Yong Li 0008, Hongzhi Yin |
ECML/PKDD (3) | 5 |
| 2025 | Pruning-Based Adaptive Federated Learning at the EdgeabstractFederated Learning (FL) is a new learning framework in which$s$clients collaboratively train a model under the guidance of a central server. Meanwhile, with the advent of the era of large models, the parameters of models are facing explosive growth. Therefore, it is important to design federated learning algorithms for edge environment. However, the edge environment is severely limited in computing, storage, and network bandwidth resources. Concurrently, adaptive gradient methods show better performance than constant learning rate in non-distributed settings. In this paper, we propose a pruning-based distributed Adam (PD-Adam) algorithm, which combines model pruning and adaptive learning steps to achieve asymptotically optimal convergence rate of$O(1/\sqrt[4]{K})$. At the same time, the algorithm can achieve convergence consistent with the centralized model. Finally, extensive experiments have confirmed the convergence of our algorithm, demonstrating its reliability and effectiveness across various scenarios. Specially, our proposed algorithm is$2$% and$18$% more accurate than the current state-of-the-art FedAvg algorithm on the ResNet and CIFAR datasets. Dongxiao Yu, Yuan Yuan 0014, Yifei Zou, Xiao Zhang 0015, Yu Liu 0085, Li-Zhen Cui 0001, Xiuzhen Cheng |
IEEE Trans. Computers | 2 |
| 2025 | BaDFL: Mitigating Model Poisoning in Decentralized Federated LearningabstractDecentralized federated learning (DFL) has gained significant attention due to its ability to facilitate collaborative model training without relying on a central server. However, it is highly vulnerable to backdoor attacks, where malicious participants can manipulate model updates to embed hidden functionalities. In this paper, we propose BaDFL, a novel Backdoor Attack defense mechanism for Decentralized Federated Learning. BaDFL enhances robustness by applying strategic model clipping at the local update level. To the best of our knowledge, BaDFL is the first decentralized federated learning algorithm with theoretical guarantees against model poisoning attacks. Specifically, BaDFL achieves an asymptotically optimal convergence rate of$O(\frac{1}{\sqrt{nT}})$, wherenis the number of nodes andTis the global maximum iteration number. Furthermore, we provide a comprehensive analysis under two different attack scenarios, showing that BaDFL maintains robustness within a specific defense radius. Extensive experimental results show that, on average, BaDFL can effectively defend against model poisoning within 6 mitigation rounds, with less than a 1% drop in accuracy. Yuan Yuan 0014, Anhao Zhou, Xiao Zhang 0015, Yifei Zou, Yangguang Shi, Dongxiao Yu |
IEEE Trans. Computers | 1 |
| 2025 | Adaptive pruning-based Newton's method for distributed learning
Shuzhen Chen 0001, Yuan Yuan 0014, Youming Tao 0001, Tianzhu Wang, Zhipeng Cai 0001, Dongxiao Yu |
Theor. Comput. Sci. | 2 |
| 2024 | Fed-MS: Fault Tolerant Federated Edge Learning with Multiple Byzantine ServersabstractDue to its decentralized framework and outdoor environments, federated edge learning (FEEL) faces significant vulnerability to malicious attacks within edge networks. Prevailing FEEL approaches typically hinge on a dependable parameter server (PS) to contend with the adversarial updates from Byzantine clients. Recognizing the inherent unreliability of PSs in edge networks, this paper delves into the security challenges of FEEL, specifically addressing Byzantine PSs. We present a Byzantine fault-tolerant FEEL algorithm, named Fed-MS, in which a multi-server technique along with a newly designed trimmed-mean-based model filter is employed. This combination ensures that each client can obtain a feasible global model for its local training, closely approximating a true model aggregated by benign PSs. Furthermore, we propose a sparse uploading strategy in Fed-MS to enhance communication efficiency for model aggregation to multiple PSs. Theoretical analysis demonstrates that, when Byzantine PSs are a minority, Fed-MS achieves an expected convergence speed of$O(1/T)$with$T$defined as the number of training rounds, akin to state-of-the-art works under non-Byzantine settings. Extensive experiments are conducted on the CIFAR-10 dataset with MobileNet V2 as the training model. The numerical results show that our Fed-MS can improve the model accuracy from 10% to at least 76% under the malicious attacks from Byzantine PSs. Our code is released at https://github.com/haoma2772/Fed-MS. Senmao Qi, Yifei Zou, Yuan Yuan 0014, Peng Li 0017, Dongxiao Yu |
ICDCS | 4 |
| 2024 | BR-DeFedRL: Byzantine-Robust Decentralized Federated Reinforcement Learning with Fast Convergence and Communication EfficiencyabstractIn this paper, we propose Byzantine-Robust Decentralized Federated Reinforcement Learning (BR-DeFedRL), an innovative framework that effectively combats the harmful influence of Byzantine agents by adaptively adjusting communication weights, thereby significantly enhancing the robustness of the learning system. By leveraging decentralized learning, our approach eliminates the dependence on a central server. Striking a harmonious balance between communication round count and sample complexity, BR-DeFedRL achieves efficient convergence with a rate of $\mathcal{O}\left( {\frac{1}{{TN}}} \right)$, where T denotes the communication rounds and N represents the local steps related to variance reduction. Notably, each agent attains an ϵ-approximation with a state-of-the-art sample complexity of $\mathcal{O}\left( {\frac{1}{{\varepsilon N}} + \frac{1}{\varepsilon }} \right)$. Extensive experimental validations further affirm the efficacy of BR-DeFedRL, making it a promising and practical solution for Byzantine-robust decentralized federated reinforcement learning. Jing Qiao, Zuyuan Zhang, Sheng Yue 0001, Yuan Yuan 0014, Zhipeng Cai 0001, Xiao Zhang 0015, Ju Ren 0001, Dongxiao Yu |
INFOCOM | 4 |
| 2024 | Fed-MoE: Efficient Federated Learning for Mixture-of-Experts Models via Empirical Pruning
Yifei Zou, Senmao Qi, Yuan Yuan 0014, Dawei Wang 0007, Shikun Shen, Shao-Yong Guo 0001, Dongxiao Yu |
PDCAT | 3 |
| 2024 | BR-FEEL: A backdoor resilient approach for federated edge learning with fragment-sharing
Senmao Qi, Yifei Zou, Yuan Yuan 0014, Peng Li 0017, Dongxiao Yu |
J. Syst. Archit. | 4 |
| 2024 | Distributed Learning for Large-Scale Models at Edge With Privacy ProtectionabstractBig data and strong computing power have promoted artificial intelligence to the era of big models. In particular, ChatGPT’s debut heralded the vigorous development of large models. It is an urgent problem to train large models with trillion-level parameters efficiently. Traditional single-machine training stores all data and model parameters in memory. However, due to the limitation of memory and communication resources, when the amount of data or model parameters increases, the problem of memory shortage and communication blocking often occurs. Therefore, distributed training is the most effective ways to solve the above problems and improve training efficiency. In this paper, we propose the algorithmDL-DP, which can achieve an asymptotically optimal convergence rate$O(1/{\sqrt{TK\Gamma^*}})$while satisfyingε-differential privacy, whereTis the local epoch number,Kis the global maximum iteration number and$\Gamma^*$is the minimum covering index. In particular, when${\Gamma ^*} = N$, DL-DP achieves a convergence rate of$O(1/{\sqrt{TKN}})$, which is equivalent to the best-known FedAvg approach implemented by training the full model at each client. When${\Gamma ^*} = 1$, DL-DP achieves a convergence rate of$O(1/{\sqrt{TK}})$, which is comparable to OAP that assumes all parameters need to be trained at least once in each iteration. Finally, our algorithm has been demonstrated to converge through extensive experiments. Yuan Yuan 0014, Shuzhen Chen 0001, Dongxiao Yu, Zengrui Zhao, Yifei Zou, Li-Zhen Cui 0001, Xiuzhen Cheng |
IEEE Trans. Computers | 1 |
| 2023 | Resource-Adaptive Newton's Method for Distributed Learning
Shuzhen Chen 0001, Yuan Yuan 0014, Youming Tao 0001, Zhipeng Cai 0001, Dongxiao Yu |
COCOON (1) | 2 |
| 2023 | Trustworthy decentralized collaborative learning for edge intelligence: A surveyabstractEdge intelligence is an emerging technology that enables artificial intelligence on connected systems and devices in close proximity to the data sources. Decentralized Collaborative Learning (DCL) is a novel edge intelligence technique that allows distributed clients to cooperatively train a global learning model without revealing their data. DCL has a wide range of applications in various domains, such as smart city and autonomous driving. However, DCL faces significant challenges in ensuring its trustworthiness, as data isolation and privacy issues make DCL systems vulnerable to adversarial attacks that aim to breach system confidentiality, undermine learning reliability or violate data privacy. Therefore, it is crucial to design DCL in a trustworthy manner, with a focus on security, robustness, and privacy. In this survey, we present a comprehensive review of existing efforts for designing trustworthy DCL systems from the three key aformentioned aspects: security, robustness, and privacy. We analyze the threats that affect the trustworthiness of DCL across different scenarios and assess specific technical solutions for achieving each aspect of Trustworthy DCL (TDCL). Finally, we highlight open challenges and future directions for advancing TDCL research and practice. Dongxiao Yu, Zhenzhen Xie 0002, Yuan Yuan 0014, Shuzhen Chen 0001, Jing Qiao, Yong Yu 0002, Yifei Zou, Xiao Zhang 0015 |
High Confid. Comput. | 3 |
| 2023 | Decentralized Parallel SGD Based on Weight-Balancing for Intelligent IoVabstractTraining machine learning models in a decentralized way has attracted tremendous attention on intelligent Internet of Vehicles (IIoV). However, it is highly dynamic and asymmetric for the connections between vehicles in IIoV due to the mobility of vehicles and the complex communication environment, which poses great challenges on designing efficient distributed learning algorithms. To address this problem, we focus on the basic stochastic gradient descent (SGD) algorithm and propose a decentralized parallel SGD algorithm (DPSGD-WB) for the complex IIoV. The algorithm is based on weight-balancing to overcome the difficulty caused by the dynamic and asymmetric connectivity in IIoV. With rigorous analysis, we show that DPSGD-WB converges on the optimal rate of$O(1/\sqrt {Kn})$, where$n$is the number of vehicle terminals and$K$is the number of iterations. To the best of our knowledge, our proposed algorithm is the first known decentralized parallel SGD algorithm that can be implemented in asymmetric and dynamic intelligent IoV systems. Finally, extensive experiments demonstrate the efficacy of our algorithm. Yuan Yuan 0014, Jiguo Yu, Xiaolu Cheng, Zongrui Zou, Dongxiao Yu, Zhipeng Cai 0001 |
IEEE Trans. Intell. Transp. Syst. | 1 |
| 2021 | NestGPU: Nested Query Processing on GPUabstractNested queries are commonly used to express complex use-cases by connecting the output of a subquery as an input to the outer query block. However, their execution is highly time-consuming. Researchers have proposed various algorithms and techniques that unnest subqueries to improve performance. Since this is a customized approach that needs high algorithmic and engineering efforts, it is largely not an open feature in most existing database systems.Our approach is general-purpose and GPU-acceleration based, aiming for high performance at a minimum development cost. We look into the major differences between nested and unnested query structures to identify their merits and limits for GPU processing. Furthermore, we focus on the nested approach that is algorithmically simple and rich in parallels, in relatively low space complexity, and generic in program structure. We create a new code generation framework that best fits GPU for the nested method. We also make several critical system optimizations including massive parallel scanning with indexing, effective vectorization to optimize join operations, exploiting cache locality for loops and efficient GPU memory management. We have implemented the proposed solutions in NestGPU, a GPU-based column-store database system that is GPU device independent. We have extensively evaluated and tested the system to show the effectiveness of our proposed methods. Sofoklis Floratos, Mengbai Xiao, Hao Wang 0002, Chengxin Guo, Yuan Yuan 0014, Rubao Lee, Xiaodong Zhang 0001 |
ICDE | 5 |
| 2021 | D-(DP)2SGD: Decentralized Parallel SGD with Differential Privacy in Dynamic NetworksabstractDecentralized machine learning has been playing an essential role in improving training efficiency. It has been applied in many real‐world scenarios, such as edge computing and IoT. However, in fact, networks are dynamic, and there is a risk of information leaking during the communication process. To address this problem, we propose a decentralized parallel stochastic gradient descent algorithm (D‐(DP)2SGD) with differential privacy in dynamic networks. With rigorous analysis, we show that D‐(DP)2SGD converges with a rate of while satisfying ε‐DP, which achieves almost the same convergence rate as previous works without privacy concern. To the best of our knowledge, our algorithm is the first known decentralized parallel SGD algorithm that can implement in dynamic networks and take privacy‐preserving into consideration. Yuan Yuan 0014, Zongrui Zou, Dongxiao Yu |
Wirel. Commun. Mob. Comput. | 1 |
| 2019 | Fast Fault-Tolerant Sampling via Random Walk in Dynamic NetworksabstractWe study the fundamental problem of fault-tolerant distributed sampling towards uniform probabilistic distribution in dynamic multi-hop wireless networks. Whereas uniform sampling has been extensively studied without concerning fault tolerance, only quite few proposals investigate how the uniform sampling algorithm tolerate Byzantine faults on dynamic networks with very special topologies, e.g., regular graphs with constant node degree. Therefore, designing fault-tolerant uniform sampling algorithms for more general graphs is still an open problem. To this end, we propose a fast and highly fault-tolerate randomized algorithm, such that nearly-uniform sampling is achieved in O(log2n) rounds, while up to O(√n/(polylog(n)·Δ)) Byzantine nodes can be tolerated, where Δ is the maximum degree of the network. Moreover, the proposed algorithm is also communication efficient in the sense that only O(log n) bits need to be exchanged on each link in every round. To show the power of distributed uniform sampling, we apply the proposed algorithm in designing polylogarithmic time distributed algorithms for two typical fundamental issues, i.e., to achieve agreement or data aggregation in Byzantine dynamic networks. Yuan Yuan 0014, Feng Li 0002, Dongxiao Yu, Jiguo Yu, Yu Wu 0010, Weifeng Lv, Xiuzhen Cheng |
ICDCS | 1 |
| 2018 | SQLoop: High Performance Iterative Processing in Data ManagementabstractIncreasingly more iterative and recursive query tasks are processed in data management systems, such as graph-structured data analytics, demanding fast response time. However, existing CTE-based recursive SQL and its implementation ineffectively respond to this intensive query processing with two major drawbacks. First, its iteration execution model is based on implicit set-oriented terminating conditions that cannot express aggregation-based tasks, such as PageRank. Second, its synchronous execution model cannot perform asynchronous computing to further accelerate execution in parallel. To address these two issues, we have designed and implemented SQLoop, a framework that extends the semantics of current SQL standard in order to accommodate iterative SQL queries. SQLoop interfaces between users and different database engines with two powerful components. First, it provides an uniform SQL expression for users to access any database engine so that they do not need to write database dependent SQL or move datasets from a target engine to process in their own sites. Second, SQLoop automatically parallelizes iterative queries that contain certain aggregate functions in both synchronous and asynchronous ways. More specifically, SQLoop is able to take advantage of intermediate results generated between different iterations and to prioritize the execution of partitions that accelerate the query processing. We have tested and evaluated SQLoop by using three popular database engines with real-world datasets and queries, and shown its effectiveness and high performance. Sofoklis Floratos, Yanfeng Zhang 0001, Yuan Yuan 0014, Rubao Lee, Xiaodong Zhang 0001 |
ICDCS | 3 |
| 2017 | Feisu: Fast Query Execution over Heterogeneous Data Sources on Large-Scale ClustersabstractFast data analytics at an increasingly large scale has become a critical task in any Internet service company. For example, in Baidu, the major search engine company in China, large volumes of Web and business data in PB-scale are timely and constantly acquired and analyzed for the purposes of evaluating product revenue, tracking product demanding activities on market, predicting user behavior, upgrading product rankings, and diagnosing spam cases, and many others. Response time for queries of various data analytics not only affects user experiences, but also has a serious impact on productivity of business operations. In this paper, to meet the challenge of fast data analytics, we present Feisu (meaning fast in Chinese), a data integration system over heterogeneous storage systems, which has been widely used in Baidu's critical and daily business analytics applications after our R&D efforts. Feisu is designed and implemented to co-work together with several heterogeneous storage systems, and exploit the query similarity embedded in complex query workloads. Our experiments using real world workloads show that Feisu can significantly improve query performance in Baidu. Feisu has been in production use in Baidu for two years to effectively manage over dozens of petabytes of data for various applications. An Qin 0001, Yuan Yuan 0014, Dai Tan, Rubao Lee, Xiaodong Zhang 0001 |
ICDE | 2 |
| 2017 | A distributed in-memory key-value store system on heterogeneous CPU-GPU cluster
Kai Zhang 0006, Kaibo Wang, Yuan Yuan 0014, Lei Guo 0004, Rubao Li, Xiaodong Zhang 0001, Bingsheng He, Bei Hua |
VLDB J. | 3 |
| 2016 | Spark-GPU: An accelerated in-memory data processing engine on clustersabstractApache Spark is an in-memory data processing system that supports both SQL queries and advanced analytics over large data sets. In this paper, we present our design and implementation of Spark-GPU that enables Spark to utilize GPU's massively parallel processing ability to achieve both high performance and high throughput. Spark-GPU transforms a general-purpose data processing system into a GPU-supported system by addressing several real-world technical challenges including minimizing internal and external data transfers, preparing a suitable data format and a batching mode for efficient GPU execution, and determining the suitability of workloads for GPU with a task scheduling capability between CPU and GPU. We have comprehensively evaluated Spark-GPU with a set of representative analytical workloads to show its effectiveness. Our results show that Spark-GPU improves the performance of machine learning workloads by up to 16.13x and the performance of SQL queries by up to 4.83x. Yuan Yuan 0014, Meisam Fathi Salmi, Yin Huai, Kaibo Wang, Rubao Lee, Xiaodong Zhang 0001 |
IEEE BigData | 1 |
| 2016 | BCC: Reducing False Aborts in Optimistic Concurrency Control with Low Cost for In-Memory DatabasesabstractThe Optimistic Concurrency Control (OCC) method has been commonly used for in-memory databases to ensure transaction serializability --- a transaction will be aborted if its read set has been changed during execution. This simple criterion to abort transactions causes a large proportion of false positives, leading to excessive transaction aborts. Transactions aborted false-positively (i.e. false aborts) waste system resources and can significantly degrade system throughput (as much as 3.68x based on our experiments) when data contention is intensive. Modern in-memory databases run on systems with increasingly parallel hardware and handle workloads with growing concurrency. They must efficiently deal with data contention in the presence of greater concurrency by minimizing false aborts. This paper presents a new concurrency control method named Balanced Concurrency Control (BCC) which aborts transactions more carefully than OCC does. BCC detects data dependency patterns which can more reliably indicate unserializable transactions than the criterion used in OCC. The paper studies the design options and implementation techniques that can effectively detect data contention by identifying dependency patterns with low overhead. To test the performance of BCC, we have implemented it in Silo and compared its performance against that of the vanilla Silo system with OCC and two-phase locking (2PL). Our extensive experiments with TPC-W-like, TPC-C-like and YCSB workloads demonstrate that when data contention is intensive, BCC can increase transaction throughput by more than 3x versus OCC and more than 2x versus 2PL; meanwhile, BCC has comparable performance with OCC for workloads with low data contention. Yuan Yuan 0014, Kaibo Wang, Rubao Lee, Xiaoning Ding, Spyros Blanas, Xiaodong Zhang 0001 |
Proc. VLDB Endow. | 1 |
| 2015 | SideWalk: A Facility of Lightweight Out-of-Band Communications for Augmenting Distributed Data Processing FlowsabstractThe foundation of a data processing engine running on a large cluster is its programming model that defines data processing operations and data movements. A special kind of communication activities that are not normally defined in the programming model but are often used in ad hoc ways in system development, is called out-of-band communications. The existing ad hoc solutions of out-of-band communications are often hard to reuse, error-prone, and not free from unwanted side effects. To address these issues, we have designed and implemented a standalone facility of out-of-band communications called SideWalk. With this facility, users can add out-of-band communication operations into their distributed data flows through a set of reusable APIs. These APIs have well defined semantics and thus, users' chances of writing error-prone programs with SideWalk are minimized. To prevent users from introducing unwanted side effects while using SideWalk, we prototype SideWalk to efficiently handle lightweight out-of-band communications and we restrict communication patterns that can be conducted through SideWalk without affecting the applicability of SideWalk on typical use cases. Our experimental results show that execution times of distributed data processing flows in a Hadoop environment with out-of-band communications implemented with SideWalk are reduced up to 1.53 times compared with that of distributed data processing flows with out-of-band communications implemented with a representative ad hoc solution. Yin Huai, Yuan Yuan 0014, Rubao Lee, Xiaodong Zhang 0001 |
CLUSTER | 2 |
| 2015 | Hetero-DB: Next Generation High-Performance Database Systems by Best Utilizing Heterogeneous Computing and Storage Resources
Kai Zhang 0006, Feng Chen 0005, Xiaoning Ding, Yin Huai, Rubao Lee, Kaibo Wang, Yuan Yuan 0014, Xiaodong Zhang 0001 |
J. Comput. Sci. Technol. | 8 |
| 2015 | Mega-KV: A Case for GPUs to Maximize the Throughput of In-Memory Key-Value StoresabstractIn-memory key-value stores play a critical role in data processing to provide high throughput and low latency data accesses. In-memory key-value stores have several unique properties that include (1) data intensive operations demanding high memory bandwidth for fast data accesses, (2) high data parallelism and simple computing operations demanding many slim parallel computing units, and (3) a large working set. As data volume continues to increase, our experiments show that conventional and general-purpose multicore systems are increasingly mismatched to the special properties of key-value stores because they do not provide massive data parallelism and high memory bandwidth; the powerful but the limited number of computing cores do not satisfy the demand of the unique data processing task; and the cache hierarchy may not well benefit to the large working set. In this paper, we make a strong case for GPUs to serve as special-purpose devices to greatly accelerate the operations of in-memory key-value stores. Specifically, we present the design and implementation of Mega-KV, a GPU-based in-memory key-value store system that achieves high performance and high throughput. Effectively utilizing the high memory bandwidth and latency hiding capability of GPUs, Mega-KV provides fast data accesses and significantly boosts overall performance. Running on a commodity PC installed with two CPUs and two GPUs, Mega-KV can process up to 160+ million key-value operations per second, which is 1.4-2.8 times as fast as the state-of-the-art key-value store system on a conventional CPU-based platform. Kai Zhang 0006, Kaibo Wang, Yuan Yuan 0014, Lei Guo 0004, Rubao Lee, Xiaodong Zhang 0001 |
Proc. VLDB Endow. | 3 |
| 2014 | Major technical advancements in apache hiveabstractApache Hive is a widely used data warehouse system for Apache Hadoop, and has been adopted by many organizations for various big data analytics applications. Closely working with many users and organizations, we have identified several shortcomings of Hive in its file formats, query planning, and query execution, which are key factors determining the performance of Hive. In order to make Hive continuously satisfy the requests and requirements of processing increasingly high volumes data in a scalable and efficient way, we have set two goals related to storage and runtime performance in our efforts on advancing Hive. First, we aim to maximize the effective storage capacity and to accelerate data accesses to the data warehouse by updating the existing file formats. Second, we aim to significantly improve cluster resource utilization and runtime performance of Hive by developing a highly optimized query planner and a highly efficient query execution engine. In this paper, we present a community-based effort on technical advancements in Hive. Our performance evaluation shows that these advancements provide significant improvements on storage efficiency and query execution performance. This paper also shows how academic research lays a foundation for Hive to improve its daily operations. Yin Huai, Ashutosh Chauhan, Alan Gates, Günther Hagleitner, Eric N. Hanson, Owen O'Malley, Jitendra Pandey, Yuan Yuan 0014, Rubao Lee, Xiaodong Zhang 0001 |
SIGMOD Conference | 8 |
| 2014 | Concurrent Analytical Query Processing with GPUsabstractIn current databases, GPUs are used as dedicated accelerators to process each individual query. Sharing GPUs among concurrent queries is not supported, causing serious resource underutilization. Based on the profiling of an open-source GPU query engine running commonly used single-query data warehousing workloads, we observe that the utilization of main GPU resources is only up to 25%. The underutilization leads to low system throughput. To address the problem, this paper proposes concurrent query execution as an effective solution. To efficiently share GPUs among concurrent queries for high throughput, the major challenge is to provide software support to control and resolve resource contention incurred by the sharing. Our solution relies on GPU query scheduling and device memory swapping policies to address this challenge. We have implemented a prototype system and evaluated it intensively. The experiment results confirm the effectiveness and performance advantage of our approach. By executing multiple GPU queries concurrently, system throughput can be improved by up to 55% compared with dedicated processing. Kaibo Wang, Kai Zhang 0006, Yuan Yuan 0014, Rubao Lee, Xiaoning Ding, Xiaodong Zhang 0001 |
Proc. VLDB Endow. | 3 |
| 2013 | The Yin and Yang of Processing Data Warehousing Queries on GPU DevicesabstractDatabase community has made significant research efforts to optimize query processing on GPUs in the past few years. However, we can hardly find that GPUs have been truly adopted in major warehousing production systems. Preparing to merge GPUs to the warehousing systems, we have identified and addressed several critical issues in a three-dimensional study of warehousing queries on GPUs by varying query characteristics, software techniques, and GPU hardware configurations. We also propose an analytical model to understand and predict the query performance on GPUs. Based on our study, we present our performance insights for warehousing query execution on GPUs. The objective of our work is to provide a comprehensive guidance for GPU architects, software system designers, and database practitioners to narrow the speed gap between the GPU kernel execution (the fast mode) and data transfer to prepare GPU execution (the slow mode) for high performance in processing data warehousing queries. The GPU query engine developed in this work is open source to the public. Yuan Yuan 0014, Rubao Lee, Xiaodong Zhang 0001 |
Proc. VLDB Endow. | 1 |