EDBT 2026 Demo / reviewers in the wild / expert
Yandong Wang 0001
dblp:54/6926-1
· DBLP profile ↗
19ranked-venue papers
8as first author
0since 2021 · last 2017
—ORCID · conflict
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 16 · 8 first-authorArtificial intelligence and machine learning · 2Databases, data management, data science and information retrieval · 2Computer networks · 1Software engineering, systems software and programming languages · 1Applied, interdisciplinary, general and emerging computing · 1
Expertise — from the expertise taxonomy: the topics of the expert's papers under the CCF categories. A weight counts papers with recency: 1 for a paper about the topic, 0.3 when the topic is its context, halved every five years.
| Computer architecture, parallel and distributed computing, and storage systems
8 papers |
Storage systems · 24% Memory systems · 21% Distributed systems · 18% | |
| Artificial intelligence
1 paper |
Efficient and distributed learning · 100% |
Topics — the 24 heaviest of 26, each with the papers that count most for it
| Topic | Weight | Papers | Last | Evidence papers |
|---|---|---|---|---|
Parallel and multicore computing › data-parallel programming
mapreduce |
0.7 | 4 | 2015 | Virtual Shuffling for Efficient Data Movement in MapReduce · IEEE Trans. Computers 2015 Design and Evaluation of Network-Levitated Merge for Hadoop Acceleration · IEEE Trans. Parallel Distributed Syst. 2014 CooMR: cross-task coordination for efficient data management in MapReduce programs · SC 2013 |
Distributed systems
fault tolerance |
0.5 | 2 | 2017 | GaDei: On Scale-Up Training as a Service for Deep Learning · ICDM 2017 HydraDB: a resilient RDMA-driven key-value middleware for in-memory cluster computing · SC 2015 |
Storage systems
key-value storage |
0.5 | 2 | 2016 | zExpander: a key-value cache with both high performance and fewer misses · EuroSys 2016 HydraDB: a resilient RDMA-driven key-value middleware for in-memory cluster computing · SC 2015 |
Cloud and datacenter computing
cluster resource management and scheduling |
0.4 | 2 | 2014 | Non-work-conserving effects in MapReduce: diffusion limit and criticality · SIGMETRICS 2014 CooMR: cross-task coordination for efficient data management in MapReduce programs · SC 2013 |
Distributed systems › distributed machine learning
parameter server |
0.3 | 1 | 2017 | GaDei: On Scale-Up Training as a Service for Deep Learning · ICDM 2017 |
Memory systems › memory compression
cache compression |
0.2 | 1 | 2016 | zExpander: a key-value cache with both high performance and fewer misses · EuroSys 2016 |
Memory systems
cache management |
0.2 | 1 | 2016 | zExpander: a key-value cache with both high performance and fewer misses · EuroSys 2016 |
Memory systems › cache management
cache replacement |
0.2 | 1 | 2016 | zExpander: a key-value cache with both high performance and fewer misses · EuroSys 2016 |
Memory systems › cache
key-value cache |
0.2 | 1 | 2016 | zExpander: a key-value cache with both high performance and fewer misses · EuroSys 2016 |
Distributed systems › distributed data processing
data shuffling |
0.2 | 1 | 2015 | Virtual Shuffling for Efficient Data Movement in MapReduce · IEEE Trans. Computers 2015 |
Storage systems › key-value storage
in-memory key-value store |
0.2 | 1 | 2015 | HydraDB: a resilient RDMA-driven key-value middleware for in-memory cluster computing · SC 2015 |
Storage systems
i/o optimization |
0.2 | 1 | 2015 | Virtual Shuffling for Efficient Data Movement in MapReduce · IEEE Trans. Computers 2015 |
Storage systems
storage reliability |
0.2 | 1 | 2015 | Virtual Shuffling for Efficient Data Movement in MapReduce · IEEE Trans. Computers 2015 |
Cloud and datacenter computing › cluster resource management and scheduling › cluster scheduling
mapreduce scheduling |
0.2 | 1 | 2014 | Non-work-conserving effects in MapReduce: diffusion limit and criticality · SIGMETRICS 2014 |
Embedded and real-time systems › real-time scheduling
non-work-conserving scheduling |
0.2 | 1 | 2014 | Non-work-conserving effects in MapReduce: diffusion limit and criticality · SIGMETRICS 2014 |
Cloud and datacenter computing › datacenter storage
i/o consolidation |
0.2 | 1 | 2013 | CooMR: cross-task coordination for efficient data management in MapReduce programs · SC 2013 |
Memory systems
data movement |
0.1 | 1 | 2011 | Hadoop acceleration through network levitated merge · SC 2011 |
Machine learning › Efficient and distributed learning
distributed training |
0.1 | 1 | 2017 | GaDei: On Scale-Up Training as a Service for Deep Learning · ICDM 2017 |
High-performance computing
cluster computing |
0.1 | 1 | 2015 | HydraDB: a resilient RDMA-driven key-value middleware for in-memory cluster computing · SC 2015 |
Cloud and datacenter computing › cluster computing framework
in-memory cluster computing |
0.1 | 1 | 2015 | HydraDB: a resilient RDMA-driven key-value middleware for in-memory cluster computing · SC 2015 |
Interconnection networks and networks-on-chip
remote direct memory access |
0.1 | 1 | 2015 | HydraDB: a resilient RDMA-driven key-value middleware for in-memory cluster computing · SC 2015 |
Cloud and datacenter computing
data movement acceleration |
0.1 | 1 | 2014 | Design and Evaluation of Network-Levitated Merge for Hadoop Acceleration · IEEE Trans. Parallel Distributed Syst. 2014 |
Performance modeling and evaluation
queueing analysis |
0.1 | 1 | 2014 | Non-work-conserving effects in MapReduce: diffusion limit and criticality · SIGMETRICS 2014 |
Interconnection networks and networks-on-chip
high-speed interconnect |
0.0 | 1 | 2011 | Hadoop acceleration through network levitated merge · SC 2011 |
Methods — techniques the papers use, named apart from their topics
mini-batch size tuning · 0.6hyperparameter tuning · 0.6data compression · 0.2compact data organization · 0.2virtual shuffling · 0.2three-level segment table · 0.2multicore awareness · 0.2balanced merging subtrees · 0.2RDMA · 0.2diffusion limit analysis · 0.2
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2017 | GaDei: On Scale-Up Training as a Service for Deep LearningabstractDeep learning (DL) training-as-a-service (TaaS) is an important emerging industrial workload. TaaS must satisfy a wide range of customers who have no experience and/or resources to tune DL hyper-parameters (e.g., mini-batch size and learning rate), and meticulous tuning for each user's dataset is prohibitively expensive. Therefore, TaaS hyper-parameters must be fixed with values that are applicable to all users. Unfortunately, few research papers have studied how to design a system for TaaS workloads. By evaluating the IBM Watson Natural Language Classfier (NLC) workloads, the most popular IBM cognitive service used by thousands of enterprise-level clients globally, we provide empirical evidence that only the conservative hyper-parameter setup (e.g., small mini-batch size) can guarantee acceptable model accuracy for a wide range of customers. Unfortunately, smaller mini-batch size requires higher communication bandwidth in a parameter-server based DL training system. In this paper, we characterize the exceedingly high communication bandwidth requirement of TaaS using representative industrial deep learning workloads. We then present GaDei, a highly optimized shared-memory based scale-up parameter server design. We evaluate GaDei using both commercial benchmarks and public benchmarks and demonstrate that GaDei significantly outperforms the state-of-the-art parameter-server based implementation while maintaining the required accuracy. GaDei achieves near-best-possible runtime performance, constrained only by the hardware limitation. Furthermore, to the best of our knowledge, GaDei is the only scale-up DL system that provides fault-tolerance. Wei Zhang 0057, Minwei Feng, Yunhui Zheng, Yufei Ren, Yandong Wang 0001, Peng Liu 0010, Bing Xiang, Li Zhang 0002, Bowen Zhou 0002, Fei Wang 0001 |
ICDM | 5 |
| 2017 | Lightweight Replication Through Remote Backup Memory Sharing for In-memory Key-Value StoresabstractMemory price will continue dropping in the next few years according to Gartner. Such trend renders it affordable for in-memory key-value stores (IMKVs) to maintain redundant memory-resident copies of each key-value pair to provision enhanced reliability and high availability services. Though contemporary IMKVs have reached unprecedented performance, delivering single-digit microsecond-scale latency with up to tens of millions queries per second throughput, existing replication protocols are unable to keep pace with such an advancement of IMKVs, either incurring unbearable latency overhead or demanding intensive resource usage. Consequently, the adoption of those replication techniques always results in substantial performance degradation.In this paper, we propose MacR, a RDMA-based high-performance and lightweight replication protocol for IMKVs. The design of MacR centers around sharing the remote backup memory to enable RDMA-based replication protocol, and synthesizes a collection of optimizations, including memory allocator cooperative replication and adaptive bulk data synchronization to control the number of network operations and to enhance the recovery performance. Performance evaluations with a variety of YCSB workloads demonstrate that MacR can efficiently outperform alternative replication methods in terms of the throughput while preserving sufficiently low latency overhead. It can also efficiently speed up the recovery process. Yandong Wang 0001, Li Zhang 0002, Michel Hack, Yufei Ren |
MASCOTS | 1 |
| 2017 | Nexus: Bringing Efficient and Scalable Training to Deep Learning FrameworksabstractDemand is mounting in the industry for scalable GPU-based deep learning systems. Unfortunately, existing training applications built atop popular deep learning frameworks, including Caffe, Theano, and Torch, etc, are incapable of conducting distributed GPU training over large-scale clusters. To remedy such a situation, this paper presents Nexus, a platform that allows existing deep learning frameworks to easily scale out to multiple machines without sacrificing model accuracy. Nexus leverages recently proposed distributed parameter management architecture to orchestrate distributed training by a large number of learners spread across the cluster. Through characterizing the run-time behavior of existing single-node based applications, Nexus is equipped with a suite of optimization schemes, including hierarchical and hybrid parameter aggregation, enhanced network and computation layer, and quality-guided communication adjustment, etc, to strengthen the communication channels and resource utilization. Empirical evaluations with a diverse set of deep learning applications demonstrate that Nexus is easy to integrate and can deliver efficient distributed training services to major deep learning frameworks. In addition, Nexus's optimization schemes are highly effective to shorten the training time with targeted accuracy bounds. Yandong Wang 0001, Li Zhang 0002, Yufei Ren, Wei Zhang 0057 |
MASCOTS | 1 |
| 2016 | zExpander: a key-value cache with both high performance and fewer missesabstractWhile key-value (KV) cache, such as memcached, dedicates a large volume of expensive memory to holding performance-critical data, it is important to improve memory efficiency, or to reduce cache miss ratio without adding more memory. As we find that optimizing replacement algorithms is of limited effect for this purpose, a promising approach is to use a compact data organization and data compression to increase effective cache size. However, this approach has the risk of degrading the cache's performance due to additional computation cost. A common perception is that a high-performance KV cache is not compatible with use of data compacting techniques. Xingbo Wu, Li Zhang 0002, Yandong Wang 0001, Yufei Ren, Michel Hack, Song Jiang 0001 |
EuroSys | 3 |
| 2016 | MEMTUNE: Dynamic Memory Management for In-Memory Data Analytic PlatformsabstractMemory is a crucial resource for big data processing frameworks such as Spark and M3R, where the memory is used both for computation and for caching intermediate storage data. Consequently, optimizing memory is the key to extracting high performance. The extant approach is to statically split the memory for computation and caching based on workload profiling. This approach is unable to capture the varying workload characteristics and dynamic memory demands. Another factor that affects caching efficiency is the choice of data placement and eviction policy. The extant LRU policy is oblivious of task scheduling information from the analytic frameworks, and thus can lead to lost optimization opportunities. In this paper, we address the above issues by designing MEMTUNE, a dynamic memory manager for in-memory data analytics. MEMTUNE dynamically tunes computation/caching memory partitions at runtime based on workload memory demand and in-memory data cache needs. Moreover, if needed, the scheduling information from the analytic framework is leveraged to evict data that will not be needed in the near future. Finally, MEMTUNE also supports task-level data prefetching with a configurable window size to more effectively overlap computation with I/O. Our experiments show that MEMTUNE improves memory utilization, yields an overall performance gain of up to 46%, and achieves cache hit ratio of up to 41% compared to standard Spark. Luna Xu, Li Zhang 0002, Ali Raza Butt, Yandong Wang 0001, Zane Zhenhua Hu |
IPDPS | 5 |
| 2015 | Cracking Down MapReduce Failure Amplification through Analytics Logging and MigrationabstractMapReduce is popular for big data analytics because it offers easy-to-use map and reduce user interfaces while hiding the complexity of system scalability and fault resiliency issues. While a large body of literature has focused on improving the performance and scalability of MapReduce, the issue of fault resiliency has thus far received little attention. In this paper, we take on an effort to investigate the fault resiliency of MapReduce using YARN (the next-generation Hadoop) as a case study. We reveal that the failures of a MapTask, a ReduceTask or a compute node can cause distinctly different impact to MapReduce programs. Particularly, YARN MapReduce is not able to gracefully handle failures that involve ReduceTasks, causing prolonged task execution, delayed job completion, and, more severely, failure amplifications due to the cascading effects to other tasks. These problems together cause the performance collapse of MapReduce jobs. In this paper, we introduce a new fault-tolerant framework that can crack down failure amplification and gracefully handle failure scenarios. It is designed with two key fault handling techniques: analytics logging and speculative fast migration. Analytics logging is a light-weight mechanism that logs the key progress information of MapReduce tasks, speculative fast migration handles node failures by proactively re-executing MapTasks, migrating ReduceTasks, and collective merging with a pipeline of shuffle/merge and reduce stages. Our performance evaluation demonstrates that these techniques can eliminate failure amplification and deliver fast job execution compared to the existing task re-execution mechanism in MapReduce. Yandong Wang 0001, Huansong Fu, Weikuan Yu |
IPDPS | 1 |
| 2015 | HydraDB: a resilient RDMA-driven key-value middleware for in-memory cluster computingabstractIn this paper, we describe our experiences and lessons learned from building a general-purpose in-memory key-value middleware, called HydraDB. HydraDB synthesizes a collection of state-of-the-art techniques, including continuous fault-tolerance, Remote Direct Memory Access (RDMA), as well as awareness for multicore systems, etc, to deliver a high-throughput, low-latency access service in a reliable manner for cluster computing applications. Yandong Wang 0001, Li Zhang 0002, Jian Tan 0001, Xavier Guerin, Xiaoqiao Meng, Shicong Meng |
SC | 1 |
| 2015 | Virtual Shuffling for Efficient Data Movement in MapReduceabstractMapReduce is a popular parallel processing framework for large-scale data analytics. To keep up with the increasing volume of datasets, it requires efficient I/O capability from the underlying computer systems to process and analyze data in two phases (mapping and reducing). Between these phases, MapReduce requires a shuffling phase to globally exchange the intermediate data generated by the mapping phase. We reveal that data shuffling, by physically moving segments of intermediate data across disks, causes significant I/O contention and compounds the I/O problem. In this paper, we propose a novel virtual shuffling strategy to enable efficient data movement and reduce I/O for MapReduce shuffling, thereby reducing power consumption and conserving energy. Virtual shuffling is realized through a combination of three techniques including a three-level segment table, near-demand merging, and dynamic and balanced merging subtrees. Our experimental results show that virtual shuffling significantly speeds up data movement in MapReduce and achieves faster job execution. Particularly, its reduction in disk I/O accesses results in as much as 12% savings in power consumption for MapReduce programs. Weikuan Yu, Yandong Wang 0001, Xinyu Que, Cong Xu 0008 |
IEEE Trans. Computers | 2 |
| 2014 | BurstMem: A high-performance burst buffer system for scientific applicationsabstractThe growth of computing power on large-scale systems requires commensurate high-bandwidth I/O systems. Many parallel file systems are designed to provide fast sustainable I/O in response to applications' soaring requirements. To meet this need, a novel system is imperative to temporarily buffer the bursty I/O and gradually flush datasets to long-term parallel file systems. In this paper, we introduce the design of BurstMem, a high-performance burst buffer system. BurstMem provides a storage framework with efficient storage and communication management strategies. Our experiments demonstrate that BurstMem is able to speed up the I/O performance of scientific applications by up to 8.5× on leadership computer systems. Teng Wang 0001, Sarp Oral, Yandong Wang 0001, Bradley W. Settlemyer, Scott Atchley, Weikuan Yu |
IEEE BigData | 3 |
| 2014 | C-Hint: An Effective and Reliable Cache Management for RDMA-Accelerated Key-Value StoresabstractRecently, many in-memory key-value stores have started using a High-Performance network protocol, Remote Direct Memory Access (RDMA), to provision ultra-low latency access services. Among various solutions, previous studies have recognized that leveraging RDMA Read to optimize GET operations and continuing using message passing for other requests can offer tremendous performance improvement while avoiding read-write races. However, although such a design can utilize the power of RDMA when there is sufficient memory space, it has also raised new challenges on the cache management that do not exist in traditional key-value stores. First, RDMA Read deprives servers of the awareness of the read operations. Therefore, how to track popular items and make replacement decisions at the server side becomes a critical issue. Second, without the access knowledge from the clients, new approaches are needed for servers to efficiently and reliably reclaim the resources. Lastly, the remote pointers hold by clients to conduct RDMA are highly susceptible to the evictions made by remote servers. Thus, any replacement algorithm that solely considers the server-side hit ratio is insufficient and can cause severe underutilization of RDMA. Yandong Wang 0001, Xiaoqiao Meng, Li Zhang 0002, Jian Tan 0001 |
SoCC | 1 |
| 2014 | Characterization and Optimization of Memory-Resident MapReduce on HPC SystemsabstractMapReduce is a widely accepted framework for addressing big data challenges. Recently, it has also gained broad attention from scientists at the U.S. leadership computing facilities as a promising solution to process gigantic simulation results. However, conventional high-end computing systems are constructed based on the compute-centric paradigm while big data analytics applications prefer a data-centric paradigm such as MapReduce. This work characterizes the performance impact of key differences between compute- and data-centric paradigms and then provides optimizations to enable a dual-purpose HPC system that can efficiently support conventional HPC applications and new data analytics applications. Using a state-of-the-art MapReduce implementation Spark and the Hyperion system at Lawrence Livermore National Laboratory, we have examined the impact of storage architectures, data locality and task scheduling to the memory-resident MapReduce jobs. Based on our characterization and findings of the performance behaviors, we have introduced two optimization techniques, namely Enhanced Load Balancer and Congestion-Aware Task Dispatching, to improve the performance of Spark applications. Yandong Wang 0001, Robin Goldstone, Weikuan Yu, Teng Wang 0001 |
IPDPS | 1 |
| 2014 | Non-work-conserving effects in MapReduce: diffusion limit and criticalityabstractSequentially arriving jobs share a MapReduce cluster, each desiring a fair allocation of computing resources to serve its associated map and reduce tasks. The model of such a system consists of a processor sharing queue for the MapTasks and a multi-server queue for the ReduceTasks. These two queues are dependent through a constraint that the input data of each ReduceTask are fetched from the intermediate data generated by the MapTasks belonging to the same job. A more generalized form of MapReduce queueing model can capture the essence of other distributed data processing systems that contain interdependent processor sharing queues and multi-server queues. Jian Tan 0001, Yandong Wang 0001, Weikuan Yu, Li Zhang 0002 |
SIGMETRICS | 2 |
| 2014 | Design and Evaluation of Network-Levitated Merge for Hadoop AccelerationabstractHadoop is a popular open source implementation of the MapReduce programming model for cloud computing. However, it faces a number of issues to achieve the best performance from the underlying systems. These include a serialization barrier that delays the reduce phase, repetitive merges, and disk accesses, and the lack of portability to different interconnects. To keep up with the increasing volume of data sets, Hadoop also requires efficient I/O capability from the underlying computer systems to process and analyze data. We describe Hadoop-A, an acceleration framework that optimizes Hadoop with plug-in components for fast data movement, overcoming the existing limitations. A novel network-levitated merge algorithm is introduced to merge data without repetition and disk access. In addition, a full pipeline is designed to overlap the shuffle, merge, and reduce phases. Our experimental results show that Hadoop-A significantly speeds up data movement in MapReduce and doubles the throughput of Hadoop. In addition, Hadoop-A significantly reduces disk accesses caused by intermediate data. Weikuan Yu, Yandong Wang 0001, Xinyu Que |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2013 | SLOAVx: Scalable LOgarithmic AlltoallV Algorithm for Hierarchical Multicore SystemsabstractScientific applications use collective communication operations in Message Passing Interface (MPI) for global synchronization and data exchanges. Alltoall and AlltoallV are two important collective operations. They are used by MPI jobs to exchange messages among all of MPI processes. AlltoallV is a generalization of Alltoall, supporting messages of varying sizes. However, the existing MPI AlltoallV implementation has linear complexity, i.e., each process has to send messages to all other processes in the job. Such linear complexity can result in sub optimal scalability of MPI applications when they are deployed on millions of cores. To address above challenge, in this paper, we introduce a new Scalable LOgarithmic AlltoallV algorithm, named SLOAV, for MPI AlltoallV collective operation. SLOAV aims to achieve global exchange of small messages of different sizes in a logarithmic number of rounds. Furthermore, given the prevalence of multicore systems with shared memory, we design a hierarchical AlltoallV algorithm based on SLOAV by leveraging the advantages of shared memory, which is referred to as SLOAVx. Compared to SLOAV, SLOAVx significantly reduces the inter-node communication, thus improving the entire system performance and mitigating the impact of message latency. We have implemented and embedded both algorithms in Open MPI. Our evaluation on large-scale computer systems shows that for the 8-byte and 1024-process MPI Alltoallv operation, the SLOAV can reduce the latency by as much as 86.4%, when compared to the state-of-the-art, and SLOAVx can further optimize the SLOAV by up to 83.1% in terms of message latency on multicore systems. In addition, experiments with NAS Parallel Benchmark (NPB) demonstrate that our algorithms are very effective for real-world applications. Cong Xu 0008, Manjunath Gorentla Venkata, Richard L. Graham, Yandong Wang 0001, Weikuan Yu |
CCGRID | 4 |
| 2013 | Profiling and Improving I/O Performance of a Large-Scale Climate Scientific ApplicationabstractExascale computing systems are soon to emerge, which will pose great challenges on the huge gap between computing and I/O performance. Many large-scale scientific applications play an important role in our daily life. The huge amounts of data generated by such applications require highly parallel and efficient I/O management policies. In this paper, we adopt a mission-critical scientific application, GEOS-5, as a case to profile and analyze the communication and I/O issues that are preventing applications from fully utilizing the underlying parallel storage systems. Through in-detail architectural and experimental characterization, we observe that current legacy I/O schemes incur significant network communication overheads and are unable to fully parallelize the data access, thus degrading applications' I/O performance and scalability. To address these inefficiencies, we redesign its I/O framework along with a set of parallel I/O techniques to achieve high scalability and performance. Evaluation results on the NASA discover cluster show that our optimization of GEOS- 5 with ADIOS has led to significant performance improvements compared to the original GEOS-5 implementation. Bin Wang 0019, Teng Wang 0001, Yuan Tian 0004, Cong Xu 0008, Yandong Wang 0001, Weikuan Yu, Carlos A. Cruz, Shujia Zhou, Thomas L. Clune, Scott Klasky |
ICCCN | 6 |
| 2013 | JVM-Bypass for Efficient Hadoop ShufflingabstractHadoop employs Java-based network transport stack on top of the Java Virtual Machine (JVM) for its data shuffling and merging purposes. Our examination reveals that JVM introduces a significant amount of overhead to data processing capability of the native interface. Furthermore, JVM constrains the use of high-performance networking mechanisms such as RDMA (Remote Direct Memory Access) which has established itself as an effective data movement technology in many networking environments because of its low-latency, high bandwidth, low CPU utilization, and energy efficiency. In this paper, we introduce a plug-in library called JVM-Bypass Shuffling (JBS) for Hadoop data shuffling. JBS helps Hadoop data shuffling by avoiding Javabased transport protocols, removing the overhead and limitations of the JVM. In addition, we design JBS as a portable library that can leverage both TCP/IP and RDMA on different network systems such as InfiniBand and 1/10 Gigabit Ethernet. We have designed and implemented JBS as part of Hadoop acceleration. It has been transferred to Mellanox as the software product UDA (Unstructured Data Accelerator) and used to enable our studies on a variety of merging algorithms. Our performance evaluation demonstrates that JBS can effectively reduce the execution time of Hadoop jobs by up to 66.3% and lower the CPU utilization by 48.1%. Yandong Wang 0001, Cong Xu 0008, Weikuan Yu |
IPDPS | 1 |
| 2013 | CooMR: cross-task coordination for efficient data management in MapReduce programsabstractHadoop is a widely adopted open source implementation of MapReduce programming model for big data processing. It represents system resources as available map and reduce slots and assigns them to various tasks. This execution model gives little regard to the need of cross-task coordination on the use of shared system resources on a compute node, which results in task interference. In addition, the existing Hadoop merge algorithm can cause excessive I/O. In this study, we undertake an effort to address both issues. Accordingly, we have designed a cross-task coordination framework called CooMR for efficient data management in MapReduce programs. CooMR consists of three component schemes including cross-task opportunistic memory sharing and log-structured I/O consolidation, which are designed to facilitate task coordination, and the key-based in-situ merge (KISM) algorithm which is designed to enable the sorting/merging of Hadoop intermediate data without actually moving the pairs. Our evaluation demonstrates that CooMR is able to increase task coordination, improve system resource utilization, and significantly speed up the execution time of MapReduce programs. Yandong Wang 0001, Yizheng Jiao, Cong Xu 0008, Weikuan Yu |
SC | 2 |
| 2011 | EDO: Improving Read Performance for Scientific Applications through Elastic Data OrganizationabstractLarge scale scientific applications are often bottlenecked due to the writing of checkpoint-restart data. Much work has been focused on improving their write performance. With the mounting needs of scientific discovery from these datasets, it is also important to provide good read performance for many common access patterns, which requires effective data organization. To address this issue, we introduce Elastic Data Organization (EDO), which can transparently enable different data organization strategies for scientific applications. Through its flexible data ordering algorithms, EDO harmonizes different access patterns with the underlying file system. Two levels of data ordering are introduced in EDO. One works at the level of data groups (a.k.a process groups). It uses Hilbert Space Filling Curves (SFC) to balance the distribution of data groups across storage targets. Another governs the ordering of data elements within a data group. It divides a data group into sub chunks and strikes a good balance between the size of sub chunks and the number of seek operations. Our experimental results demonstrate that EDO is able to achieve balanced data distribution across all dimensions and improve the read performance of multidimensional datasets in scientific applications. Yuan Tian 0004, Scott Klasky, Hasan Abbasi, Jay F. Lofstead, Ray W. Grout, Norbert Podhorszki, Qing Liu 0002, Yandong Wang 0001, Weikuan Yu |
CLUSTER | 8 |
| 2011 | Hadoop acceleration through network levitated mergeabstractHadoop is a popular open-source implementation of the MapReduce programming model for cloud computing. However, it faces a number of issues to achieve the best performance from the underlying system. These include a serialization barrier that delays the reduce phase, repetitive merges and disk access, and lack of capability to leverage latest high speed interconnects. We describe Hadoop-A, an acceleration framework that optimizes Hadoop with plugin components implemented in C++ for fast data movement, overcoming its existing limitations. A novel network-levitated merge algorithm is introduced to merge data without repetition and disk access. In addition, a full pipeline is designed to overlap the shuffle, merge and reduce phases. Our experimental results show that Hadoop-A doubles the data processing throughput of Hadoop, and reduces CPU utilization by more than 36%. Yandong Wang 0001, Xinyu Que, Weikuan Yu, Dror Goldenberg, Dhiraj Sehgal |
SC | 1 |