Ioan Raicu

dblp:94/3254 · DBLP profile ↗
← Back
54ranked-venue papers
6as first author
7since 2021 · last 2025
0000-0002-5477-439XORCID · verified

Domains — the database's venue-derived domains; a paper can count in several

Systems, architecture and hardware · 44 · 6 first-author · 6 since 2021Applied, interdisciplinary, general and emerging computing · 7 · 1 since 2021Artificial intelligence and machine learning · 5Databases, data management, data science and information retrieval · 5Software engineering, systems software and programming languages · 4 · 1 since 2021
YearPublicationVenuePosition
2025 CryptexLLM: How LLM Generalizability Forecasts High Volatility
abstract
We present CryptexLLM, an approach to time series forecasting that adapts large language models (LLMs) for predicting high volatility data. We extend the TimeLLM framework with an adaptive weighted loss function, feature engineering, and sentiment analysis. Our experiments show that the approach we took outperforms traditional LSTM models and statistical methods, with our best performing model being Llama 3.1. The adaptations we made improve directional accuracy, which is particularly useful for financial applications, while still maintaining computational efficiency. Our results suggest that LLMs have the potential to effectively generalize to volatile time series domains.
Logan Rao, Mary Tsaryk, Anwar Benhnini, Ioan Raicu
eScience4
2025 Optimizing Fine-Grained Parallelism Through Dynamic Load Balancing on Multi-Socket Many-Core Systems
abstract
Achieving efficient task parallelism on many-core architectures is an important challenge. The widely used GNU OpenMP implementation of the popular OpenMP parallel programming model incurs high overhead for fine-grained, shortrunning tasks due to time spent on runtime synchronization. In this work, we introduce and analyze three key advances that collectively achieve significant performance gains. First, we introduce XQueue, a lock-less concurrent queue implementation to replace GNU's priority task queue and remove the global task lock. Second, we develop a scalable, efficient, and hybrid lock-free/lock-less distributed tree barrier to address the high hardware synchronization overhead from GNU's centralized barrier. Third, we develop two lock-less and NUMA-aware load balancing strategies. We evaluate our implementation using Barcelona OpenMP Task Suite (BOTS) benchmarks. We show that the use of XQueue and the distributed tree barrier can improve performance by up to$1522.8 \times$compared to the original GNU OpenMP. We further show that lock-less load balancing can improve performance by up to$4 \times$compared to GNU OpenMP using XQueue.
Maxime Gonthier, Poornima Nookala, Haochen Pan, Ian T. Foster, Ioan Raicu, Kyle Chard
IPDPS6
2024 X-OpenMP - eXtreme fine-grained tasking using lock-less work stealing
Poornima Nookala, Kyle Chard, Ioan Raicu
Future Gener. Comput. Syst.3
2024 SCIPIS: Scalable and concurrent persistent indexing and search in high-end computing systems
Alexandru Iulian Orhean, Anna Giannakou, Lavanya Ramakrishnan, Kyle Chard, Boris Glavic, Ioan Raicu
J. Parallel Distributed Comput.6
2022 SCANNS: Towards Scalable and Concurrent Data Indexing and Searching in High-End Computing System
Alexandru Iulian Orhean, Anna Giannakou, Lavanya Ramakrishnan, Kyle Chard, Ioan Raicu
CCGRID5
2022 Evaluation of a scientific data search infrastructure
abstract
Summary The ability to search over large scientific datasets has become crucial to next‐generation scientific discoveries as data generated from scientific facilities grow dramatically. In previous work, we developed and deployed ScienceSearch, a search infrastructure for scientific data which uses machine learning to automate metadata creation. Our current deployment is deployed atop a container based platform at a HPC center. In this article, we present an evaluation and discuss our experiences with the ScienceSearch infrastructure. Specifically, we present a performance evaluation of ScienceSearch's infrastructure focusing on scalability trends. The obtained results show that ScienceSearch is able to serve up to 130 queries/min with latency under 3 s. We discuss our infrastructure setup and evaluation results to provide our experiences and a perspective on opportunities and challenges of our search infrastructure.
Alexandru Iulian Orhean, Anna Giannakou, Katie Antypas, Ioan Raicu, Lavanya Ramakrishnan
Concurr. Comput. Pract. Exp.4
2021 Enabling Extremely Fine-grained Parallelism via Scalable Concurrent Queues on Modern Many-core Architectures
abstract
Enabling efficient fine-grained task parallelism is a significant challenge for hardware platforms with increasingly many cores. Existing techniques do not scale to hundreds of threads due to the high cost of synchronization in concurrent data structures. To overcome these limitations we present XQueue, a novel lock-less concurrent queuing system with relaxed ordering semantics that is geared towards realizing scalability up to hundreds of concurrent threads. We demonstrate the scalability of XQueue using microbenchmarks and show that XQueue can deliver concurrent operations with latencies as low as 110 cycles at scales of up to 192 cores (up to 6900× improvement compared to traditional synchronization mechanisms) across our diverse hardware, including x86, ARM, and Power9. The reduced latency allows XQueue to provide orders of magnitude (3300×) better throughput that existing techniques. To evaluate the real-world benefits of XQueue, we integrated XQueue with LLVM OpenMP and evaluated five unmodified benchmarks from the Barcelona OpenMP Task Suite (BOTS) as well as a graph traversal benchmark from the GAP benchmark suite. We compared the XQueue-enabled LLVM OpenMP implementation with the native LLVM and GNU OpenMP versions. Using fine-grained task workloads, XQueue can deliver 4× to 6× speedup compared to native GNU OpenMP and LLVM OpenMP in many cases, with speedups as high as 116× in some cases.
Poornima Nookala, Peter A. Dinda, Kyle C. Hale, Kyle Chard, Ioan Raicu
MASCOTS5
2019 A High-Performance Distributed Relational Database System for Scalable OLAP Processing
abstract
The scalability of systems such as Hive and Spark SQL that are built on top of big data platforms have enabled query processing over very large data sets. However, the per-node performance of these systems is typically low compared to traditional relational databases. Conversely, Massively Parallel Processing (MPP) databases do not scale as well as these systems. We present HRDBMS, a fully implemented distributed shared-nothing relational database developed with the goal of improving the scalability of OLAP queries. HRDBMS achieves high scalability through a principled combination of techniques from relational and big data systems with novel communication and work-distribution techniques. While we also support serializable transactions, the system has not been optimized for this use case. HRDBMS runs on a custom distributed and asynchronous execution engine that was built from the ground up to support highly parallelized operator implementations. Our experimental comparison with Hive, Spark SQL, and Greenplum confirms that HRDBMS's scalability is on par with Hive and Spark SQL (up to 96 nodes) while its per-node performance can compete with MPP databases like Greenplum.
Jason Arnold, Boris Glavic, Ioan Raicu
IPDPS3
2018 New scheduling approach using reinforcement learning for heterogeneous distributed systems
Alexandru Iulian Orhean, Florin Pop, Ioan Raicu
J. Parallel Distributed Comput.3
2017 Understanding the Performance and Potential of Cloud Computing for Scientific Applications
abstract
Commercial clouds bring a great opportunity to the scientific computing area. Scientific applications usually require significant resources, however not all scientists have access to sufficient high-end computing systems. Cloud computing has gained the attention of scientists as a competitive resource to run HPC applications at a potentially lower cost. But as a different infrastructure, it is unclear whether clouds are capable of running scientific applications with a reasonable performance per money spent. This work provides a comprehensive evaluation of EC2 cloud in different aspects. We first analyze the potentials of the cloud by evaluating the raw performance of different services of AWS such as compute, memory, network and I/O. Based on the findings on the raw performance, we then evaluate the performance of the scientific applications running in the cloud. Finally, we compare the performance of AWS with a private cloud, in order to find the root cause of its limitations while running scientific applications. This paper aims to assess the ability of the cloud to perform well, as well as to evaluate the cost of the cloud in terms of both raw performance and scientific applications performance. Furthermore, we evaluate other services including S3, EBS and DynamoDB among many AWS services in order to assess the abilities of those to be used by scientific applications and frameworks. We also evaluate a real scientific computing application through the Swift parallel scripting system at scale. Armed with both detailed benchmarks to gauge expected performance and a detailed monetary cost analysis, we expect this paper will be a recipe cookbook for scientists to help them decide where to deploy and run their scientific applications between public clouds, private clouds, or hybrid clouds.
Iman Sadooghi, Jesus Hernandez Martin, Tonglin Li, Kevin Brandstatter, Ketan Maheshwari, Tiago Pais Pitta De Lacerda Ruivo, Gabriele Garzoglio, Steven Timm, Yong Zhao 0009, Ioan Raicu
IEEE Trans. Cloud Comput.10
2016 Albatross: An efficient cloud-enabled task scheduling and execution framework using distributed message queues
abstract
Data Analytics has become very popular on large datasets in different organizations. It is inevitable to use distributed resources such as Clouds for Data Analytics and other types of data processing at larger scales. To effectively utilize all system resources, an efficient scheduler is needed, but the traditional resource managers and job schedulers are centralized and designed for larger batch jobs which are fewer in number. Frameworks such as Hadoop and Spark, which are mainly designed for Big Data analytics, have been able to allow for more diversity in job types to some extent. However, even these systems have centralized architectures and will not be able to perform well on large scales and under heavy task loads. Modern applications generate tasks at very high rates that can cause significant slowdowns on these frameworks. Additionally, over-decomposition has shown to be very useful in increasing the system utilization. In order to achieve high efficiency, scalability, and better system utilization, it is critical for a modern scheduler to be able to handle over-decomposition and run highly granular tasks. Further, to achieve high performance, Albatross is written in C/C++, which imposes a minimal overhead to the workload process as compared to languages like Java or Python. We propose Albatross, a task level scheduling and execution framework that uses a Distributed Message Queue (DMQ) for task distribution among its workers. Unlike most scheduling systems, Albatross uses a pulling approach as opposed to the common push approach. The former would let Albatross achieve a good load balancing and scalability. Furthermore, the framework has built in support for task execution dependency on workflows. Therefore, Albatross is able to run various types of workloads, including Data Analytics and HPC applications. Finally, Albatross provides data locality support. This allows the framework to achieve higher performance through minimizing the amount of unnecessary data movement on the network. Our evaluations show that Albatross outperforms Spark and Hadoop at larger scales and in the case of running higher granularity workloads.
Iman Sadooghi, Geet Kumar, Ke Wang 0012, Dongfang Zhao 0001, Tonglin Li, Ioan Raicu
eScience6
2016 A convergence of key-value storage systems from clouds to supercomputers
abstract
Summary This paper presents a convergence of distributed key‐value storage systems in clouds and supercomputers. It specifically presents ZHT, a zero‐hop distributed key‐value store system, which has been tuned for the requirements of high‐end computing systems. ZHT aims to be a building block for future distributed systems, such as parallel and distributed file systems, distributed job management systems, and parallel programming systems. ZHT has some important properties, such as being lightweight, dynamically allowing nodes join and leave, fault tolerant through replication, persistent, scalable, and supporting unconventional operations such as append, compare and swap, callback in addition to the traditional insert/lookup/remove. We have evaluated ZHT's performance under a variety of systems, ranging from a Linux cluster with 64 nodes, an Amazon EC2 virtual cluster up to 96 nodes, to an IBM Blue Gene/P supercomputer with 8K nodes. We compared ZHT against other key‐value stores and found it offers superior performance for the features and portability it supports. This paper also presents several real systems that have adopted ZHT, namely, FusionFS (a distributed file system), IStore (a storage system with erasure coding), MATRIX (distributed scheduling), Slurm++ (distributed HPC job launch), Fabriq (distributed message queue management); all of these real systems have been simplified because of key‐value storage systems and have been shown to outperform other leading systems by orders of magnitude in some cases. It is important to highlight that some of these systems are rooted in HPC systems from supercomputers, while others are rooted in clouds and ad hoc distributed systems; through our work, we have shown how versatile key‐value storage systems can be in such a variety of environments. Copyright © 2015 John Wiley & Sons, Ltd.
Tonglin Li, Xiaobing Zhou, Ke Wang 0012, Dongfang Zhao 0001, Iman Sadooghi, Zhao Zhang 0007, Ioan Raicu
Concurr. Comput. Pract. Exp.7
2016 Load-balanced and locality-aware scheduling for data-intensive workloads at extreme scales
abstract
Summary Data‐driven programming models such as many‐task computing (MTC) have been prevalent for running data‐intensive scientific applications. MTC applies over‐decomposition to enable distributed scheduling. To achieve extreme scalability, MTC proposes a fully distributed task scheduling architecture that employs as many schedulers as the compute nodes to make scheduling decisions. Achieving distributed load balancing and best exploiting data locality are two important goals for the best performance of distributed scheduling of data‐intensive applications. Our previous research proposed a data‐aware work‐stealing technique to optimize both load balancing and data locality by using both dedicated and shared task ready queues in each scheduler. Tasks were organized in queues based on the input data size and location. Distributed key‐value store was applied to manage task metadata. We implemented the technique in MATRIX, a distributed MTC task execution framework. In this work, we devise an analytical suboptimal upper bound of the proposed technique, compare MATRIX with other scheduling systems, and explore the scalability of the technique at extreme scales. Results show that the technique is not only scalable but can achieve performance within 15% of the suboptimal solution. Copyright © 2015 John Wiley & Sons, Ltd.
Ke Wang 0012, Kan Qiao, Iman Sadooghi, Xiaobing Zhou, Tonglin Li, Michael Lang 0003, Ioan Raicu
Concurr. Comput. Pract. Exp.7
2016 Exploiting multi-cores for efficient interchange of large messages in distributed systems
abstract
Summary Conventional data serialization tools assume that objects to be coded are usually small in size so a single CPU core can encode it in a timely manner. In the era of Big Data, however, object gets increasingly complex and larger, which makes data serialization become a new performance bottleneck. This paper describes an approach to parallelize data serialization by leveraging multiple cores. Parallelizing data serialization introduces new questions such as how to split the (sub)objects, how to allocate the available cores, and how to minimize its overhead in practice. In this paper we design a framework for parallelly serializing large objects and analyze the design tradeoffs under different scenarios. To validate the proposed approach, we implemented parallel protocol buffers—the parallel version of Google's Protocol Buffers, a widely‐used data serialization utility. Experimental results confirm the effectiveness of Parallel Protocol Buffers: multiple cores employed in data serialization achieve highly scalable performance and incur negligible overhead. Copyright © 2015 John Wiley & Sons, Ltd.
Dongfang Zhao 0001, Kan Qiao, Zhou Zhou 0006, Tonglin Li, Xiaobing Zhou, Ioan Raicu
Concurr. Comput. Pract. Exp.6
2016 Toward high-performance key-value stores through GPU encoding and locality-aware encoding
Dongfang Zhao 0001, Ke Wang 0012, Kan Qiao, Tonglin Li, Iman Sadooghi, Ioan Raicu
J. Parallel Distributed Comput.6
2016 Guest Editors Introduction: Special Issue on Scientific Cloud Computing
abstract
The papers in this special section contribute important advances towards leveraging clouds for scientific applications. The contributions focus on a broad range of topics, including: performance modeling and optimization, data management, resource allocation and scheduling, elasticity, reconfiguration, cost prediction and optimization. Most papers revolve around general techniques and approaches that are agnostic of the applications, while two contributions demonstrate how domain specific scientific applications can be migrated to the cloud.
Kate Keahey, Ioan Raicu, Kyle Chard, Bogdan Nicolae
IEEE Trans. Cloud Comput.2
2016 Exploring the Design Tradeoffs for Extreme-Scale High-Performance Computing System Software
abstract
Owing to the extreme parallelism and the high component failure rates of tomorrow's exascale, high-performance computing (HPC) system software will need to be scalable, failure-resistant, and adaptive for sustained system operation and full system utilizations. Many of the existing HPC system software are still designed around a centralized server paradigm and hence are susceptible to scaling issues and single points of failure. In this article, we explore the design tradeoffs for scalable system software at extreme scales. We propose a general system software taxonomy by deconstructing common HPC system software into their basic components. The taxonomy helps us reason about system software as follows: (1) it gives us a systematic way to architect scalable system software by decomposing them into their basic components; (2) it allows us to categorize system software based on the features of these components, and finally (3) it suggests the configuration space to consider for design evaluation via simulations or real implementations. Further, we evaluate different design choices of a representative system software, i.e. key-value store, through simulations up to millions of nodes. Finally, we show evaluation results of two distributed system software, Slurm++ (a distributed HPC resource manager) and MATRIX (a distributed task execution framework), both developed based on insights from this work. We envision that the results in this article help to lay the foundations of developing next-generation HPC system software for extreme scales.
Ke Wang 0012, Abhishek Kulkarni, Michael Lang 0003, Dorian C. Arnold, Ioan Raicu
IEEE Trans. Parallel Distributed Syst.5
2016 Towards Exploring Data-Intensive Scientific Applications at Extreme Scales through Systems and Simulations
abstract
The state-of-the-art storage architecture of high-performance computing systems was designed decades ago, and with today's scale and level of concurrency, it is showing significant limitations. Our recent work proposed a new architecture to address the I/O bottleneck of the conventional wisdom, and the system prototype (FusionFS) demonstrated its effectiveness on up to 16 K nodes-the scale on par with today's largest supercomputers. The main objective of this paper is to investigate FusionFS's scalability towards exascale. Exascale computers are predicted to emerge by 2018, comprising millions of cores and billions of threads. We built an event-driven simulator (FusionSim) according to the FusionFS architecture, and validated it with FusionFS's traces. FusionSim introduced less than 4 percent error between its simulation results and FusionFS traces. With FusionSim we simulated workloads on up to two million nodes and find out almost linear scalability of I/O performance; results justified FusionFS's viability for exascale systems. In addition to the simulation work, this paper extends the FusionFS system prototype in the following perspectives: (1) the fault tolerance of file metadata is supported, (2) the limitations of the current system design is discussed, and (3) a more thorough performance evaluation is conducted, such as N-to-1 metadata write, system efficiency, and more platforms such as Amazon Cloud.
Dongfang Zhao 0001, Ning Liu 0008, Dries Kimpe, Robert B. Ross, Xian-He Sun, Ioan Raicu
IEEE Trans. Parallel Distributed Syst.6
2016 Dynamic Virtual Chunks: On Supporting Efficient Accesses to Compressed Scientific Data
abstract
Data compression could ameliorate the I/O pressure of data-intensive scientific applications. Unfortunately, the conventional wisdom of naively applying data compression to the file or block brings the dilemma between efficient random accesses and high compression ratios. File-level compression barely supports efficient random accesses to the compressed data: any retrieval request need trigger the decompression from the beginning of the compressed file. Block-level compression provides flexible random accesses to the compressed blocks, but introduces extra overhead when applying the compressor to each and every block that results in a degraded overall compression ratio. This paper extends our prior work that introduces virtual chunks offering efficient random accesses to the compressed scientific data without sacrificing the compression ratio. Virtual chunks are logical blocks pointed at by appended references without breaking the physical continuity of the file content. These references allow the decompression to start from an arbitrary position (efficient random accesses), while no per-block overhead is introduced because the file's physical entirety is retained (high compression ratio). One limitation of virtual chunk is it only supports static references. This paper presents the algorithms, analysis, and evaluations of dynamic virtual chunks to deal with the cases where the references are updated dynamically.
Dongfang Zhao 0001, Kan Qiao, Jian Yin 0002, Ioan Raicu
IEEE Trans. Serv. Comput.4
2015 A flexible QoS fortified distributed key-value storage system for the cloud
abstract
In the era of big data and cloud, distributed key-value stores are increasingly used as building blocks of large-scale applications. Comparing to traditional relational databases, key-value stores are particularly compelling due to their low latency and excellent scalability. Many big companies, such as Facebook and Amazon, run multiple different applications and services on top of a single key-value store deployment to reduce the deployment and maintenance complexity as well as economic cost. However, every application has its performance requirement but most current key-value store systems are designed to serve every application request equally. This design works well when a single application accesses the key-value store, but it is not as good for the emerging concurrent multi-application scenario. In this paper, we present ZHT/Q, a flexible QoS (Quality of Service) fortified distributed key-value storage system for clouds and data centers. It improves the overall throughput by an order of magnitude and still satisfies different applications' latency requirements with QoS using dynamic and adaptive request batching mechanisms. The experiment results show that our new system delivers up to 28 times higher throughput than the base solution while more than 99% of requests' latency requirements are satisfied.
Tonglin Li, Ke Wang 0012, Dongfang Zhao 0001, Kan Qiao, Iman Sadooghi, Xiaobing Zhou, Ioan Raicu
IEEE BigData7
2015 MHT: A light-weight scalable zero-hop MPI enabled distributed key-value store
abstract
In this paper, we propose and implement a key-value store that supports MPI while allowing application access at any time without having to declaring in the same MPI communication world. This feature may significantly simplify the application design and allow programmers leverage the power of key-value store in an intuitive way. In our preliminary experiment results captured from a supercomputer at Los Alamos National Laboratory, our prototype shows linear scalability at up to 256 nodes.
Xiaobing Zhou, Tonglin Li, Ke Wang 0012, Dongfang Zhao 0001, Iman Sadooghi, Ioan Raicu
IEEE BigData6
2015 HRDBMS: A NewSQL Database for Analytics
abstract
HRDBMS is a novel distributed relational database that uses a hybrid model combining the best of traditional distributed relational databases and Big Data analytics platforms such as Hive. This allows HRDBMS to leverage years worth of research regarding query optimization, while also taking advantage of the scalability of Big Data platforms. The system uses an execution framework that is tailored for relational processing, thus addressing some of the performance challenges of running SQL on top of platforms such as MapReduce and Spark. These include excessive materialization of intermediate results, lack of a global cost-based optimization, unnecessary sorting, lack of index support, no statistics, no support for DML and ACID, and excessive communication caused by the rigid communication patterns enforced by these platforms.
Jason Arnold, Boris Glavic, Ioan Raicu
CLUSTER3
2015 GRAPH/Z: A Key-Value Store Based Scalable Graph Processing System
abstract
The emerging applications in big data and social networks issue rapidly increasing demands on graph processing. Graph query operations that involve a large number of vertices and edges can be tremendously slow on traditional databases. The state-of-the-art graph processing systems and databases usually adopt master/slave architecture that potentially impairs their The contributions of this paper are as follows: scalability. This work describes the design and implementation of a new graph processing system based on Bulk Synchronous Parallel model. Our system is built on top of ZHT, a scalable distributed key-value store, which benefits the graph processing in terms of scalability, performance and persistency. The experiment results imply excellent scalability.
Tonglin Li, Chaoqi Ma, Xiaobing Zhou, Ke Wang 0012, Dongfang Zhao 0001, Iman Sadooghi, Ioan Raicu
CLUSTER8
2015 Evaluating the Support of MTC Applications on Intel Xeon Phi Many-Core Accelerators
abstract
As Many-Task Computing (MTC) is becoming common-place on clusters, grids, and supercomputers, research that aims to take advantage of the new advances in hardware for MTC workloads is becoming more relevant. A good example is the design of frameworks like GeMTC that incorporate general purpose GPU hardware to improve the concurrency of executing tasks. This work attempts to support MTC workloads on the Intel Xeon Phi accelerators. Our plan is to develop two frameworks that will achieve that goal. One based on OpenMP and the other one based on Intel's Symmetric Communication Interface (SCIF) provided for Many-Integrated Core (MIC) accelerators like the Xeon Phi. Both frameworks aim to provide the same interface as GeMTC, leveraging the integration efforts with the Swift parallel programming system. Our end-goal is to present how programming many-core computing processors can be made easier and more productive using OpenMP or SCIF, and enable the execution of MTC workloads hybrid accelerator-based systems.
Poornima Nookala, Serapheim Dimitropoulos, Karl Stough, Ioan Raicu
CLUSTER4
2015 Overcoming Hadoop Scaling Limitations through Distributed Task Execution
abstract
Data driven programming models like MapReduce have gained the popularity in large-scale data processing. Although great efforts through the Hadoop implementation and framework decoupling (e.g. YARN, Mesos) have allowed Hadoop to scale to tens of thousands of commodity cluster processors, the centralized designs of the resource manager, task scheduler and metadata management of HDFS file system adversely affect Hadoop's scalability to tomorrow's extreme-scale data centers. This paper aims to address the YARN scaling issues through a distributed task execution framework, MATRIX, which was originally designed to schedule the executions of data-intensive scientific applications of many-task computing on supercomputers. We propose to leverage the distributed design wisdoms of MATRIX to schedule arbitrary data processing applications in cloud. We compare MATRIX with YARN in processing typical Hadoop workloads, such as WordCount, TeraSort, Grep and RandomWriter, and the Ligand application in Bioinformatics on the Amazon Cloud. Experimental results show that MATRIX outperforms YARN by 1.27X for the typical workloads, and by 2.04X for the real application. We also run and simulate MATRIX with fine-grained sub-second workloads. With the simulation results giving the efficiency of 86.8% at 64K cores for the 150ms workload, we show that MATRIX has the potential to enable Hadoop to scale to extreme-scale data centers for fine-grained workloads.
Ke Wang 0012, Ning Liu 0008, Iman Sadooghi, Xi Yang 0002, Xiaobing Zhou, Tonglin Li, Michael Lang 0003, Xian-He Sun, Ioan Raicu
CLUSTER9
2015 Towards Scalable Distributed Workload Manager with Monitoring-Based Weakly Consistent Resource Stealing
abstract
One way to efficiently utilize the coming exascale machines is to support a mixture of applications in various domains, such as traditional large-scale HPC, the ensemble runs, and the fine-grained many-task computing (MTC). Delivering high performance in resource allocation, scheduling and launching for all types of jobs has driven us to develop Slurm++, a distributed workload manager directly extended from the Slurm centralized production system. Slurm++ employs multiple controllers with each one managing a partition of compute nodes and participating in resource allocation through resource balancing techniques. In this paper, we propose a monitoring-based weakly consistent resource stealing technique to achieve resource balancing in distributed HPC job launch, and implement the technique in Slurm++. We compare Slurm++ with Slurm using micro-benchmark workloads with different job sizes. Slurm++ showed 10X faster than Slurm in allocating resources and launching jobs -- we expect the performance gap to grow as the job sizes and system scales increase in future high-end computing systems.
Ke Wang 0012, Xiaobing Zhou, Kan Qiao, Michael Lang 0003, Benjamin McClelland, Ioan Raicu
HPDC6
2015 Enabling scalable scientific workflow management in the Cloud
Yong Zhao 0009, Youfu Li 0002, Ioan Raicu, Shiyong Lu, Wenhong Tian, Heng Liu 0004
Future Gener. Comput. Syst.3
2015 A Service Framework for Scientific Workflow Management in the Cloud
abstract
Cloud computing is an emerging computing paradigm that can offer unprecedented scalability and resources on demand, and is getting more and more adoption in the science community, while scientific workflow management systems provide essential support such as management of data and task dependencies, job scheduling and execution, provenance tracking, etc., to scientific computing. As we are entering into a “big data” era, it is imperative to migrate scientific workflow management systems into the cloud to manage the ever increasing data scale and analysis complexity. We propose a reference service framework for integrating scientific workflow management systems into various cloud platforms, which consists of eight major components, including Cloud Workflow Management Service, Cloud Resource Manager, etc., and six interfaces between them. We also present a reference framework for the implementation of Cloud Resource Manager, which is responsible for the provisioning and management of virtual resources in the cloud. We discuss our implementation of the framework by integrating the Swift scientific workflow management system with the OpenNebula and Eucalyptus cloud platforms, and demonstrate the capability of the solution using a NASA MODIS image processing workflow and a production deployment on the Science@Guoshi network with support for the Montage image mosaic workflow.
Yong Zhao 0009, Youfu Li 0002, Ioan Raicu, Shiyong Lu, Cui Lin, Wenhong Tian, Ruini Xue
IEEE Trans. Serv. Comput.3
2014 Optimizing load balancing and data-locality with data-aware scheduling
abstract
Load balancing techniques (e.g. work stealing) are important to obtain the best performance for distributed task scheduling systems that have multiple schedulers making scheduling decisions. In work stealing, tasks are randomly migrated from heavy-loaded schedulers to idle ones. However, for data-intensive applications where tasks are dependent and task execution involves processing a large amount of data, migrating tasks blindly yields poor data-locality and incurs significant data-transferring overhead. This work improves work stealing by using both dedicated and shared queues. Tasks are organized in queues based on task data size and location. We implement our technique in MATRIX, a distributed task scheduler for many-task computing. We leverage distributed key-value store to organize and scale the task metadata, task dependency, and data-locality. We evaluate the improved work stealing technique with both applications and micro-benchmarks structured as direct acyclic graphs. Results show that the proposed data-aware work stealing technique performs well.
Ke Wang 0012, Xiaobing Zhou, Tonglin Li, Dongfang Zhao 0001, Michael Lang 0003, Ioan Raicu
IEEE BigData6
2014 Virtual chunks: On supporting random accesses to scientific data in compressible storage systems
abstract
Data compression could ameliorate the I/O pressure of scientific applications on high-performance computing systems. Unfortunately, the conventional wisdom of naively applying data compression to the file or block brings the dilemma between efficient random accesses and high compression ratios. Filelevel compression can barely support efficient random accesses to the compressed data: any retrieval request need trigger the decompression from the beginning of the compressed file. Block-level compression provides flexible random accesses to the compressed data, but introduces extra overhead when applying the compressor to each every block that results in a degraded overall compression ratio. This paper introduces a concept called virtual chunks aiming to support efficient random accesses to the compressed scientific data without sacrificing its compression ratio. In essence, virtual chunks are logical blocks identified by appended references without breaking the physical continuity of the file content. These additional references allow the decompression to start from an arbitrary position (efficient random access), and retain the file's physical entirety to achieve high compression ratio on par with file-level compression. One potential concern of virtual chunks lies on its space overhead (from the additional references) that degrades the compression ratio, but our analytic study and experimental results demonstrate that such overhead is negligible. We have implemented virtual chunks in two forms: a middleware to the GPFS parallel file system, and a module in the FusionFS distributed file system. Large-scale evaluations on up to 1,024 cores showed that virtual chunks could help improve the I/O throughput by 2X speedup.
Dongfang Zhao 0001, Jian Yin 0002, Kan Qiao, Ioan Raicu
IEEE BigData4
2014 FusionFS: Toward supporting data-intensive scientific applications on extreme-scale high-performance computing systems
abstract
State-of-the-art, yet decades-old, architecture of high-performance computing systems has its compute and storage resources separated. It thus is limited for modern data-intensive scientific applications because every I/O needs to be transferred via the network between the compute and storage resources. In this paper we propose an architecture that hss a distributed storage layer local to the compute nodes. This layer is responsible for most of the I/O operations and saves extreme amounts of data movement between compute and storage resources. We have designed and implemented a system prototype of this architecture - which we call the FusionFS distributed file system - to support metadata-intensive and write-intensive operations, both of which are critical to the I/O performance of scientific applications. FusionFS has been deployed and evaluated on up to 16K compute nodes of an IBM Blue Gene/P supercomputer, showing more than an order of magnitude performance improvement over other popular file systems such as GPFS, PVFS, and HDFS.
Dongfang Zhao 0001, Zhao Zhang 0007, Xiaobing Zhou, Tonglin Li, Ke Wang 0012, Dries Kimpe, Philip H. Carns, Robert B. Ross, Ioan Raicu
IEEE BigData9
2014 Towards In-Order and Exactly-Once Delivery Using Hierarchical Distributed Message Queues
abstract
In today's world, distributed message queues are used in many systems and play different roles (e.g. content delivery, notification system and message delivery tools). It is important for the queue services to be able to deliver messages at large scales with a variety of message sizes with high concurrency. An example of a commercial state of the art distributed message queue is Amazon Simple Queuing Service (SQS). SQS is a distributed message delivery fabric that is highly scalable. It can queue unlimited number of short messages (maximum size: 256 KB) and deliver them to multiple users in parallel. In order to be able to provide such high throughput at large scales, SQS omits some of features that are provided by traditional queues. SQS does not guarantee the order of the messages, nor does it guarantee the exactly once delivery. This paper addresses these limitations through the design and implementation of HDMQ, a hierarchical distributed message queue. HDMQ consist of collection of area message nodes that can be used to store messages up to 512 KB. It utilizes a round robin local load balancer to save the message and scale across the area region accordingly. HDMQ provides replication for high reliability of messages. It also provides SQS-like APIs in order to provide compatibility with current systems that currently use SQS. We performed a detailed performance evaluation and compared HDMQ to the commonly used commercial distributed queues measuring throughput, latency and price per request. We found HDMQ to outperform SQS, Windows Azure Service bus, and Iron MQ by up to 2-15x times in throughput, 1.6-39x times in latency, and all this for 13%-80% less costs.
Dharmit Patel, Faraj Khasib, Iman Sadooghi, Ioan Raicu
CCGRID4
2014 Exploring Infiniband Hardware Virtualization in OpenNebula towards Efficient High-Performance Computing
abstract
It has been widely accepted that software virtualization has a big negative impact on high-performance computing (HPC) application performance. This work explores the potential use of Infiniband hardware virtualization in an Open Nebula cloud towards the efficient support of MPI-based workloads. We have implemented, deployed, and tested an Infiniband network on the Fermi Cloud private Infrastructure-as-a-Service (IaaS) cloud. To avoid software virtualization towards minimizing the virtualization overhead, we employed a technique called Single Root Input/Output Virtualization (SRIOV). Our solution spanned modifications to the Linux's Hypervisor as well as the Open Nebula manager. We evaluated the performance of the hardware virtualization on up to 56 virtual machines connected by up to 8 DDR Infiniband network links, with micro-benchmarks (latency and bandwidth) as well as with a MPI-intensive application (the HPL Linpack benchmark).
Tiago Pais Pitta De Lacerda Ruivo, Gerard Bernabeu Altayo, Gabriele Garzoglio, Steven Timm, Seo-Young Noh, Ioan Raicu
CCGRID7
2014 Achieving Efficient Distributed Scheduling with Message Queues in the Cloud for Many-Task Computing and High-Performance Computing
abstract
Task scheduling and execution over large scale, distributed systems plays an important role on achieving good performance and high system utilization. Due to the explosion of parallelism found in today's hardware, applications need to perform over-decomposition to deliver good performance, this over-decomposition is driving job management systems' requirements to support applications with a growing number of tasks with finer granularity. Our goal in this work is to provide a compact, light-weight, scalable, and distributed task execution framework (Cloud Kon) that builds upon cloud computing building blocks (Amazon EC2, SQS, and Dynamo DB). Most of today's state-of-the-art job execution systems have predominantly Master/Slaves architectures, which have inherent limitations, such as scalability issues at extreme scales and single point of failures. On the other hand distributed job management systems are complex, and employ non-trivial load balancing algorithms to maintain good utilization. Cloud Kon is a distributed job management system that can support both HPC and MTC workloads with millions of tasks/jobs. We compare our work with other state-of-the-art job management systems including Sparrow and MATRIX. The results show that Cloud Kon delivers better scalability compared to other state-of-the-art systems for some metrics - all with a significantly smaller code-base (5%).
Iman Sadooghi, Sandeep Palur, Ajay Anthony, Isha Kapur, Karthik Belagodu, Pankaj Purandare, Kiran Ramamurty, Ke Wang 0012, Ioan Raicu
CCGRID9
2014 HyCache+: Towards Scalable High-Performance Caching Middleware for Parallel File Systems
abstract
The ever-growing gap between the computation and I/O is one of the fundamental challenges for future computing systems. This computation-I/O gap is even larger for modern large scale high-performance systems due to their state-of-the-art yet decades long architecture: the compute and storage resources form two cliques that are interconnected with shared networking infrastructure. This paper presents a distributed storage middleware, called HyCache+, right on the compute nodes, which allows I/O to effectively leverage the high bi-section bandwidth of the high-speed interconnect of massively parallel high-end computing systems. HyCache+ provides the POSIX interface to end users with the memory-class I/O throughput and latency, and transparently swap the cached data with the existing slow speed but high-capacity networked attached storage. HyCache+ has the potential to achieve both high performance and low cost large capacity, the best of both worlds. To further improve the caching performance from the perspective of the global storage system, we propose a 2-phase mechanism to cache the hot data for parallel applications, called 2-Layer Scheduling (2LS), which minimizes the file size to be transferred between compute nodes and heuristically replaces files in the cache. We deploy HyCache+ on the IBM Blue Gene/P supercomputer, and observe two orders of magnitude faster I/O throughput than the default GPFS parallel file system. Furthermore, the proposed heuristic caching approach shows 29X speedup over the traditional LRU algorithm.
Dongfang Zhao 0001, Kan Qiao, Ioan Raicu
CCGRID3
2014 Design and evaluation of the gemtc framework for GPU-enabled many-task computing
abstract
We present the design and first performance and usability evaluation of GeMTC, a novel execution model and runtime system that enables accelerators to be programmed with many concurrent and independent tasks of potentially short or variable duration. With GeMTC, a broad class of such "many-task" applications can leverage the increasing number of accelerated and hybrid high-end computing systems. GeMTC overcomes the obstacles to using GPUs in a many-task manner by scheduling and launching independent tasks on hardware designed for SIMD-style vector processing. We demonstrate the use of a high-level MTC programming model (the Swift parallel dataflow language) to run tasks on many accelerators and thus provide a high-productivity programming model for the growing number of supercomputers that are accelerator-enabled. While still in an experimental stage, GeMTC can already support tasks of fine (subsecond) granularity and execute concurrent heterogeneous tasks on 86,000 independent GPU warps spanning 2.7M GPU threads on the Blue Waters supercomputer.
Scott J. Krieder, Justin M. Wozniak, Timothy G. Armstrong, Michael Wilde, Daniel S. Katz, Benjamin Grimmer, Ian T. Foster, Ioan Raicu
HPDC8
2014 Next generation job management systems for extreme-scale ensemble computing
abstract
With the exponential growth of supercomputers in parallelism, applications are growing more diverse, including traditional large-scale HPC MPI jobs, and ensemble workloads such as finer-grained many-task computing (MTC) applications. Delivering high throughput and low latency for both workloads requires developing a distributed job management system that is magnitudes more scalable than today's centralized ones. In this paper, we present a distributed job launch prototype, SLURM++, which is comprised of multiple controllers with each one managing a partition of SLURM daemons, while ZHT (a distributed key-value store) is used to store the job and resource metadata. We compared SLURM++ with SLURM using micro-benchmarks of different job sizes up to 500 nodes, with excellent results showing 10X higher throughput. We also studied the potential of distributed scheduling through simulations up to millions of nodes.
Ke Wang 0012, Xiaobing Zhou, Michael Lang 0003, Ioan Raicu
HPDC5
2014 Devising a Cloud Scientific Workflow Platform for Big Data
abstract
Scientific workflow management systems (SWFMSs) are facing unprecedented challenges from big data deluge. As revising all the existing workflow applications to fit into Cloud computing paradigm is impractical, thus migrating SWFMSs into the Cloud to leverage the functionalities of both Cloud computing and SWFMSs may provide a viable approach to big data processing. In this paper, we first discuss the challenges for scientific workflow applications and the available solutions in details, and analyze the essential requirements for a scientific computing Cloud platform. Then we propose a service framework to normalize the integration of SWFMS with Cloud computing. Meanwhile, we also present our implementation experience based on the service Framework. At last, we set up a series of experiments to demonstrate the capability of our implementation and use a Montage Image Mosaic Workflow as a showcase of the implementation.
Yong Zhao 0009, Youfu Li 0002, Shiyong Lu, Ioan Raicu, Cui Lin
SERVICES4
2013 Optimizing Large Data Transfers over 100Gbps Wide Area Networks
abstract
The Advanced Networking Initiative (ANI) project from the Energy Services Network provides a 100 Gbps test bed, which offers the opportunity for evaluating applications and middleware used by scientific experiments. This test bed is a prototype of a 100 Gbps wide-area network backbone, which links several Department of Energy (DOE) national laboratories, universities and other research institutions. These scientific experiments involve movement of large datasets for collaborations among researchers at different sites and thus require advanced infrastructure for supporting large and fast data transfers. A 100 Gbps network test bed is a key component of the ANI project and is used for DOE's science research programs. This work presents results towards obtaining maximum throughput in large data transfers by optimizing and fine-tuning scientific applications and middleware to use this advanced infrastructure efficiently. A detailed performance evaluation is discussed measuring both applications, from High Energy Physics (HEP) and from data transfer middleware (GridFTP, Globus Online, Storage Resource Management, XrootD and Squid) at 100 Gbps speeds and 53 ms of latency. Results show that up to 97% efficiency of such high bandwidth high latency network is possible, achieving 80-90 Gbps in most test cases with a peak transfer rate of 100 Gbps.
Anupam Rajendran, Parag Mhashilkar, David Dykstra, Gabriele Garzoglio, Ioan Raicu
CCGRID6
2013 Towards high-performance and cost-effective distributed storage systems with information dispersal algorithms
abstract
Reliability is one of the most fundamental challenges for high performance computing (HPC) and cloud computing. Data replication is the de facto mechanism to achieve high reliability, even though it has been criticized for its high cost and low efficiency. Recent research showed promising results by switching the traditional data replication to a software-based RAID. In order to systematically study the effectiveness of this new method, we built two storage systems from the ground up: a POSIX-compliant distributed file system (FusionFS) and a distributed key-value store (IStore), both supporting information dispersal algorithms (IDA) for data redundancy. FusionFS is crafted to have excellent throughput and scalability for HPC, whereas IStore is architected mainly as a light-weight key-value storage in cloud computing. We evaluated both systems with a large number of parameter combinations. Results show that, for both HPC and cloud computing communities, IDA-based methods with current commodity hardware could outperform data replication in some cases, and would completely surpass data replication with the growing computational capacity through multi/many-core processors (e.g. Intel Xeon Phi, NVIDIA GPU).
Dongfang Zhao 0001, Kent Burlingame, Corentin Debains, Pedro Alvarez-Tabio, Ioan Raicu
CLUSTER5
2013 Distributed data provenance for large-scale data-intensive computing
abstract
It has become increasingly important to capture and understand the origins and derivation of data (its provenance). A key issue in evaluating the feasibility of data provenance is its performance, overheads, and scalability. In this paper, we explore the feasibility of a general metadata storage and management layer for parallel file systems, in which metadata includes both file operations and provenance metadata. We experimentally investigate the design optimality—whether provenance metadata should be loosely-coupled or tightly integrated with a file metadata storage systems. We consider two systems that have applied similar distributed concepts to metadata management, but focusing singularly on kind of metadata: (i) FusionFS, which implements a distributed file metadata management based on distributed hash tables, and (ii) SPADE, which uses a graph database to store audited provenance data and provides distributed module for querying provenance. Our results on a 32-node cluster show that FusionFS+SPADE is a promising prototype with negligible provenance overhead and has promise to scale to petascale and beyond. Furthermore, FusionFS with its own storage layer for provenance capture is able to scale up to 1K nodes on BlueGene/P supercomputer.
Dongfang Zhao 0001, Chen Shou, Tanu Malik, Ioan Raicu
CLUSTER4
2013 ZHT: A Light-Weight Reliable Persistent Dynamic Scalable Zero-Hop Distributed Hash Table
abstract
This paper presents ZHT, a zero-hop distributed hash table, which has been tuned for the requirements of high-end computing systems. ZHT aims to be a building block for future distributed systems, such as parallel and distributed file systems, distributed job management systems, and parallel programming systems. The goals of ZHT are delivering high availability, good fault tolerance, high throughput, and low latencies, at extreme scales of millions of nodes. ZHT has some important properties, such as being light-weight, dynamically allowing nodes join and leave, fault tolerant through replication, persistent, scalable, and supporting unconventional operations such as append (providing lock-free concurrent key/value modifications) in addition to insert/lookup/remove. We have evaluated ZHT's performance under a variety of systems, ranging from a Linux cluster with 512-cores, to an IBM Blue Gene/P supercomputer with 160K-cores. Using micro-benchmarks, we scaled ZHT up to 32K-cores with latencies of only 1.1ms and 18M operations/sec throughput. This work provides three real systems that have integrated with ZHT, and evaluate them at modest scales. 1) ZHT was used in the FusionFS distributed file system to deliver distributed meta-data management at over 60K operations (e.g. file create) per second at 2K-core scales. 2) ZHT was used in the IStore, an information dispersal algorithm enabled distributed object storage system, to manage chunk locations, delivering more than 500 chunks/sec at 32-nodes scales. 3) ZHT was also used as a building block to MATRIX, a distributed job scheduling system, delivering 5000 jobs/sec throughputs at 2K-core scales. We compared ZHT against other distributed hash tables and key/value stores and found it offers superior performance for the features and portability it supports.
Tonglin Li, Xiaobing Zhou, Kevin Brandstatter, Dongfang Zhao 0001, Ke Wang 0012, Anupam Rajendran, Zhao Zhang 0007, Ioan Raicu
IPDPS8
2013 Using simulation to explore distributed key-value stores for extreme-scale system services
abstract
Owing to the significant high rate of component failures at extreme scales, system services will need to be failure-resistant, adaptive and self-healing. A majority of HPC services are still designed around a centralized paradigm and hence are susceptible to scaling issues. Peer-to-peer services have proved themselves at scale for wide-area internet workloads. Distributed key-value stores (KVS) are widely used as a building block for these services, but are not prevalent in HPC services. In this paper, we simulate KVS for various service architectures and examine the design trade-offs as applied to HPC service workloads to support extreme-scale systems. The simulator is validated against existing distributed KVS-based services. Via simulation, we demonstrate how failure, replication, and consistency models affect performance at scale. Finally, we emphasize the general use of KVS to HPC services by feeding real HPC service workloads into the simulator and presenting a KVS-based distributed job launch prototype.
Ke Wang 0012, Abhishek Kulkarni, Michael Lang 0003, Dorian C. Arnold, Ioan Raicu
SC5
2012 ADAPT: Availability-Aware MapReduce Data Placement for Non-dedicated Distributed Computing
abstract
The MapReduce programming paradigm is gaining more and more popularity recently due to its merits of ease of programming, data distribution and fault tolerance. The low barrier of adoption of MapReduce makes it a promising framework for non-dedicated distributed computing environments. However, the variability of hosts resources and availability could substantially degrade the performance of MapReduce applications. The replication-based fault tolerance mechanism helps to alleviate some problems at the cost of inefficient storage space utilization. Intelligent solutions that guarantee the performance of MapReduce applications with low data replication degree are needed to promote the idea of running MapReduce applications in non-dedicated environment at lower costs. In this research, we propose an Availability-aware Data Placement (ADAPT) strategy to improve the application performance without extra storage cost. The basic idea of ADAPT is to dispatch data based on the availability of each node, reduce network traffic, improve data locality, and optimize the application performance. We implement the prototype of ADAPT within the Hadoop framework, an open-source implementation of MapReduce. The performance of ADAPT is evaluated in an emulated non-dedicated distributed environment. The experimental results show that ADAPT can improve the performance by more than 30%. ADAPT achieves high reliability without the need for additional data replication. ADAPT has also been evaluated for large-scale computing environment through simulations, with promising results.
Hui Jin 0001, Xi Yang 0002, Xian-He Sun, Ioan Raicu
ICDCS4
2012 Guest Editors' Introduction: Special Issue on Data-Intensive Computing in the Clouds
Tevfik Kosar, Ioan Raicu
J. Grid Comput.2
2011 Guest Editors' Introduction: Special Section on Many-Task Computing
abstract
IT is our honor to serve as guest editors of this special section of the IEEE Transactions on Parallel and Distributed Systems (TPDS) on many-task computing (MTC). This section focuses on the methods required to manage and execute large multiple program multiple data (MPMD) computations on large clusters, grids, clouds, and supercomputers. We are pleased to present 10 high-quality contributions chosen from 42 submissions, on resource management, data-intensive computing, applications, and MTC on supercomputers, grids, and clouds. We introduce the term many-task computing (MTC) [2] for computations that bridge the gap between high-performance computing (HPC) and high-throughput computing (HTC) [1]. MTC differs from HTC in its emphasis on using many computing resources over short periods of time to accomplish many computational tasks (both dependent and independent), for which primary metrics are measured in seconds (e.g., FLOPS, tasks/sec., MB/s I/O rates), as opposed to jobs per month. MTC computations comprise multiple distinct activities, coupled via files, shared memory, or message passing. Tasks may be small or large, uniprocessor or multiprocessor, or compute-intensive or data-intensive. The set of tasks may be static or dynamic, homogeneous or heterogeneous, or loosely coupled or tightly coupled. The number of tasks, quantity of computing, and volumes of data may be large. Today’sHPCsystemsareaviableplatformforMTC[3],but large MTC applications can stress HPC hardware and sotware. Challenges include local resource manager scalability and granularity, efficient utilization of raw hardware, parallel file system contention and scalability, data management, I/O management, reliability at scale, application scalability, and understanding the limitations of HPC systems in order to identify good candidate MTC applications [4]. MTC applications can also be executed on cloud systems, but face other challenges there, for example, relating to internode communication performance. Three recent MTC workshops (MTAGS, http://dsl.cs. uchicago.edu/MTAGS10/) and this special section attracted 142 abstracts and 110 paper submissions, from which 41 papers were accepted. Papers covered resource management, data-intensive computing, applications, and MTC on supercomputers, grids, and clouds. More than 1,000 people have participated as coauthors, program committee members, reviewers, and attendees in these venues. We are well beyond a critical mass for a new, thriving community, which is quickly expanding.
Ioan Raicu, Ian T. Foster, Yong Zhao 0009
IEEE Trans. Parallel Distributed Syst.1
2009 The quest for scalable support of data-intensive workloads in distributed systems
abstract
Data-intensive applications involving the analysis of large datasets often require large amounts of compute and storage resources, for which locality can be crucial to high throughput and performance. We propose a data approach that acquires compute and storage resources dynamically, replicates in response to demand, and schedules computations close to data. As demand increases, more resources are acquired, thus allowing faster response to subsequent requests that refer to the same data; when demand drops, resources are released. This approach can provide the benefits of dedicated hardware without the associated high costs, depending on workload and resource characteristics. To explore the feasibility of diffusion, we offer both a theoretical and an empirical analysis. We define an abstract model for diffusion, introduce new scheduling policies with heuristics to optimize real-world performance, and develop a competitive online cache eviction policy. We also offer many empirical experiments to explore the benefits of dynamically expanding and contracting resources based on load, to improve system responsiveness while keeping wasted resources small. We show performance improvements of one to two orders of magnitude across three diverse workloads when compared to the performance of parallel file systems with throughputs approaching 80 Gb/s on a modest cluster of 200 processors. We also compare diffusion with a best model for active storage, contrasting the difference between a pull-model found in diffusion and a push-model found in active storage.
Ioan Raicu, Ian T. Foster, Yong Zhao 0009, Philip Little, Christopher Moretti, Amitabh Chaudhary, Douglas Thain
HPDC1
2008 Toward loosely coupled programming on petascale systems
abstract
We have extended the Falkon lightweight task execution framework to make loosely coupled programming on petascale systems a practical and useful programming model. This work studies and measures the performance factors involved in applying this approach to enable the use of petascale systems by a broader user community, and with greater ease. Our work enables the execution of highly parallel computations composed of loosely coupled serial jobs with no modifications to the respective applications. This approach allows a new-and potentially far larger-class of applications to leverage petascale systems, such as the IBM Blue Gene/P supercomputer. We present the challenges of I/O performance encountered in making this model practical, and show results using both microbenchmarks and real applications from two domains: economic energy modeling and molecular dynamics. Our benchmarks show that we can scale up to 160 K processor-cores with high efficiency, and can achieve sustained execution rates of thousands of tasks per second.
Ioan Raicu, Zhao Zhang 0007, Michael Wilde, Ian T. Foster, Pete Beckman, Kamil Iskra, Ben Clifford
SC1
2007 Falkon: a Fast and Light-weight tasK executiON framework
abstract
To enable the rapid execution of many tasks on compute clusters, we have developed Falkon, a Fast and Light-weight tasK executiON framework. Falkon integrates (1) multi-level scheduling to separate resource acquisition (via, e.g., requests to batch schedulers) from task dispatch, and (2) a streamlined dispatcher. Falkon's integration of multi-level scheduling and streamlined dispatchers delivers performance not provided by any other system. We describe Falkon architecture and implementation, and present performance results for both microbenchmarks and applications. Microbenchmarks show that Falkon throughput (487 tasks/sec) and scalability (to 54,000 executors and 2,000,000 tasks processed in just 112 minutes) are one to two orders of magnitude better than other systems used in production Grids. Large-scale astronomy and medical applications executed under Falkon by the Swift parallel programming system achieve up to 90% reduction in end-to-end run time, relative to versions that execute tasks via separate scheduler submissions.
Ioan Raicu, Yong Zhao 0009, Catalin Dumitrescu, Ian T. Foster, Michael Wilde
SC1
2007 Usage SLA-based scheduling in Grids
abstract
Abstract Managing usage service level agreements (uSLAs) within environments that integrate participants and resources spanning multiple physical institutions is a challenging problem. Running workloads in such environments is often a similarly challenging problem owing to the scale of the environment, and to the resource partitioning based on various sharing strategies. Also, a resource may be taken down during a job execution, be improperly set up or fail job execution. Such elements have to be taken into account whenever targeting a Grid environment for problem solving. In this paper we explore uSLA‐based scheduling on a real Grid, Grid3, by means of a specific workload (the BLAST workload) and a specific scheduling framework, GRUBER (an architecture and toolkit for resource uSLA specification and enforcement). The paper provides extensive experimental results and comparisons with other scheduling strategies. We also address, in great detail, the performance of different uSLA‐based site selection strategies and the overall performance in scheduling workloads over Grid3 with workload sizes ranging from 10 to 10 000 jobs. Copyright © 2006 John Wiley & Sons, Ltd.
Catalin Dumitrescu, Ioan Raicu, Ian T. Foster
Concurr. Comput. Pract. Exp.2
2007 The Design, Usage, and Performance of GRUBER: A Grid Usage Service Level Agreement based BrokERing Infrastructure
Catalin Dumitrescu, Ioan Raicu, Ian T. Foster
J. Grid Comput.2
2006 Poster reception - Harnessing grid resources to enable the dynamic analysis of large astronomy datasets
abstract
Astronomy datasets are generally terabytes in size and contain hundreds of millions of objects separated into millions of files-factors which makes many analyses impractical to perform on small computers. The key question we answer in this paper is: How can we leverage Grid resources to make the analysis of large astronomy datasets a reality for the astronomy community? To address this question, we have developed a Web Services-based system, AstroPortal, that uses grid computing to federate large computing and storage resources for dynamic analysis of large datasets. Building on the GT4, we have built a prototype and implemented a first analysis, stacking, that sums multiple regions of the sky, a function that can help both identify variable sources and detect faint objects. AstroPortal gives the astronomy community a new tool to advance their research and to open new doors to opportunities never before possible on such a large scale.
Ioan Raicu, Ian T. Foster, Alex Szalay
SC1
2006 The Design, Performance, and Use of DiPerF: An automated DIstributed PERformance evaluation Framework
Ioan Raicu, Catalin Dumitrescu, Matei Ripeanu, Ian T. Foster
J. Grid Comput.1
2005 DI-GRUBER: A Distributed Approach to Grid Resource Brokering
abstract
Managing usage service level agreements (USLAs) within environments that integrate participants and resources spanning multiple physical institutions is a challenging problem. Maintaining a single unified USLA management decision point over hundreds to thousands of jobs and sites can become a bottleneck in terms of reliability as well as performance. DIGRUBER, an extension to our GRUBER brokering framework, was developed as a distributed grid USLAbased resource broker that allows multiple decision points to coexist and cooperate in real-time. DIGRUBER addresses issues regarding how USLAs can be stored, retrieved, and disseminated efficiently in a large distributed environment. The key question this paper addresses is the scalability and performance of DI-GRUBER in large Grid environments. We conclude that as little as three to five decision points can be sufficient in an environment with 300 sites and 60 VOs, an environment ten times larger than today’s Open Science Grid.
Catalin Dumitrescu, Ioan Raicu, Ian T. Foster
SC2