Ewa Deelman

dblp:d/EwaDeelman · DBLP profile ↗
← Back
127ranked-venue papers
17as first author
30since 2021 · last 2026
0000-0001-5106-503XORCID · conflict

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

Systems, architecture and hardware · 71 · 11 first-author · 14 since 2021Applied, interdisciplinary, general and emerging computing · 39 · 4 first-author · 13 since 2021Software engineering, systems software and programming languages · 37 · 4 first-author · 13 since 2021Artificial intelligence and machine learning · 9 · 2 since 2021Databases, data management, data science and information retrieval · 6 · 1 first-author · 2 since 2021Graphics, computer vision, multimedia, augmented reality and games · 2Computer networks · 1 · 1 since 2021Theory of computation · 1
YearPublicationVenuePosition
2026 AeroResQ: Edge-accelerated UAV framework for scalable, resilient and collaborative escape route planning in wildfire scenarios
Suman Raj, Radhika Mittal, Rajiv Mayani, Pawel Zuk, Anirban Mandal, Michael Zink, Yogesh L. Simmhan, Ewa Deelman
Future Gener. Comput. Syst.8
2026 A terminology for scientific workflow systems
Frédéric Suter, Tainã Coleman, Ilkay Altintas, Rosa M. Badia, Bartosz Balis, Kyle Chard, Iacopo Colonnelli, Ewa Deelman, Paolo Di Tommaso, Thomas Fahringer, Carole A. Goble, Shantenu Jha, Daniel S. Katz, Johannes Köster, Ulf Leser, Kshitij Mehta, Hilary Oliver, Jayson Luc Peterson, Giovanni Pizzi, Loïc Pottier, Raül Sirvent, Eric Suchyta, Douglas Thain, Sean R. Wilkinson, Justin M. Wozniak, Rafael Ferreira da Silva
Future Gener. Comput. Syst.8
2026 Advances in algorithms, models, hardware, and software for next-generations HPC systems, volume 2
Roman Wyrzykowski, Ewa Deelman
Future Gener. Comput. Syst.2
2025 A Greedy Consensus-Based Approach to Distributed Job Selection: Toward Fully-Decentralized Workload Management System
abstract
Current approaches to resilience for highly distributed, heterogeneous, large-scale scientific workflows are limited. Most existing workflow and resource management systems have a single point of failure and resilience strategies are often static, depend on a centralized control, and require considerable design effort from experts. The increasing scale and complexity of workflows coupled with limited resilience capabilities in centralized systems necessitates a fully decentralized, adaptive resource management approach. This paper addresses a very important slice of the overall problem by leveraging the advances in multi-agent systems (MAS). In particular, we explore the suitability of a MAS consisting of globally distributed agents to perform distributed job selection from a dynamic job pool in a truly decentralized, performant, and resilient manner. We present a novel consensus formulation of the distributed job selection problem. By introducing a cost function encapsulating the requirements and constraints of the job and resource loads, we design a novel, greedy consensus algorithm leveraging the Practical Byzantine Fault Tolerance (PBFT)-based consensus method, allowing agents to collectively select jobs in a resilient manner. We compared our algorithms with other state of the art approaches by deploying them in a network testbed infrastructure to emulate distributed job selection. Our evaluation results demonstrated that our greedy consensus algorithm employing the cost-function and PBFT-based consensus method outperforms the ones using the vanilla PBFT-based consensus method - improving scheduling latency by as much as 63.5 % and reducing resource idle time by as much as 63.8 %, with benefits increasing with higher numbers of agents emulated.
Komal Thareja, Raghavan Krishnan, Anirban Mandal, Pawel Zuk, Imtiaz Mahmud, Mariam Kiran, Ewa Deelman
CCGrid7
2025 Advancing anomaly detection in computational workflows with active learning
Raghavan Krishnan, George Papadimitriou 0002, Anirban Mandal, Mariam Kiran, Prasanna Balaprakash, Ewa Deelman
Future Gener. Comput. Syst.7
2024 DISTRI: Development and Integration of Simulation Tools for Resilient Infrastructure
abstract
In contemporary scientific research, data acquisition and analysis platforms have grown increasingly complex, often spanning multiple facilities with diverse internal structures. Efficiently managing the interactions between job scheduling, resource allocation, and networking across these distributed systems requires a robust simulation framework. However, existing simulators fall short in capturing the detailed interactions necessary for comprehensive analysis of large-scale distributed environments. To address this gap, we introduce DISTRI, a versatile framework specifically designed for the development and testing of distributed multi-facility workflows. DISTRI allows for customizable facility configurations and includes built-in support for distributed, resilient scheduling and resource management, alongside detailed network simulation for data communication. Key features of DISTRI encompass inter- and intra-facility resource management, agent-based distributed scheduling, and extensive performance metrics logging for both resource and network management. By providing these essential tools, DISTRI enables thorough analysis and optimization, thereby advancing research in the resilience and efficiency of multi-facility systems.
Imtiaz Mahmud, Pawel Zuk, Cong Wang 0014, Mariam Kiran, Kesheng Wu, Komal Thareja, Raghavan Krishnan, Anirban Mandal, Ewa Deelman
IEEE Big Data9
2024 Dynamic Tracking, MLOps, and Workflow Integration: Enabling Transparent Reproducibility in Machine Learning
abstract
Workflow management systems (WMS) provide a robust solution for automating and ensuring the reproducibility of scientific and engineering experiments. However, reproducing machine learning (ML) experiments requires replicating every aspect of the process, including code implementation, workflow execution, data, and the execution environment. Traditionally, tracking these components is done manually, if at all, before execution. In this work, we propose an approach for on-demand and dynamic tracking of ML workflows. Our approach extends the ML workflow automatically and introduces steps for tracking, organizing, and versioning all elements, such as code, data, the main workflow steps and the execution environment for each job. This tracking approach includes two modes: a custom mode, where user-tagged elements will be tracked by the WMS, and an automatic mode, where the WMS automatically tracks and organizes all necessary elements. We implemented this solution by extending the Pegasus-WMS system and tested it on two types of workflows: traditional scientific ML pipelines and Federated Learning (FL) applications. Our findings demonstrate that this tracking approach does not interfere with the normal execution of user-designed workflows and the execution time. Additionally, we show how this approach can integrate a WMS with versioned data in remote storage (such as S3 or Google Drive) and ML lifecycle solutions like MLflow, ensuring reproducibility and transparency of the computational experiments.
Hamza Safri, George Papadimitriou 0002, Ewa Deelman
e-Science3
2024 Large Language Models for Anomaly Detection in Computational Workflows: From Supervised Fine-Tuning to In-Context Learning
abstract
Anomaly detection in computational workflows is critical for ensuring system reliability and security. However, traditional rule-based methods struggle to detect novel anomalies. This paper leverages large language models (LLMs) for workflow anomaly detection by exploiting their ability to learn complex data patterns. Two approaches are investigated: (1) supervised fine-tuning (SFT), where pretrained LLMs are fine-tuned on labeled data for sentence classification to identify anomalies, and (2) in-context learning (ICL), where prompts containing task descriptions and examples guide LLMs in few-shot anomaly detection without fine-tuning. The paper evaluates the performance, efficiency, and generalization of SFT models and explores zeroshot and few-shot ICL prompts and interpretability enhancement via chain-of-thought prompting. Experiments across multiple workflow datasets demonstrate the promising potential of LLMs for effective anomaly detection in complex executions.
George Papadimitriou 0002, Raghavan Krishnan, Pawel Zuk, Prasanna Balaprakash, Cong Wang 0014, Anirban Mandal, Ewa Deelman
SC8
2024 Paving the way to hybrid quantum-classical scientific workflows
Sandeep Suresh Cranganore, Vincenzo De Maio, Ivona Brandic, Ewa Deelman
Future Gener. Comput. Syst.4
2024 Preface of special issue on advances in algorithms, models, hardware, and software for next-generations HPC systems
Roman Wyrzykowski, Ewa Deelman
Future Gener. Comput. Syst.2
2023 How is Artificial Intelligence Changing Science?
abstract
Artificial Intelligence (AI) is changing our lives. We talk with Alexa to find out about the weather or buy groceries, we use ChatGPT to help us be creative, Google Translate to communicate with people from other cultures, and Grammarly to fix spelling and grammar. AI has also entered the science arena where it is used to classify celestial galaxies, predict the weather, detect cancer cells, and design new materials among many others. Although AI has shown incredible results, for example creating proteins that have not been seen in nature, it raises a number of societal, ethical, and technological issues. AI is also challenging traditional science methods. In this paper, we explore how the scientific lifecycle is changing due to the growing capabilities of AI.
Ewa Deelman
e-Science1
2023 FlyPaw: Optimized Route Planning for Scientific UAVMissions
abstract
Many Internet of Things (IoT) applications require compute resources that cannot be provided by the devices themselves. On the other hand, processing of the data generated by IoT devices and sensors often has to be performed in real- or near real-time, i.e., with stringent latency requirements in constrained environments (e.g., intermittent network connectivity and limited power envelopes). Examples of such scenarios are autonomous vehicles in the form of cars and drones where the processing and analysis of observational data (e.g., video feeds) need to be performed expeditiously to allow for safe operation of the vehicles and to deliver the results in a timely fashion to the stakeholders of the mission. To support the compute and timeliness requirements of such applications, it is essential to include suitable edge resources to process these workflows, and to develop an end-to-end system that can route the vehicles dynamically and process and deliver mission-critical data and analyzed results. In this paper, we develop and evaluate a dynamic scheduling approach that considers complex tradeoffs between real-time constraints, network availability, and latency sensitivity of the mission. We devise an optimized route planning and data transmission schedule for drone flights. The scheduling algorithm is encapsulated in a novel end-to-end architecture (FlyPaw) and an associated adaptive drone mission control system, which enables deployment and management of an integrated cyberphysical system (CPS) – from real drone testbed to base stations to edge-to-cloud resources. The planning algorithm takes into account measured network communication characteristics, estimated uncertainties of future data link connectivity, and data timeliness requirements of the mission to prioritize candidate decision tree solutions based on a risk metric derived from Sharpe's ratio. Our results show that for given task sets, Net Time to Retrieve, our metric describing the time required to perform end-to-end collection and downstream processing of data, can be significantly reduced compared to other naive approaches. The theoretical improvement provided by our algorithm over other naive approaches is dependent on several factors — task locations, network connectivity, processing times and available resources, and is bounded by the duration of the drone flight.
Andrew Grote, Eric Lyons 0001, Komal Thareja, George Papadimitriou 0002, Ewa Deelman, Anirban Mandal, Prasad Calyam, Michael Zink
e-Science5
2023 Online Boosted Gaussian Learners for In-Situ Detection and Characterization of Protein Folding States in Molecular Dynamics Simulations
abstract
Molecular Dynamics (MD) simulations are a crucial tool for understanding how proteins fold. In its easiest form, MD simulations can be scaled through data parallelism, this means that multiple folding trajectories can be spawned and executed in parallel, facilitating a more efficient exploration of the protein folding space. However, due to data dependencies, the analysis of MD simulations remains largely as a centralized process. In this work, we propose a data parallel, lightweight technique to learn the characteristics of protein folding states in MD simulations. Contrary to other methods, ours can differentiate relevant states in a single protein folding trajectory without requiring centralized global knowledge of the protein dynamics. As its processing and memory overheads are negligible (in the order of milliseconds per window of frames, and kilo bytes respectively) this technique can be coupled with the simulation for in-situ analysis.
Harshita Sahni, Hector Carrillo-Cabada, Ekaterina D. Kots, Silvina Caíno-Lores, Jack D. Marquez, Ewa Deelman, Michel A. Cuendet, Harel Weinstein, Michela Taufer, Trilce Estrada
e-Science6
2023 Performance assessment of ensembles of in situ workflows under resource constraints
abstract
Summary Scientific breakthroughs in biomolecular methods and improvements in hardware technology have shifted from a long‐running simulation to a large set of shorter simulations running simultaneously, called an ensemble. In an ensemble, simulations are usually coupled with analyses of data produced by the simulations. In situ methods can be used to analyze large volumes of data generated by scientific simulations at runtime (i.e., simulations and analyses are performed concurrently). In this work, we study the execution of ensemble‐based simulations paired with in situ analyses using in‐memory staging methods. Using an ensemble of molecular dynamics in situ workflows with multiple simulations and analyses, we first show that collecting traditional metrics such as makespan, instructions per cycle, memory usage, or cache miss ratio is not sufficient to characterize complex behaviors of ensembles. We propose a method to evaluate the performance of ensembles of workflows that captures multiple resource usage aspects: resource efficiency, resource allocation, and resource provisioning. Experimental results demonstrate that the proposed method can effectively distinguish the performance of different component placements in an ensemble with up to 32 ensemble members. By evaluating different co‐location scenarios, our proposed performance indicators demonstrate benefits of co‐locating simulation and coupled analyses within a compute node.
Tu Mai Anh Do, Loïc Pottier, Rafael Ferreira da Silva, Silvina Caíno-Lores, Michela Taufer, Ewa Deelman
Concurr. Comput. Pract. Exp.6
2022 Accelerating Scientific Workflows on HPC Platforms with In Situ Processing
abstract
Scientific workflows drive most modern large-scale science breakthroughs by allowing scientists to define their computations as a set of jobs executed in a given order based on their data dependencies. Workflow management systems (WMSs) have become key to automating scientific workflows-executing computational jobs and orchestrating data transfers between those jobs running on complex high-performance computing (HPC) platforms. Traditionally, WMSs use files to communicate between jobs: a job writes out files that are read by other jobs. However, HPC machines face a growing gap between their storage and compute capabilities. To address that concern, the scientific community has adopted a new approach called in situ, which bypasses costly parallel filesystem I/O operations with faster in-memory or in-network communications. When using in situ approaches, communication and computations can be interleaved. In this work, we leverage the Decaf in situ dataflow framework to accelerate task-based scientific workflows managed by the Pegasus WMS, by replacing file communications with faster MPI messaging. We propose a new execution engine that uses Decaf to manage communications within a sub-workflow (i.e., set of jobs) to optimize inter-job communications. We consider two workflows in this study: (i) a synthetic workflow that benchmarks and compares file- and MPI-based communication; and (ii) a realistic bioinformatics workflow that computes mu-tational overlaps in the human genome. Experiments show that in situ communication can improve the bioinformatics workflow execution time by 22% to 30% compared with file communication. Our results motivate further opportunities and challenges for bridging traditional WMSs with in situ frameworks.
Tu Mai Anh Do, Loïc Pottier, Orcun Yildiz, Karan Vahi, Patrycja Krawczuk, Tom Peterka, Ewa Deelman
CCGRID7
2022 Automating Edge-to-cloud Workflows for Science: Traversing the Edge-to-cloud Continuum with Pegasus
abstract
In this paper, we describe how we extended the Pegasus Workflow Management System to support edge-to-cloud workflows in an automated fashion. We discuss how Pegasus and HTCondor (its job scheduler) work together to enable this automation. We use HTCondor to form heterogeneous pools of compute resources and Pegasus to plan the workflow onto these resources and manage containers and data movement for executing workflows in hybrid edge-cloud environments. We then show how Pegasus can be used to evaluate the execution of workflows running on edge only, cloud only, and edge-cloud hybrid environments. Using the Chameleon Cloud testbed to set up and configure an edge-cloud environment, we use Pegasus to benchmark the executions of one synthetic workflow and two production workflows: CASA-Wind and the Ocean Observatories Initiative Orcasound workflow, all of which derive their data from edge devices. We present the performance impact on workflow runs of job and data placement strategies employed by Pegasus when configured to run in the above three execution environments. Results show that the synthetic workflow performs best in an edge only environment, while the CASA - Wind and Orcasound workflows see significant improvements in overall makespan when run in a cloud only environment. The results demonstrate that Pegasus can be used to automate edge-to-cloud science workflows and the workflow provenance data collection capabilities of the Pegasus monitoring daemon enable computer scientists to conduct edge-to-cloud research.
Ryan Tanaka, George Papadimitriou 0002, Sai Charan Viswanath, Cong Wang 0014, Eric Lyons 0001, Komal Thareja, Chengyi Qu, Alicia Esquivel Morel, Ewa Deelman, Anirban Mandal, Prasad Calyam, Michael Zink
CCGRID9
2022 Application of Edge-to-Cloud Methods Toward Deep Learning
abstract
Scientific workflows are important in modern computational science and are a convenient way to represent complex computations, which are often geographically distributed among several computers. In many scientific domains, scientists use sensors (e.g., edge devices) to gather data such as CO2 level or temperature, that are usually sent to a central processing facility (e.g., a cloud). However, these edge devices are often not powerful enough to perform basic computations or machine learning inference computations and thus applications need the power of cloud platforms to generate scientific results. This work explores the execution and deployment of a complex workflow on an edge-to-cloud architecture in a use case of the detection and classification of plankton. In the original application, images were captured by cameras attached to buoys floating in Lake Greifensee (Switzerland). We developed a workflow based on that application. The workflow aims to pre-process images locally on the edge devices (i.e., buoys) then transfer data from each edge device to a cloud platform. Here, we developed a Pegasus workflow that runs using HTCondor and leveraged the Chameleon cloud platform and its recent CHI@Edge feature to mimic such deployment and study its feasibility in terms of performance and deployment.
Khushi Choudhary, Nona Nersisyan, Edward Lin, Shobana Chandrasekaran, Rajiv Mayani, Loïc Pottier, Angela P. Murillo, Nicole K. Virdone, Kerk F. Kee, Ewa Deelman
e-Science10
2022 Molecular Dynamics Workflow Decomposition for Hybrid Classic/Quantum Systems
abstract
Since we are entering the Post-Moore Law era and consequently the limit of Von Neumann's architecture, the scientific community is looking for alternatives to satisfy the growing computing power demands of scientific applications. Quantum computing promises to achieve a computational advantage over the classic Von Neumann architecture. However, the limited capabilities of current noisy intermediate-scale quantum (NISQ) devices require quantum computers to interoperate with classic systems, forming the so-called hybrid quantum systems. Research on hybrid quantum systems led to the design of Variational Quantum Algorithms, currently the most promising way to move towards quantum advantage. However, execution time and accuracy of variational quantum algorithms are affected by different hyperparameters, including selected cost functions and parametrized quantum circuits. Consequently, providing developers with methods to select the right set of parameters is of paramount importance. In this work, we provide a formal method for the selection of hyperparameters in variational quantum algorithms, which will support quantum algorithms developers in the design of quantum applications, and evaluate it on a real-world scientific application, showing a reduction of error up to 31%.
Sandeep Suresh Cranganore, Vincenzo De Maio, Ivona Brandic, Tu Mai Anh Do, Ewa Deelman
e-Science5
2022 SIM-SITU: A Framework for the Faithful Simulation of in situ Processing
abstract
The amount of data generated by numerical simulations in various scientific domains led to a fundamental redesign of how the analysis and visualization of simulation outputs are performed. The throughput and capacity of storage subsystems have not evolved as fast as the computing power in extreme-scale supercomputers, making the classical post-hoc approach highly inefficient. In situ processing has then emerged as a solution in which simulation and data analysis/visualization are intertwined for better performance and greater interactivity. Determining the best allocation, i.e., how many resources to allocate to simulation and analysis respectively, mapping, i.e., where and at which frequency to run the analysis/visualization, and data transfer mode is a complex task whose performance assessment is crucial to the efficient execution of in situ processing. However, such a performance evaluation of different strategies usually relies either on directly running them on the targeted execution environments, which can rapidly become extremely time- and resource-consuming, or on resorting to simplified models of the components of an in situ application, which can lack of realism. In both cases, the validity of the performance evaluation is limited. In this paper, we present SIM-SITU, a framework for the faithful performance evaluation of in situ processing strategies. We designed SIM-SITU to reflect the typical features of in situ processing systems. Thanks to its modular design, Sim-Situ has the necessary flexibility to easily and faithfully evaluate the behavior and performance of various allocation, mapping, and data transfer strategies. We illustrate the capabilities of SIM-SITU on a Molecular Dynamics use case. We study the impact of different strategies on performance and show how users can leverage SIM-SITU to determine interesting tradeoffs when adding analysis/visualization components to their application.
Valentin Honoré, Tu Mai Anh Do, Loïc Pottier, Rafael Ferreira da Silva, Ewa Deelman, Frédéric Suter
e-Science5
2022 Data Integrity Error Localization in Networked Systems with Missing Data
abstract
Most recent network failure diagnosis systems focused on data center networks where complex measurement systems can be deployed to derive routing information and ensure network coverage in order to achieve accurate and fast fault localization. In this paper, we target wide-area networks that support data-intensive distributed applications. We first present a new multi-output prediction model that directly maps the application level observations to localize the system component failures. In reality, this application-centric approach may face the missing data challenge as some input (feature) data to the inference models may be missing due to incomplete or lost measurements in wide area networks. We show that the presented prediction model naturally allows the multivariate imputation to recover the missing data. We evaluate multiple imputation algorithms and show that the prediction performance can be improved significantly in a large-scale network. As far as we know, this is the first study on the missing data issue and applying imputation techniques in network failure localization.
Yufeng Xin, Shih-Wen Fu, Anirban Mandal, Ryan Tanaka, Mats Rynge, Karan Vahi, Ewa Deelman
ICC7
2022 WfCommons: A framework for enabling scientific workflow research and development
Tainã Coleman, Henri Casanova, Loïc Pottier, Manav Kaushik, Ewa Deelman, Rafael Ferreira da Silva
Future Gener. Comput. Syst.5
2021 A Roadmap to Robust Science for High-throughput Applications: The Developers' Perspective
abstract
Scientists using the high-throughput computing (HTC) paradigm for scientific discovery rely on complex software systems and heterogeneous architectures that must deliver robust science (i.e., ensuring performance scalability in space and time; trust in technology, people, and infrastructures; and reproducible or confirmable research). Developers must overcome a variety of obstacles to pursue workflow interoperability, identify tools and libraries for robust science, port codes across different architectures, and establish trust in non-deterministic results. This poster presents recommendations to build a roadmap to overcome these challenges and enable robust science for HTC applications and workflows. The findings were collected from an international community of software developers during a Virtual World Cafe in May 2021.
Michela Taufer, Ewa Deelman, Rafael Ferreira da Silva, Trilce Estrada, Mary W. Hall, Miron Livny
CLUSTER2
2021 Serverless Containers - Rising Viable Approach to Scientific Workflows
abstract
The increasing popularity of the serverless computing approach has led to the emergence of new cloud infrastructures working in Container-as-a-Service (CaaS) model like AWS Fargate, Google Cloud Run, or Azure Container Instances. New infrastructures facilitate an innovative approach to running cloud containers where developers are freed from managing underlying resources. In this paper, we focus on evaluating the capabilities of elastic containers and their usefulness for scientific computing in the scientific workflow paradigm using AWS Fargate and Google Cloud Run infrastructures. For the experimental evaluation of our approach, we extended the HyperFlow engine to support these CaaS platforms, together with adapting four scientific workflows composed of several dozen to hundreds of tasks organized into a dependency graph. Studied applications are used to create cost-performance benchmarks and flow execution plots, delay, elasticity, and scalability measurements. Results show that serverless containers can be successfully utilized for running scientific workflows. Moreover, the results allow for gaining insight into the specific advantages and limits of the studied platforms.
Krzysztof Burkat, Maciej Pawlik, Bartosz Balis, Maciej Malawski, Karan Vahi, Mats Rynge, Rafael Ferreira da Silva, Ewa Deelman
e-Science8
2021 A Case Study in Scientific Reproducibility from the Event Horizon Telescope (EHT)
abstract
This poster presents the first results of an interdisciplinary project aiming to develop and share sustainable knowledge necessary to analyze, understand, and use published scientific results to advance reproducibility in multi-messenger astrophysics. Specifically, the project targets breakthrough work associated with the First M87 Event Horizon Telescope (EHT) and delivers recommendations on how the published results of the first black hole can be effectively reproduced. The project has the potential to advance new discovery in multi-messenger astrophysics by providing guidance for generalizing methods and findings from use cases.
Ross Ketron, Jacob Leonard, Brandan Roachell, Ria Patel, R. White, Silvina Caíno-Lores, Nigel Tan, Patrick R. Miles, Karan Vahi, Ewa Deelman, Duncan A. Brown, Michela Taufer
e-Science10
2021 CrisisFlow: Multimodal Representation Learning Workflow for Crisis Computing
abstract
An increasing number of people use social media (SM) platforms like Twitter and Instagram to report critical emergencies or disaster events. Multimodal data shared on these platforms often contain useful information about the scale of the event, victims, and infrastructure damage. The data can provide local authorities and humanitarian organizations with a big-picture understanding of the emergency (situational awareness). Moreover, it can be used to effectively and timely plan relief responses. In our project, we aim to address the challenge of finding relevant information among the vast amount of published SM posts. Specifically, we use deep learning algorithms to produce embeddings that encode the informativeness of multimodal SM data in the context of disaster events. Our method improves the state-of-the-art performance on the informative vs. non-informative classification task for the CrisisMMD dataset. To ensure the reliability and scalability of our solution in real-world scenarios, we implement the resulting crisis computing workflow in the Pegasus Workflow Management System (WMS).
Patrycja Krawczuk, Shubham Nagarkar, Ewa Deelman
e-Science3
2021 Predicting Flash Floods in the Dallas-Fort Worth Metroplex Using Workflows and Cloud Computing
abstract
Accurate and timely prediction of flash flooding events can be a very useful tool for stormwater officials and first responders. Having lead time with which to issue evacuation directives, to close flood prone roadways, to deploy rescue gear and personnel, and to fortify areas against flooding is essential to minimize property damage and risk of casualties. In this poster, we are presenting a flash flooding prediction workflow based on the Hydrology Lab-Research Distributed Hydrologic Model (HL-RDHM). This workflow leverages cloud computing and the Pegasus Workflow Management System to provide continuous high resolution flood predictions for the Dallas-Fort Worth Metroplex area in North Texas, and can be easily expanded to other regions.
Eric Lyons 0001, Dong-Jun Seo, Sunghee Kim, Hamideh Habibi, George Papadimitriou 0002, Ryan Tanaka, Ewa Deelman, Michael Zink, Anirban Mandal
e-Science7
2021 A Roadmap to Robust Science for High-throughput Applications: The Scientists' Perspective
abstract
This poster presents our first steps to define a roadmap to robust science for high-throughput applications used in scientific discovery. These applications combine multiple components into increasingly complex multi-modal workflows that are often executed in concert on heterogeneous systems. The increasing complexity hinders the ability of scientists to generate robust science (i.e., ensuring performance scalability in space and time; trust in technology, people, and infrastructures; and reproducible or confirmable research). Scientists must withstand and overcome adverse conditions such as heterogeneous and unreliable architectures at all scales (including extreme scale), rigorous testing under uncertainties, unexplainable algorithms in machine learning, and black-box methods. This poster presents findings and recommendations to build a roadmap to overcome these challenges and enable robust science. The data was collected from an international community of scientists during a virtual world café in February 2021.
Michela Taufer, Ewa Deelman, Rafael Ferreira da Silva, Trilce Estrada, Mary W. Hall
e-Science2
2021 Mining Workflows for Anomalous Data Transfers
abstract
Modern scientific workflows are data-driven and are often executed on distributed, heterogeneous, high-performance computing infrastructures. Anomalies and failures in the work-flow execution cause loss of scientific productivity and inefficient use of the infrastructure. Hence, detecting, diagnosing, and mitigating these anomalies are immensely important for reliable and performant scientific workflows. Since these workflows rely heavily on high-performance network transfers that require strict QoS constraints, accurately detecting anomalous network performance is crucial to ensure reliable and efficient workflow execution. To address this challenge, we have developed X-FLASH, a network anomaly detection tool for faulty TCP workflow transfers. X-FLASH incorporates novel hyperparameter tuning and data mining approaches for improving the performance of the machine learning algorithms to accurately classify the anomalous TCP packets. X-FLASH leverages XGBoost as an ensemble model and couples XGBoost with a sequential optimizer, FLASH, borrowed from search-based Software Engineering to learn the optimal model parameters. X-FLASH found configurations that outperformed the existing approach up to 28%, 29%, and 40% relatively for F-measure, G-score, and recall in less than 30 evaluations. From (1) large improvement and (2) simple tuning, we recommend future research to have additional tuning study as a new standard, at least in the area of scientific workflow anomaly detection.
Huy Tu, George Papadimitriou 0002, Mariam Kiran, Cong Wang 0014, Anirban Mandal, Ewa Deelman, Tim Menzies
MSR6
2021 End-to-end online performance data capture and analysis for scientific workflows
George Papadimitriou 0002, Cong Wang 0014, Karan Vahi, Rafael Ferreira da Silva, Anirban Mandal, Zhengchun Liu, Rajiv Mayani, Mats Rynge, Mariam Kiran, Vickie E. Lynch, Rajkumar Kettimuthu, Ewa Deelman, Jeffrey S. Vetter, Ian T. Foster
Future Gener. Comput. Syst.12
2021 Artificial Intelligence for Modeling Complex Systems: Taming the Complexity of Expert Models to Improve Decision Making
abstract
Major societal and environmental challenges involve complex systems that have diverse multi-scale interacting processes. Consider, for example, how droughts and water reserves affect crop production and how agriculture and industrial needs affect water quality and availability. Preventive measures, such as delaying planting dates and adopting new agricultural practices in response to changing weather patterns, can reduce the damage caused by natural processes. Understanding how these natural and human processes affect one another allows forecasting the effects of undesirable situations and study interventions to take preventive measures. For many of these processes, there are expert models that incorporate state-of-the-art theories and knowledge to quantify a system's response to a diversity of conditions. A major challenge for efficient modeling is the diversity of modeling approaches across disciplines and the wide variety of data sources available only in formats that require complex conversions. Using expert models for particular problems requires integration of models with third-party data as well as integration of models across disciplines. Modelers face significant heterogeneity that requires resolving semantic, spatiotemporal, and execution mismatches, which are largely done by hand today and may take more than 2 years of effort. We are developing a modeling framework that uses artificial intelligence (AI) techniques to reduce modeling effort while ensuring utility for decision making. Our work to date makes several innovative contributions: (1) an intelligent user interface that guides analysts to frame their modeling problem and assists them by suggesting relevant choices and automating steps along the way; (2) semantic metadata for models, including their modeling variables and constraints, that ensures model relevance and proper use for a given decision-making problem; and (3) semantic representations of datasets in terms of modeling variables that enable automated data selection and data transformations. This framework is implemented in the MINT (Model INTegration) framework, and currently includes data and models to analyze the interactions between natural and human systems involving climate, water availability, agricultural production, and markets. Our work to date demonstrates the utility of AI techniques to accelerate modeling to support decision-making and uncovers several challenging directions for future work.
Yolanda Gil, Daniel Garijo, Deborah Khider, Craig A. Knoblock, Varun Ratnakar, Maximiliano Osorio, Hernán Vargas, Minh Pham 0004, Jay Pujara, Basel Shbita, Yao-Yi Chiang, Dan Feldman, Yijun Lin 0001, Hayley Song, Vipin Kumar 0001, Ankush Khandelwal, Michael S. Steinbach, Kshitij Tayal, Shaoming Xu, Suzanne A. Pierce, Lissa Pearson, Daniel Hardesty-Lewis, Ewa Deelman, Rafael Ferreira da Silva, Rajiv Mayani, Armen R. Kemanian, Lorne Leonard, Scott D. Peckham, Maria Stoica 0001, Kelly M. Cobourn, Zeya Zhang, Christopher J. Duffy, Lele Shu
ACM Trans. Interact. Intell. Syst.24
2020 Modeling the Performance of Scientific Workflow Executions on HPC Platforms with Burst Buffers
abstract
Scientific domains ranging from bioinformatics to astronomy and earth science rely on traditional high-performance computing (HPC) codes, often encapsulated in scientific workflows. In contrast to traditional HPC codes that employ a few programming and runtime approaches that are highly optimized for HPC platforms, scientific workflows are not necessarily optimized for these platforms. As an effort to reduce the gap between compute and I/O performance, HPC platforms have adopted intermediate storage layers known as burst buffers. A burst buffer (BB) is a fast storage layer positioned between the global parallel file system and the compute nodes. Two designs currently exist: (i) shared, where the BBs are located on dedicated nodes; and (ii) on-node, in which each compute node embeds a private BB. In this paper, using accurate simulations and realworld experiments, we study how to best use these new storage layers when executing scientific workflows. These applications are not necessarily optimized to run on HPC systems, and thus can exhibit I/O patterns that differ from that of HPC codes. Thus, we first characterize the I/O behaviors of a real-world workflow under different configuration scenarios on two leadership-class HPC systems (Cori at NERSC and Summit at ORNL). Then, we use these characterizations to calibrate a simulator for workflow executions on HPC systems featuring shared and private BBs. Last, we evaluate our approach against a large I/O-intensive workflow, and we provide insights on the performance levels and the potential limitations of these two BBs architectures.
Loïc Pottier, Rafael Ferreira da Silva, Henri Casanova, Ewa Deelman
CLUSTER4
2020 Detecting anomalous packets in network transfers: investigations using PCA, autoencoder and isolation forest in TCP
Mariam Kiran, Cong Wang 0014, George Papadimitriou 0002, Anirban Mandal, Ewa Deelman
Mach. Learn.5
2020 The Workflow Trace Archive: Open-Access Data From Public and Private Computing Infrastructures
abstract
Realistic, relevant, and reproducible experiments often need input traces collected from real-world environments. In this work, we focus on traces of workflows-common in datacenters, clouds, and HPC infrastructures. We show that the state-of-the-art in using workflow-traces raises important issues: (1) the use of realistic traces is infrequent and (2) the use of realistic, open-access traces even more so. Alleviating these issues, we introduce the Workflow Trace Archive (WTA), an open-access archive of workflow traces from diverse computing infrastructures and tooling to parse, validate, and analyze traces. The WTA includes > 48 million workflows captured from > 10 computing infrastructures, representing a broad diversity of trace domains and characteristics. To emphasize the importance of trace diversity, we characterize the WTA contents and analyze in simulation the impact of trace diversity on experiment results. Our results indicate significant differences in characteristics, properties, and workflow structures between workload sources, domains, and fields.
Laurens Versluis, Roland Mathá, Sacheendra Talluri, Tim Hegeman, Radu Prodan, Ewa Deelman, Alexandru Iosup
IEEE Trans. Parallel Distributed Syst.6
2019 Exploration of Workflow Management Systems Emerging Features from Users Perspectives
abstract
There has been a recent emergence of new workflow applications focused on data analytics and machine learning. This emergence has precipitated a change in the workflow management landscape, causing the development of new dataoriented workflow management systems (WMSs) in addition to the earlier standard of task-oriented WMSs. In this paper, we summarize three general workflow use-cases and explore the unique requirements of each use-case in order to understand how WMSs from both workflow management models meet the requirements of each workflow use-case from the user’s perspective. We analyze the applicability of the two models by carefully describing each model and by providing an examination of the different variations of WMSs that fall under the task driven model. To illustrate the strengths and weaknesses of each workflow management model, we summarize the key features of four production-ready WMSs: Pegasus, Makeflow, Apache Airflow, and Pachyderm. To deepen our analysis of the four WMSs examined in this paper,we implement three real-world use-cases to highlight the specifications and features of each WMS. We present our final assessment of each WMS after considering the following factors: usability, performance, ease of deployment, and relevance. The purpose of this work is to offer insights from the user’s perspective into the research challenges that WMSs currently face due to the evolving workflow landscape.
Ryan Mitchell, Loïc Pottier, Steve Jacobs, Rafael Ferreira da Silva, Mats Rynge, Karan Vahi, Ewa Deelman
IEEE BigData7
2019 Empowering Agroecosystem Modeling with HTC Scientific Workflows: The Cycles Model Use Case
abstract
Scientific workflows have enabled large-scale scientific computations and data analysis, and lowered the entry barrier for performing computations in distributed heterogeneous platforms (e.g., HTC and HPC). In spite of impressive achievements to date, large-scale modeling, simulation, and data analytics in the long-tail still face several challenges such as efficient scheduling and execution of large-scale workflows (O(106)) with very short-running tasks (few seconds). While the current trend to support next-generation workflows on leadership class machines have gained much attention in the past years, at the other end of the spectrum scientific workflows from the long-tail science have become larger and require processing massive volumes of data. In this paper, we report on our experience in designing and implementing an HTC workflow for agroecosystem modeling. We leverage well-known (task clustering and co-scheduling) and emerging (hierarchical workflows and containers) workflow optimization techniques to make the workflow planning problem tractable, and maximize resource utilization and the degree of task parallelism. Experimental results, via the implementation of a use case, show that by strategically combining the above strategies and defining an appropriate set of optimization parameters, the overall workflow makespan can be improved by 3.5 orders of magnitude when compared to a regular (non-optimized) execution of the workflow.
Rafael Ferreira da Silva, Rajiv Mayani, Armen R. Kemanian, Mats Rynge, Ewa Deelman
IEEE BigData6
2019 Cyberinfrastructure Center of Excellence Pilot: Connecting Large Facilities Cyberinfrastructure
abstract
The National Science Foundation's Large Facilities are major, multi-user research facilities that operate and manage sophisticated and diverse research instruments and platforms (e.g., large telescopes, interferometers, distributed sensor arrays) that serve a variety of scientific disciplines, from astronomy and physics to geology and biology and beyond. Large Facilities are increasingly dependent on advanced cyberinfrastructure (i.e., computing, data, and software systems; networking; and associated human capital) to enable the broad delivery and analysis of facility-generated data. These cyberinfrastructure tools enable scientists and the public to gain new insights into fundamental questions about the structure and history of the universe, the world we live in today, and how our environment may change in the coming decades. This paper describes a pilot project that aims to develop a model for a Cyberinfrastructure Center of Excellence (CI CoE) that facilitates community building and knowledge sharing and that disseminates and applies best practices and innovative solutions for facility CI.
Ewa Deelman, Ryan Mitchell, Loïc Pottier, Mats Rynge, Erik Scott, Karan Vahi, Marina Kogan, Jasmine Mann, Tom Gulbransen, Daniel Allen, David Barlow, Anirban Mandal, Santiago Bonarrigo, Chris Clark, Leslie Goldman, Tristan Goulden, Phil Harvey, David Hulsander, Steve Jacobs, Christine Laney, Ivan Lobo-Padilla, Jeremy Sampson, Valerio Pascucci, John Staarmann, Steve Stone, Susan Sons, Jane Wyngaard, Charles Vardeman, Steve Petruzza, Ilya Baldin, Laura Christopherson
eScience1
2019 Toward a Dynamic Network-Centric Distributed Cloud Platform for Scientific Workflows: A Case Study for Adaptive Weather Sensing
abstract
Computational science today depends on complex, data-intensive applications operating on datasets from a variety of scientific instruments. A major challenge is the integration of data into the scientist's workflow. Recent advances in dynamic, networked cloud resources provide the building blocks to construct reconfigurable, end-to-end infrastructure that can increase scientific productivity. However, applications have not adequately taken advantage of these advanced capabilities. In this work, we have developed a novel network-centric platform that enables high-performance, adaptive data flows and coordinated access to distributed cloud resources and data repositories for atmospheric scientists. We demonstrate the effectiveness of our approach by evaluating time-critical, adaptive weather sensing workflows, which utilize advanced networked infrastructure to ingest live weather data from radars and compute data products used for timely response to weather events. The workflows are orchestrated by the Pegasus workflow management system and were chosen because of their diverse resource requirements. We show that our approach results in timely processing of Nowcast workflows under different infrastructure configurations and network conditions. We also show how workflow task clustering choices affect throughput of an ensemble of Nowcast workflows with improved turnaround times. Additionally, we find that using our network-centric platform powered by advanced layer2 networking techniques results in faster, more reliable data throughput, makes cloud resources easier to provision, and the workflows easier to configure for operational use and automation.
Eric Lyons 0001, Anirban Mandal, George Papadimitriou 0002, Cong Wang 0014, Komal Thareja, Paul Ruth, Juan J. Villalobos, Ivan Rodero, Ewa Deelman, Michael Zink
eScience9
2019 Characterizing In Situ and In Transit Analytics of Molecular Dynamics Simulations for Next-Generation Supercomputers
abstract
Molecular Dynamics (MD) simulations executed on state-of-the-art supercomputers are producing data at rates faster than it can be written out to disk. In situ and in transit analysis of data generated by MD simulations reduce the original volume of information by several orders of magnitude, thereby alleviating the negative impact of I/O bottlenecks. This work focuses on characterizing the impact of in situ and in transit analytics on the overall MD workflow performance, and the capability for capturing rapid, rare events in the simulated molecular system. The MD simulation and analysis processes share data via remote direct memory access (RDMA) using DataSpaces. Our metrics of interest are time spent waiting in I/O by the MD simulation, lost frames of the MD simulation, and idle time of the analysis. We measure these metrics for a diverse set of molecular systems and characterize their trends for in situ and in transit configurations. We then model which frames are dropped and which ones are analyzed for a real use case. The insights gained from this study are generally applicable for in situ and in transit workflows that require optimization of parameters to minimize loss in workflow performance and analytic accuracy.
Michela Taufer, Ewa Deelman, Michael R. Wyatt II, Tu Mai Anh Do, Loïc Pottier, Rafael Ferreira da Silva, Harel Weinstein, Michel A. Cuendet, Trilce Estrada
eScience2
2019 Custom Execution Environments with Containers in Pegasus-Enabled Scientific Workflows
abstract
Science reproducibility is a cornerstone feature in scientific workflows. In most cases, this has been implemented as a way to exactly reproduce the computational steps taken to reach the final results. While these steps are often completely described, including the input parameters, datasets, and codes, the environment in which these steps are executed is only described at a higher level with endpoints and operating system name and versions. Though this may be sufficient for reproducibility in the short term, systems evolve and are replaced over time, breaking the underlying workflow reproducibility. A natural solution to this problem is containers, as they are well defined, have a lifetime independent of the underlying system, and can be user-controlled so that they can provide custom environments if needed. This paper highlights some unique challenges that may arise when using containers in distributed scientific workflows. Further, this paper explores how the Pegasus Workflow Management System implements container support to address such challenges.
Karan Vahi, Michael Zink, Mats Rynge, George Papadimitriou 0002, Duncan A. Brown, Rajiv Mayani, Rafael Ferreira da Silva, Ewa Deelman, Anirban Mandal, Eric Lyons 0001
eScience8
2019 Collaborative circuit designs using the CRAFT repository
Adam Brinckman, Ewa Deelman, Sandeep Gupta 0001, Jarek Nabrzyski, Soowang Park, Rafael Ferreira da Silva, Ian J. Taylor, Karan Vahi
Future Gener. Comput. Syst.2
2019 Measuring the impact of burst buffers on data-intensive scientific workflows
Rafael Ferreira da Silva, Scott Callaghan, Tu Mai Anh Do, George Papadimitriou 0002, Ewa Deelman
Future Gener. Comput. Syst.5
2019 Using simple PID-inspired controllers for online resilient resource management of distributed scientific workflows
Rafael Ferreira da Silva, Rosa Filgueira, Ewa Deelman, Erola Pairo-Castineira, Ian Michael Overton, Malcolm P. Atkinson 0001
Future Gener. Comput. Syst.3
2018 A Job Sizing Strategy for High-Throughput Scientific Workflows
abstract
The user of a computing facility must make a critical decision when submitting jobs for execution: how many resources (such as cores, memory, and disk) should be requested for each job? If the request is too small, the job may fail due to resource exhaustion; if the request is too large, the job may succeed, but resources will be wasted. This decision is especially important when running hundreds of thousands of jobs in a high throughput workflow, which may exhibit complex, long tailed distributions of resource consumption. In this paper, we present a strategy for solving the job sizing problem: (1) applications are monitored and measured in user-space as they run; (2) the resource usage is collected into an online archive; and (3) jobs are automatically sized according to historical data in order to maximize throughput or minimize waste. We evaluate the solution analytically, and present case studies of applying the technique to high throughput physics and bioinformatics workflows consisting of hundreds of thousands of jobs, demonstrating an increase in throughput of 10-400 percent compared to naive approaches.
Benjamín Tovar, Rafael Ferreira da Silva, Gideon Juve, Ewa Deelman, William E. Allcock, Douglas Thain, Miron Livny
IEEE Trans. Parallel Distributed Syst.4
2017 Toward Prioritization of Data Flows for Scientific Workflows Using Virtual Software Defined Exchanges
abstract
Recent advances in cloud systems, on-demand circuits and software-defined networking have created new opportunities to enable complex, data-intensive scientific applications to run on dynamic networked cloud infrastructures. In this work, we present an end-to-end framework for autonomic adaptation for scientific workflows on networked cloud systems, which leverages novel network provisioning technologies. We present an application-independent controller framework called Mobius++ that includes dynamic network adaptation capabilities using Software-Defined Networking (SDN) mechanisms, which enables workflow management systems to address competing priorities of workflow operations, data movements in particular. We use a representative, data-intensive bioinformatics workflow as a driving use case to showcase the above capabilities. Experimental results show that the Mobius++ framework, in conjunction with a novel virtual Software Defined Exchange (SDX) platform, is able to dynamically prioritize bandwidths between different end-points, on-demand, and being driven by priority directives from a workflow management system. We show that data transfer jobs from two workflows with different priorities are accurately arbitrated as the relative priorities change.
Anirban Mandal, Paul Ruth, Ilya Baldin, Rafael Ferreira da Silva, Ewa Deelman
eScience5
2017 Reproducibility of execution environments in computational science using Semantics and Clouds
Idafen Santana-Pérez, Rafael Ferreira da Silva, Mats Rynge, Ewa Deelman, María S. Pérez 0001, Óscar Corcho
Future Gener. Comput. Syst.4
2017 A characterization of workflow management systems for extreme-scale applications
Rafael Ferreira da Silva, Rosa Filgueira, Ilia Pietri, Ming Jiang 0005, Rizos Sakellariou, Ewa Deelman
Future Gener. Comput. Syst.6
2016 Automating environmental computing applications with scientific workflows
abstract
Computational environmental science applications have evolved and become more complex over the last decade. In order to cope with the needs of such applications, computational methods and technologies have emerged to support the execution of these applications on heterogeneous, distributed systems. Among them are workflow management systems such as Pegasus. Pegasus is being used by researchers to model seismic wave propagation, to discover new celestial objects, to study RNA critical to human brain development, and to investigate other important research questions. This paper provides an introduction to scientific workflows and describes Pegasus and its main features. The paper highlights how the environmental science community has used Pegasus to automate their scientific workflow executions on high performance and high throughput computing systems by presenting three use cases: two Earth science workflows, and a climate science workflow.
Rafael Ferreira da Silva, Ewa Deelman, Rosa Filgueira, Karan Vahi, Mats Rynge, Rajiv Mayani, Benjamin Mayer
eScience2
2016 Science automation in practice: Performance data farming in workflows
abstract
This paper describes an approach to conduct large-scale parameter studies, where each data point in the study requires the execution of a whole scientific workflow. We show how a parameter studies system can be integrated with a workflow management system to seamlessly execute a large number of workflows, each with different input parameter values using large-scale computing infrastructure. The work is motivated by a need to collect performance-related data to conduct a sensitivity analysis in the context of relation between workflow input parameters and the performance of tasks in the workflow developed for the Spallation Neutron Source facility at the Oak Ridge National Laboratory.
Dariusz Król 0002, Jacek Kitowski, Rafael Ferreira da Silva, Gideon Juve, Karan Vahi, Mats Rynge, Ewa Deelman
ETFA7
2016 Consecutive Job Submission Behavior at Mira Supercomputer
abstract
Understanding user behavior is crucial for the evaluation of scheduling and allocation performances in HPC environments. This paper aims to further understand the dynamic user reaction to different levels of system performance by performing a comprehensive analysis of user behavior in recorded data in the form of delays in the subsequent job submission behavior. Therefore, we characterize a workload trace covering one year of job submissions from the Mira supercomputer at ALCF (Argonne Leadership Computing Facility). We perform an in-depth analysis of correlations between job characteristics, system performance metrics, and the subsequent user behavior. Analysis results show that the user behavior is significantly influenced by long waiting times, and that complex jobs (number of nodes and CPU hours) lead to longer delays in subsequent job submissions.
Stephan Schlagkamp, Rafael Ferreira da Silva, William E. Allcock, Ewa Deelman, Uwe Schwiegelshohn
HPDC4
2016 Storage-aware Algorithms for Scheduling of Workflow Ensembles in Clouds
abstract
This paper focuses on data-intensive workflows and addresses the problem of scheduling workflow ensembles under cost and deadline constraints in Infrastructure as a Service (IaaS) clouds. Previous research in this area ignores file transfers between workflow tasks, which, as we show, often have a large impact on workflow ensemble execution. In this paper we propose and implement a simulation model for handling file transfers between tasks, featuring the ability to dynamically calculate bandwidth and supporting a configurable number of replicas, thus allowing us to simulate various levels of congestion. The resulting model is capable of representing a wide range of storage systems available on clouds: from in-memory caches (such as memcached), to distributed file systems (such as NFS servers) and cloud storage (such as Amazon S3 or Google Cloud Storage). We observe that file transfers may have a significant impact on ensemble execution; for some applications up to 90 % of the execution time is spent on file transfers. Next, we propose and evaluate a novel scheduling algorithm that minimizes the number of transfers by taking advantage of data caching and file locality. We find that for data-intensive applications it performs better than other scheduling algorithms. Additionally, we modify the original scheduling algorithms to effectively operate in environments where file transfers take non-zero time.
Piotr Bryk, Maciej Malawski, Gideon Juve, Ewa Deelman
J. Grid Comput.4
2016 Dynamic and Fault-Tolerant Clustering for Scientific Workflows
abstract
Task clustering has proven to be an effective method to reduce execution overhead and to improve the computational granularity of scientific workflow tasks executing on distributed resources. However, a job composed of multiple tasks may have a higher risk of suffering from failures than a single task job. In this paper, we conduct a theoretical analysis of the impact of transient failures on the runtime performance of scientific workflow executions. We propose a general task failure modeling framework that uses a maximum likelihood estimation-based parameter estimation process to model workflow performance. We further propose three fault-tolerant clustering strategies to improve the runtime performance of workflow executions in faulty execution environments. Experimental results show that failures can have significant impact on executions where task clustering policies are not fault-tolerant, and that our solutions yield makespan improvements in such scenarios. In addition, we propose a dynamic task clustering strategy to optimize the workflow's makespan by dynamically adjusting the clustering granularity when failures arise. A trace-based simulation of five real workflows shows that our dynamic method is able to adapt to unexpected behaviors, and yields better makespans when compared to static methods.
Weiwei Chen 0002, Rafael Ferreira da Silva, Ewa Deelman, Thomas Fahringer
IEEE Trans. Cloud Comput.3
2015 Practical Resource Monitoring for Robust High Throughput Computing
abstract
Robust high throughput computing requires effective monitoring and enforcement of a variety of resources including CPU cores, memory, disk, and network traffic. Without effective monitoring and enforcement, it is easy to overload machines, causing failures and slowdowns, or underutilize machines, which results in wasted opportunities. This paper explores how to describe, measure, and enforce resources used by computational tasks. We focus on tasks running in distributed execution systems, in which a task requests the resources it needs, and the execution system ensures the availability of such resources. This presents two non-trivial problems: how to measure the resources consumed by a task, and how to monitor and report resource exhaustion in a robust and timely manner. For both of these tasks, operating systems have a variety of mechanisms with different degrees of availability, accuracy, overhead, and intrusiveness. We describe various forms of monitoring and the available mechanisms in contemporary operating systems. We then present two specific monitoring tools that choose different tradeoffs in overhead and accuracy, and evaluate them on a selection of benchmarks.
Gideon Juve, Benjamín Tovar, Rafael Ferreira da Silva, Dariusz Król 0002, Douglas Thain, Ewa Deelman, William E. Allcock, Miron Livny
CLUSTER6
2015 High Impact Computing: Computing for Science and the Science of Computing
abstract
Modern science often requires the processing and analysis of vast amounts of data in search of postulated phenomena, and the validation of core principles through the simulation of complex system behaviors and interactions. This is the case in fields such as astronomy, bioinformatics, physics, and climate and ocean modeling, and others. In order to support the computational and data needs of today's science, new knowledge must be gained on how to deliver the growing high-performance and distributed computing resources to the scientist's desktop in an accessible, reliable and scalable way. In over a decade of working with domain scientists, the Pegasus project [1,2] has developed tools and techniques that automate the computational processes used in data- and compute-intensive research. Among them is the scientific workflow management system, Pegasus, which is being used by researchers to model seismic wave propagation, to discover new celestial objects, to study RNA critical to human brain development, and to investigate other important research questions.
Ewa Deelman
HPDC1
2015 HUBzero and Pegasus: integrating scientific workflows into science gateways
abstract
Summary In this paper, we described the benefits and the challenges of integrating existing scientific workflow technologies into science gateways. Scientific workflow managers are powerful tools for handling large computational tasks. Domain scientists find it difficult to create new workflows, so many tasks that could benefit from workflow automation are often avoided or performed by hand. Two technologies have come together to bring the benefits of workflow to the masses. The Pegasus Workflow Management System can manage workflows comprised of millions of tasks, all the while recording data about the execution and intermediate results so that the provenance of the final result is clear. The HUBzero platform for scientific collaboration provides a venue for building and delivering tools to researchers and educators. With the press of a button, these tools can launch Pegasus workflows on national computing infrastructures and bring results back for plotting and visualization. As a result, the combination of Pegasus and HUBzero is bringing high‐throughput computing to a much wider audience. Copyright © 2014 John Wiley & Sons, Ltd.
Michael McLennan, Steven M. Clark, Ewa Deelman, Mats Rynge, Karan Vahi, Frank McKenna, Derrick Kearney, Carol X. Song
Concurr. Comput. Pract. Exp.3
2015 Using imbalance metrics to optimize task clustering in scientific workflow executions
Weiwei Chen 0002, Rafael Ferreira da Silva, Ewa Deelman, Rizos Sakellariou
Future Gener. Comput. Syst.3
2015 Pegasus, a workflow management system for science automation
Ewa Deelman, Karan Vahi, Gideon Juve, Mats Rynge, Scott Callaghan, Philip Maechling, Rajiv Mayani, Weiwei Chen 0002, Rafael Ferreira da Silva, Miron Livny, R. Kent Wenger
Future Gener. Comput. Syst.1
2015 Algorithms for cost- and deadline-constrained provisioning for scientific workflow ensembles in IaaS clouds
Maciej Malawski, Gideon Juve, Ewa Deelman, Jarek Nabrzyski
Future Gener. Comput. Syst.3
2014 Community Resources for Enabling Research in Distributed Scientific Workflows
abstract
A significant amount of recent research in scientific workflows aims to develop new techniques, algorithms and systems that can overcome the challenges of efficient and robust execution of ever larger workflows on increasingly complex distributed infrastructures. Since the infrastructures, systems and applications are complex, and their behavior is difficult to reproduce using physical experiments, much of this research is based on simulation. However, there exists a shortage of realistic datasets and tools that can be used for such studies. In this paper we describe a collection of tools and data that have enabled research in new techniques, algorithms, and systems for scientific workflows. These resources include: 1) execution traces of real workflow applications from which workflow and system characteristics such as resource usage and failure profiles can be extracted, 2) a synthetic workflow generator that can produce realistic synthetic workflows based on profiles extracted from execution traces, and 3) a simulator framework that can simulate the execution of synthetic workflows on realistic distributed infrastructures. This paper describes how we have used these resources to investigate new techniques for efficient and robust workflow execution, as well as to provide improvements to the Pegasus Workflow Management System or other workflow tools. Our goal in describing these resources is to share them with other researchers in the workflow research community. All of the tools and data are freely available online for the community at http://www.workflowarchive.org. These data have already been leveraged for a number of studies.
Rafael Ferreira da Silva, Weiwei Chen 0002, Gideon Juve, Karan Vahi, Ewa Deelman
eScience5
2013 Rethinking data management for big data scientific workflows
abstract
Scientific workflows consist of tasks that operate on input data to generate new data products that are used by subsequent tasks. Workflow management systems typically stage data to computational sites before invoking the necessary computations. In some cases data may be accessed using remote I/O. There are limitations with these approaches, however. First, the storage at a computational site may be limited and not able to accommodate the necessary input and intermediate data. Second, even if there is enough storage, it is sometimes managed by a filesystem with limited scalability. In recent years, object stores have been shown to provide a scalable way to store and access large datasets, however, they provide a limited set of operations (retrieve, store and delete) that do not always match the requirements of the workflow tasks. In this paper, we show how scientific workflows can take advantage of the capabilities of object stores without requiring users to modify their workflow-based applications or scientific codes. We present two general approaches, one that exclusively uses object stores to store all the files accessed and generated by a workflow, while the other relies on the shared filesystem for caching intermediate data sets. We have implemented both of these approaches in the Pegasus Workflow Management System and have used them to execute workflows in variety of execution environments ranging from traditional supercomputing environments that have a shared filesystem to dynamic environments like Amazon AWS and the Open Science Grid that only offer remote object stores. As a result, Pegasus users can easily migrate their applications from a shared filesystem deployment to one using object stores without changing their application codes.
Karan Vahi, Mats Rynge, Gideon Juve, Rajiv Mayani, Ewa Deelman
IEEE BigData5
2013 Introducing PRECIP: An API for Managing Repeatable Experiments in the Cloud
abstract
Cloud computing with its on-demand access to resources has emerged as a tool used by researchers from a wide range of domains to run computer-based experiments. In this paper we introduce a flexible experiment management API, written in Python that simplifies and formalizes the execution of scientific experiments on cloud infrastructures. We describe the features and functionality of PRECIP (Pegasus Repeatable Experiments for the Cloud in Python), and how PRECIP can be used to set up experiments on academic clouds such as OpenStack Eucalyptus, Nimbus, and commercial clouds such as Amazon EC2.
Sepideh Azarnoosh, Mats Rynge, Gideon Juve, Ewa Deelman, Michal Niec, Maciej Malawski, Rafael Ferreira da Silva
CloudCom (2)4
2013 Balanced Task Clustering in Scientific Workflows
abstract
Scientific workflows can be composed of many fine computational granularity tasks. The runtime of these tasks may be shorter than the duration of system overheads, for example, when using multiple resources of a cloud infrastructure. Task clustering is a runtime optimization technique that merges multiple short tasks into a single job such that the scheduling overhead is reduced and the overall runtime performance is improved. However, existing task clustering strategies only provide a coarse-grained approach that relies on an over-simplified workflow model. In our work, we examine the reasons that cause Runtime Imbalance and Dependency Imbalance in task clustering. Next, we propose quantitative metrics to evaluate the severity of the two imbalance problems respectively. Furthermore, we propose a series of task balancing methods to address these imbalance problems. Finally, we analyze their relationship with the performance of these task balancing methods. A trace-based simulation shows our methods can significantly improve the runtime performance of two widely used workflows compared to the actual implementation of task clustering.
Weiwei Chen 0002, Rafael Ferreira da Silva, Ewa Deelman, Rizos Sakellariou
e-Science3
2013 Imbalance optimization in scientific workflows
abstract
Scientific workflows are a means of defining and orchestrating large, complex, multi-stage computations that perform data analysis and/or simulation. Task clustering is a runtime optimization technique that merges multiple short workflow tasks into a single job such that the job execution overhead is reduced and the overall runtime performance of the workflow is significantly improved. However, current task clustering strategies fail to consider the imbalance problem of both task runtime and task dependency. In our work, we first investigate the different causes of runtime imbalance and dependency imbalance. We then introduce a series of metrics based on our prior work to measure the severity of runtime and dependency imbalance respectively. Finally, we study a wide range of real scientific workflows to generalize the relationship between these metrics and balancing methods.
Weiwei Chen 0002, Ewa Deelman, Rizos Sakellariou
ICS2
2013 Characterizing and profiling scientific workflows
Gideon Juve, Ann L. Chervenak, Ewa Deelman, Shishir Bharathi, Gaurang Mehta, Karan Vahi
Future Gener. Comput. Syst.3
2013 A Case Study into Using Common Real-Time Workflow Monitoring Infrastructure for Scientific Workflows
Karan Vahi, Ian Harvey, Taghrid Samak, Dan Gunter, Kieran Evans, David Rogers, Ian J. Taylor, Monte Goode, Fabio Silva, Eddie Al-Shakarchi, Gaurang Mehta, Ewa Deelman, Andrew C. Jones
J. Grid Comput.12
2012 Integration of Workflow Partitioning and Resource Provisioning
abstract
The recent increased use of workflow management systems by large scientific collaborations presents the challenge of scheduling large-scale workflows onto distributed resources. This work aims to partition large-scale scientific workflows in conjunction with resources provisioning to reduce the workflow make span and resource cost.
Weiwei Chen 0002, Ewa Deelman
CCGRID2
2012 Using Clouds for Science, is it just Kicking the Can down the Road?
Ewa Deelman, Gideon Juve, G. Bruce Berriman
CLOSER1
2012 Failure analysis of distributed scientific workflows executing in the cloud
Taghrid Samak, Dan Gunter, Monte Goode, Ewa Deelman, Gideon Juve, Fabio Silva, Karan Vahi
CNSM4
2012 WorkflowSim: A toolkit for simulating scientific workflows in distributed environments
abstract
Simulation is one of the most popular evaluation methods in scientific workflow studies. However, existing workflow simulators fail to provide a framework that takes into consideration heterogeneous system overheads and failures. They also lack the support for widely used workflow optimization techniques such as task clustering. In this paper, we introduce WorkflowSim, which extends the existing CloudSim simulator by providing a higher layer of workflow management. We also indicate that to ignore system overheads and failures in simulating scientific workflows could cause significant inaccuracies in the predicted workflow runtime. To further validate its value in promoting other research work, we introduce two promising research areas for which WorkflowSim provides a unique and effective evaluation platform.
Weiwei Chen 0002, Ewa Deelman
eScience2
2012 Cost- and deadline-constrained provisioning for scientific workflow ensembles in IaaS clouds
abstract
Large-scale applications expressed as scientific workflows are often grouped into ensembles of inter-related workflows. In this paper, we address a new and important problem concerning the efficient management of such ensembles under budget and deadline constraints on Infrastructure- as-aService (IaaS) clouds. We discuss, develop, and assess algorithms based on static and dynamic strategies for both task scheduling and resource provisioning. We perform the evaluation via simulation using a set of scientific workflow ensembles with a broad range of budget and deadline parameters, taking into account uncertainties in task runtime estimations, provisioning delays, and failures. We find that the key factor determining the performance of an algorithm is its ability to decide which workflows in an ensemble to admit or reject for execution. Our results show that an admission procedure based on workflow structure and estimates of task runtimes can significantly improve the quality of solutions.
Maciej Malawski, Gideon Juve, Ewa Deelman, Jarek Nabrzyski
SC3
2012 Fault Tolerant Clustering in Scientific Workflows
abstract
Task clustering has been proven to be an effective method to reduce execution overhead and increase the computational granularity of workflow tasks executing on distributed resources. However, a job composed of multiple tasks may have a greater risk of suffering from failures than a job composed of a single task. Our theoretic analysis and simulation results demonstrate that failures can have a significant impact on the runtime performance of workflows that use existing clustering policies that ignore failures. We therefore propose two general failure modeling frameworks (task failure model and job failure model) to address these performance issues. We show the necessity to consider the fault tolerance in the task failure model. Based on the task failure model, we propose three methods to improve the workflow performance in dynamic environments. A simulation-based evaluation is performed and it shows that our approach can improve the workflow makespan significantly for two important applications.
Weiwei Chen 0002, Ewa Deelman
SERVICES2
2012 An Evaluation of the Cost and Performance of Scientific Workflows on Amazon EC2
Gideon Juve, Ewa Deelman, G. Bruce Berriman, Benjamin P. Berman, Philip Maechling
J. Grid Comput.2
2011 Automating Application Deployment in Infrastructure Clouds
abstract
Cloud computing systems are becoming an important platform for distributed applications in science and engineering. Infrastructure as a Service (IaaS) clouds provide the capability to provision virtual machines (VMs) on demand with a specific configuration of hardware resources, but they do not provide functionality for managing resources once they are provisioned. In order for such clouds to be used effectively, tools need to be developed that can help users to deploy their applications in the cloud. In this paper we describe a system we have developed to provision, configure, and manage virtual machine deployments in the cloud. We also describe our experiences using the system to provision resources for scientific workflow applications, and identify areas for further research.
Gideon Juve, Ewa Deelman
CloudCom2
2011 Online workflow management and performance analysis with Stampede
Dan Gunter, Ewa Deelman, Taghrid Samak, Christopher X. Brooks, Monte Goode, Gideon Juve, Gaurang Mehta, Priscilla Moraes, Fabio Silva, D. Martin Swany, Karan Vahi
CNSM2
2011 A Cloud-based Dynamic Workflow for Mass Spectrometry Data Analysis
abstract
There is a growing interest in the use of cloud computing for scientific applications, including scientific workflows. Key attractions of the cloud include the pay-as-you-go model and elasticity. While the elasticity offered by clouds can be beneficial for many applications and use-scenarios, it also imposes significant challenges in the development of applications or services. For example, no general framework exists that can enable a scientific workflow to execute in a dynamic fashion, i.e. exploiting elasticity of clouds and automatically allocating and deal locating resources to meet time and/or cost constraints. This paper presents a case-study in creating a dynamic cloud workflow implementation of a scientific application. We work with Mass Matrix, an application which searches proteins and peptides from tandem mass spectrometry data. In order to use cloud resources, we first parallelize the search method used in this algorithm. Next, we create a flexible workflow using the Pegasus Workflow Management System. Finally, we add a new dynamic resource allocation module, which can use fewer or a larger number of resources based on a time constraint specified by the user. We evaluate our implementation using several different datasets, and show that the application scales quite well, and that our dynamic framework is effective in meeting time constraints.
Ashish Nagavaram, Gagan Agrawal, Michael A. Freitas, Kelly H. Telu, Gaurang Mehta, Rajiv Mayani, Ewa Deelman
eScience7
2011 Experiences Using GlideinWMS and the Corral Frontend across Cyberinfrastructures
abstract
Even with Grid technologies, the main mode of access for the current High Performance Computing and High Throughput Computing infrastructures today is logging in via ssh. This mode of access locks scientists to particular machines as it is difficult to move the codes and environments between hosts. In this paper we show how switching the resource access mode to a Condor glide in-based overlay can bring together computational resources from multiple cyber infrastructures. This approach provides scientists with a computational infrastructure anchored around the familiar environment of the desktop computer. Additionally, the approach enhances the reliability of applications and workflows by automatically rerouting jobs to functioning infrastructures. Two different science applications were used to demonstrate applicability, one from the field of astronomy and the other one from earth sciences. We demonstrate that a desktop computer is viable as a submit host and central manager for these kind of glide in overlays. However, issues of ease of use and security need to be considered.
Mats Rynge, Gideon Juve, Gaurang Mehta, Ewa Deelman, Krista Larson, Burt Holzman, Igor Sfiligoi, Frank Würthwein, G. Bruce Berriman, Scott Callaghan
eScience4
2011 Online Fault and Anomaly Detection for Large-Scale Scientific Workflows
abstract
Scientific workflows are an enabler of complex scientific analyses. Large-scale scientific workflows are executed on complex parallel and distributed resources, where many things can fail. Application scientists need to track the status of their workflows in real time, detect execution anomalies automatically, and perform troubleshooting -- without logging into remote nodes or searching through thousands of log files. As part of the NSF-funded Synthesized Tools for Archiving Monitoring Performance and Enhanced DEbugging (STAMPEDE) project, we have developed an infrastructure to answer these needs by integrating detailed workflow and resource monitoring. On top of this infrastructure, we have developed analysis techniques for online detection of a wide variety of "hard" and "soft" types of failures. We use these detected failures to derive higher-level statistics about the status of the resources and the workflow as a whole. In this paper, we describe our techniques and evaluate their effectiveness in the context of real application logs.
Taghrid Samak, Dan Gunter, Monte Goode, Ewa Deelman, Gideon Juve, Gaurang Mehta, Fabio Silva, Karan Vahi
HPCC4
2011 Wrangler: virtual cluster provisioning for the cloud
abstract
Cloud computing systems are becoming an important platform for science applications. Infrastructure as a Service (IaaS) clouds provide the capability to provision virtual machines (VMs) on demand with a specific configuration of hardware resources, but they do not provide functionality for managing those resources once provisioned. In order for such clouds to be used effectively for parallel and distributed scientific applications, tools need to be developed that can help users to deploy their applications in the cloud. This paper describes a system we have developed to provision, configure, and manage clusters of virtual machines.
Gideon Juve, Ewa Deelman
HPDC2
2011 RseqFlow: workflows for RNA-Seq data analysis
abstract
SUMMARY: We have developed an RNA-Seq analysis workflow for single-ended Illumina reads, termed RseqFlow. This workflow includes a set of analytic functions, such as quality control for sequencing data, signal tracks of mapped reads, calculation of expression levels, identification of differentially expressed genes and coding SNPs calling. This workflow is formalized and managed by the Pegasus Workflow Management System, which maps the analysis modules onto available computational resources, automatically executes the steps in the appropriate order and supervises the whole running process. RseqFlow is available as a Virtual Machine with all the necessary software, which eliminates any complex configuration and installation steps. AVAILABILITY AND IMPLEMENTATION: http://genomics.isi.edu/rnaseq CONTACT: [email protected]; [email protected]; [email protected]; [email protected] SUPPLEMENTARY INFORMATION: Supplementary data are available at Bioinformatics online.
Gaurang Mehta, Rajiv Mayani, Jingxi Lu, Tade Souaiaia, Yangho Chen, Andrew P. Clark, Hee Jae Yoon, Oleg V. Evgrafov, James A. Knowles, Ewa Deelman
Bioinform.12
2011 The interplay of resource provisioning and workflow optimization in scientific applications
abstract
Abstract In this paper, we propose resource provisioning as the means to reduce the completion time of scientific workflows in a Grid environment. We propose task clustering as a form of workflow optimization that can be used along with provisioning in order to achieve this reduction in completion time. Provisioning can be done statically using advance reservations (ARs) or using dynamic provisioning mechanisms. A simulation is done using the Maui simulator, a workload trace collected from the NCSA Teragrid cluster and 13 workflows to study the effect of provisioning on the completion time of the scientific workflows. The results show in general a reduction of about 50% in the workflow completion time using provisioning for the First In First Out and fair share scheduling policies. In this paper, we also examine the cost of resource provisioning and propose a utilization‐based metric that can be used to guide the provisioning decisions in order to reduce the cost. Finally, we present the results of a survey on the support of AR at the Grid sites. Copyright © 2011 John Wiley & Sons, Ltd.
Ewa Deelman
Concurr. Comput. Pract. Exp.2
2011 BTS: Resource capacity estimate for time-targeted science workflows
Eun-Kyu Byun, Yang-Suk Kee, Jin-Soo Kim 0001, Ewa Deelman, Seung Ryoul Maeng
J. Parallel Distributed Comput.4
2010 Bridging the Gap between Business and Scientific Workflows: Humans in the Loop of Scientific Workflows
abstract
Due to their different target applications business and scientific workflow systems provide different sets of features to their users. Significant amount of research is currently being done to employ the business workflow technology in the scientific domain. This usually means extending the workflow language and thus the modeling tool and execution engine. In this paper we aim to bring business and scientific workflows together in order to exploit the advantages of both. We explore the interplay between business and scientific workflows in the context of human interactions with the management of workflow execution. We present an approach and implementation based on BPEL and Pegasus and show that the approach can be beneficial to scientists.
Mirko Sonntag, Dimka Karastoyanova, Ewa Deelman
eScience3
2010 BPEL4Pegasus: Combining Business and Scientific Workflows
Mirko Sonntag, Dimka Karastoyanova, Ewa Deelman
ICSOC3
2010 Data Sharing Options for Scientific Workflows on Amazon EC2
abstract
Efficient data management is a key component in achieving good performance for scientific workflows in distributed environments. Workflow applications typically communicate data between tasks using files. When tasks are distributed, these files are either transferred from one computational node to another, or accessed through a shared storage system. In grids and clusters, workflow data is often stored on network and parallel file systems. In this paper we investigate some of the ways in which data can be managed for workflows in the cloud. We ran experiments using three typical workflow applications on Amazon's EC2. We discuss the various storage and file systems we used, describe the issues and problems we encountered deploying them on EC2, and analyze the resulting performance and cost of the workflows.
Gideon Juve, Ewa Deelman, Karan Vahi, Gaurang Mehta, G. Bruce Berriman, Benjamin P. Berman, Philip Maechling
SC2
2010 Scaling up workflow-based applications
Scott Callaghan, Ewa Deelman, Dan Gunter, Gideon Juve, Philip Maechling, Christopher X. Brooks, Karan Vahi, Kevin Milner 0001, Robert Graves, Edward Field, David Okaya, Thomas H. Jordan
J. Comput. Syst. Sci.2
2009 An integrated framework for performance-based optimization of scientific workflows
abstract
Data analysis processes in scientific applications can be expressed as coarse-grain workflows of complex data processing operations with data flow dependencies between them. Performance optimization of these workflows can be viewed as a search for a set of optimal values in a multi-dimensional parameter space. While some performance parameters such as grouping of workflow components and their mapping to machines do not a ect the accuracy of the output, others may dictate trading the output quality of individual components (and of the whole workflow) for performance. This paper describes an integrated framework which is capable of supporting performance optimizations along multiple dimensions of the parameter space. Using two real-world applications in the spatial data analysis domain, we present an experimental evaluation of the proposed framework.
Vijay S. Kumar, P. Sadayappan, Gaurang Mehta, Karan Vahi, Ewa Deelman, Varun Ratnakar, Jihie Kim, Yolanda Gil, Mary W. Hall, Tahsin M. Kurç, Joel H. Saltz
HPDC5
2009 Adaptive workflow processing and execution in Pegasus
abstract
Abstract Workflows are widely used in applications that require coordinated use of computational resources. Workflow definition languages typically abstract over some aspects of the way in which a workflow is to be executed, such as the level of parallelism to be used or the physical resources to be deployed. As a result, a workflow management system has the responsibility of establishing how best to execute a workflow given the available resources. The Pegasus workflow management system compiles abstract workflows into concrete execution plans, and has been widely used in large‐scale e‐Science applications. This paper describes an extension to Pegasus whereby resource allocation decisions are revised during workflow evaluation, in the light of feedback on the performance of jobs at runtime. The contributions of this paper include: (i) a description of how adaptive processing has been retrofitted to an existing workflow management system; (ii) a scheduling algorithm that allocates resources based on runtime performance; and (iii) an experimental evaluation of the resulting infrastructure using grid middleware over clusters. Copyright © 2009 John Wiley & Sons, Ltd.
Kevin Lee 0006, Norman W. Paton, Rizos Sakellariou, Ewa Deelman, Alvaro A. A. Fernandes, Gaurang Mehta
Concurr. Comput. Pract. Exp.4
2009 Workflows and e-Science: An overview of workflow system features and capabilities
Ewa Deelman, Dennis Gannon, Matthew S. Shields, Ian J. Taylor
Future Gener. Comput. Syst.1
2008 Data Management Challenges of Data-Intensive Scientific Workflows
abstract
Scientific workflows play an important role in today's science. Many disciplines rely on workflow technologies to orchestrate the execution of thousands of computational tasks. Much research to-date focuses on efficient, scalable, and robust workflow execution, especially in distributed environments. However, many challenges remain in the area of data management related to workflow creation, execution, and result management. In this paper we examine some of these issues in the context of the entire workflow lifecycle.
Ewa Deelman, Ann L. Chervenak
CCGRID1
2008 Estimating Resource Needs for Time-Constrained Workflows
abstract
Workflow technologies have become a major vehicle for the easy and efficient development of science applications. At the same time new computing environments such as the Cloud are now available. A challenge is to determine the right amount of resources to provision for an application. This paper introduces an algorithm named balanced time scheduling (BTS), which estimates the minimum number of virtual processors required to execute a workflow within a user-specified finish time. The resource estimate of BTS is abstract, so it can be easily integrated with any resource description language or any resource provisioning system. The experimental results with a number of synthetic workflows demonstrate that BTS can estimate the computing capacity close to the optimal. The algorithm is scalable so that its turnaround time is only tens of seconds even with workflows having thousands of tasks and edges.
Eun-Kyu Byun, Yang-Suk Kee, Ewa Deelman, Karan Vahi, Gaurang Mehta, Jin-Soo Kim 0001
eScience3
2008 Reducing Time-to-Solution Using Distributed High-Throughput Mega-Workflows - Experiences from SCEC CyberShake
abstract
Researchers at the Southern California Earthquake Center (SCEC) use large-scale grid-based scientific workflows to perform seismic hazard research as a part of SCEC's program of earthquake system science research. The scientific goal of the SCEC CyberShake project is to calculate probabilistic seismic hazard curves for sites in Southern California. For each site of interest, the CyberShake platform includes two large-scale MPI calculations and approximately 840,000 embarrassingly parallel post-processing jobs. In this paper, we describe the computational requirements of CyberShake and detail how we meet these requirements using grid-based, high-throughput, scientific workflow tools. We describe the specific challenges we encountered and we discuss workflow throughput optimizations we developed that reduced our time to solution by a factor of three and we present runtime statistics and propose further optimizations.
Scott Callaghan, Philip Maechling, Ewa Deelman, Karan Vahi, Gaurang Mehta, Gideon Juve, Kevin Milner 0001, Robert Graves, Edward Field, David Okaya, Dan Gunter, Keith Beattie, Thomas H. Jordan
eScience3
2008 On the Use of Cloud Computing for Scientific Workflows
abstract
This paper explores the use of cloud computing for scientific workflows, focusing on a widely used astronomy application-Montage. The approach is to evaluate from the point of view of a scientific workflow the tradeoffs between running in a local environment, if such is available, and running in a virtual environment via remote, wide-area network resource access. Our results show that for Montage, a workflow with short job runtimes, the virtual environment can provide good compute time performance but it can suffer from resource scheduling delays and widearea communications.
Christina Hoffa, Gaurang Mehta, Timothy Freeman 0001, Ewa Deelman, Kate Keahey, G. Bruce Berriman, John Good
eScience4
2008 Resource Provisioning Options for Large-Scale Scientific Workflows
abstract
Scientists in many fields are developing large-scale workflows containing millions of tasks and requiring thousands of hours of aggregate computation time. Acquiring the computational resources to execute these workflows poses many challenges for application developers. Although the grid provides ready access to large pools of computational resources, the traditional approach to accessing these resources suffers from many overheads that lead to poor performance. In this paper we examine several techniques based on resource provisioning that may be used to reduce these overheads. These techniques include: advance reservations, multi-level scheduling, and infrastructure as a service (IaaS). We explain the advantages and disadvantages of these techniques in terms of cost, performance and usability.
Gideon Juve, Ewa Deelman
eScience2
2008 Designing and parameterizing a workflow for optimization: A case study in biomedical imaging
abstract
This paper describes our experience to date employing the systematic mapping and optimization of large- scale scientific application workflows to current and future parallel platforms. The overall goal of the project is to integrate a set of system layers - application program, compiler, run-time environment, knowledge representation, optimization framework, and workflow manager - and through a systematic strategy for workflow mapping, our approach will exploit the vast machine resources available in such parallel platforms to dramatically increase the productivity of application programmers. In this paper, we describe the representation of a biomedical imaging application as a workflow, our early experiences in integrating the set of tools brought together for this project, and implications for future applications.
Vijay S. Kumar, Mary W. Hall, Jihie Kim, Yolanda Gil, Tahsin M. Kurç, Ewa Deelman, Varun Ratnakar, Joel H. Saltz
IPDPS6
2008 The cost of doing science on the cloud: the Montage example
abstract
Utility grids such as the Amazon EC2 cloud and Amazon S3 offer computational and storage resources that can be used on-demand for a fee by compute and data-intensive applications. The cost of running an application on such a cloud depends on the compute, storage and communication resources it will provision and consume. Different execution plans of the same application may result in significantly different costs. Using the Amazon cloud fee structure and a real-life astronomy application, we study via simulation the cost performance tradeoffs of different execution and resource provisioning plans. We also study these trade-offs in the context of the storage and communication fees of Amazon S3 when used for long-term application data archival. Our results show that by provisioning the right amount of storage and compute resources, cost can be significantly reduced with no significant impact on application performance.
Ewa Deelman, Miron Livny, G. Bruce Berriman, John Good
SC1
2008 Provenance trails in the Wings/Pegasus system
abstract
Abstract Our research focuses on creating and executing large‐scale scientific workflows that often involve thousands of computations over distributed, shared resources. We describe an approach to workflow creation and refinement that uses semantic representations to (1) describe complex scientific applications in a data‐independent manner, (2) automatically generate workflows of computations for given data sets, and (3) map the workflows to available computing resources for efficient execution. Our approach is implemented in the Wings/Pegasus workflow system and has been demonstrated in a variety of scientific application domains. This paper illustrates the application‐level provenance information generated Wings during workflow creation and the refinement provenance by the Pegasus mapping system for execution over grid computing environments. We show how this information is used in answering the queries of the First Provenance Challenge. Copyright © 2007 John Wiley & Sons, Ltd.
Jihie Kim, Ewa Deelman, Yolanda Gil, Gaurang Mehta, Varun Ratnakar
Concurr. Comput. Pract. Exp.2
2008 Special Issue: The First Provenance Challenge
abstract
Abstract The first Provenance Challenge was set up in order to provide a forum for the community to understand the capabilities of different provenance systems and the expressiveness of their provenance representations. To this end, a functional magnetic resonance imaging workflow was defined, which participants had to either simulate or run in order to produce some provenance representation, from which a set of identified queries had to be implemented and executed. Sixteen teams responded to the challenge, and submitted their inputs. In this paper, we present the challenge workflow and queries, and summarize the participants' contributions. Copyright © 2007 John Wiley & Sons, Ltd.
Luc Moreau 0001, Bertram Ludäscher, Ilkay Altintas, Roger S. Barga, Shawn Bowers, Steven P. Callahan, George Chin, Ben Clifford, Shirley Cohen, Sarah Cohen Boulakia, Susan B. Davidson, Ewa Deelman, Luciano A. Digiampietri, Ian T. Foster, Juliana Freire, James Frew, Joe Futrelle, Tara Gibson, Yolanda Gil, Carole A. Goble, Jennifer Golbeck, Paul Groth, David A. Holland, Jihie Kim, David Koop, Ales Krenek, Timothy M. McPhillips, Gaurang Mehta, Simon Miles, Dominic Metzger, Steve Munroe, James D. Myers, Beth Plale, Norbert Podhorszki, Varun Ratnakar, Emanuele Santos, Carlos Scheidegger, Karen Schuchardt, Margo I. Seltzer, Yogesh L. Simmhan, Cláudio T. Silva, Peter Slaughter, Eric G. Stephan, Robert Stevens 0001, Daniele Turi, Huy T. Vo, Michael Wilde, Jun Zhao 0003, Yong Zhao 0009
Concurr. Comput. Pract. Exp.12
2007 Wings for Pegasus: Creating Large-Scale Scientific Applications Using Semantic Representations of Computational Workflows
Yolanda Gil, Varun Ratnakar, Ewa Deelman, Gaurang Mehta, Jihie Kim
AAAI3
2007 Scheduling Data-IntensiveWorkflows onto Storage-Constrained Distributed Resources
abstract
In this paper we examine the issue of optimizing disk usage and of scheduling large-scale scientific workflows onto distributed resources where the workflows are data- intensive, requiring large amounts of data storage, and where the resources have limited storage resources. Our approach is two-fold: we minimize the amount of space a workflow requires during execution by removing data files at runtime when they are no longer required and we schedule the workflows in a way that assures that the amount of data required and generated by the workflow fits onto the individual resources. For a workflow used by gravitational- wave physicists, we were able to improve the amount of storage required by the workflow by up to 57 %. We also designed an algorithm that can not only find feasible solutions for workflow task assignment to resources in disk- space constrained environments, but can also improve the overall workflow performance.
Arun Ramakrishnan, Henan Zhao, Ewa Deelman, Rizos Sakellariou, Karan Vahi, Kent Blackburn, David Meyers, Michael Samidi
CCGRID4
2007 Connecting Scientific Data to Scientific Experiments with Provenance
abstract
As scientific workflows and the data they operate on, grow in size and complexity, the task of defining how those workflows should execute (which resources to use, where the resources must be in readiness for processing etc.) becomes proportionally more difficult. While "workflow compilers", such as Pegasus, reduce this burden, a further problem arises: since specifying details of execution is now automatic, a workflow's results are harder to interpret, as they are partly due to specifics of execution. By automating steps between the experiment design and its results, we lose the connection between them, hindering interpretation of results. To reconnect the scientific data with the original experiment, we argue that scientists should have access to the full provenance of their data, including not only parameters, inputs and intermediary data, but also the abstract experiment, refined into a concrete execution by the "workflow compiler". In this paper, we describe preliminary work on adapting Pegasus to capture the process of workflow refinement in the PASOA provenance system.
Simon Miles, Ewa Deelman, Paul Groth, Karan Vahi, Gaurang Mehta, Luc Moreau 0001
eScience2
2007 A provisioning model and its comparison with best-effort for performance-cost optimization in grids
abstract
The resource availability in Grids is generally unpredictable due to the autonomous and shared nature of the Grid resources and stochastic nature of the workload resulting in a best effort quality of service. The resource providers optimize for throughput and utilization whereas the users optimize for application performance. We present a cost-based model where the providers advertise resource availability to the user community. We also present a multi-objective genetic algorithm formulation for selecting the set of resources to be provisioned that optimizes the application performance while minimizing the resource costs. We use trace-based simulations to compare the application performance and cost using the provisioned and the best effort approach with a number of artificially generated workflow-structured applications and a seismic hazard application from the earthquake science community. The provisioned approach shows promising results when the resources are under high utilization and/or the applications have significant resource requirements.
Carl Kesselman, Ewa Deelman
HPDC3
2007 Intelligent Optimization of Parallel and Distributed Applications
abstract
This paper describes a new project that systematically addresses the enormous complexity of mapping applications to current and future parallel platforms. By integrating the system layers - domain-specific environment, application program, compiler, run-time environment, performance models and simulation, and workflow manager - and through a systematic strategy for application mapping, our approach exploit the vast machine resources available in such parallel platforms to dramatically increase the productivity of application programmers. This project brings together computer scientists in the areas represented by the system layers (i.e., language extensions, compilers, run-time systems, workflows) together with expertise in knowledge representation and machine learning. With expert domain scientists in molecular dynamics (MD) simulation, we are developing our approach in the context of a specific application class which already targets environments consisting of several hundreds of processors. In this way, we gain valuable insight into a generalizable strategy, while simultaneously producing performance benefits for existing and important applications.
Bhupesh Bansal, Ümit V. Çatalyürek, Jacqueline Chame, Chun Chen 0002, Ewa Deelman, Yolanda Gil, Mary W. Hall, Vijay S. Kumar, Tahsin M. Kurç, Kristina Lerman, Aiichiro Nakano, Yoon-Ju Lee Nelson, Joel H. Saltz, Ashish Sharma 0001, Priya Vashishta
IPDPS5
2006 Managing Large-Scale Workflow Execution from Resource Provisioning to Provenance Tracking: The CyberShake Example
abstract
This paper discusses the process of building an environment where large-scale, complex, scientific analysis can be scheduled onto a heterogeneous collection of computational and storage resources. The example application is the Southern California Earthquake Center (SCEC) CyberShake project, an analysis designed to compute probabilistic seismic hazard curves for sites in the Los Angeles area. We explain which software tools were used to build to the system, describe their functionality and interactions. We show the results of running the CyberShake analysis that included over 250,000 jobs using resources available through SCEC and the TeraGrid.
Ewa Deelman, Scott Callaghan, Edward Field, Hunter Francoeur, Robert Graves, Vipin Gupta, Thomas H. Jordan, Carl Kesselman, Philip Maechling, John Mehringer, Gaurang Mehta, David Okaya, Karan Vahi
e-Science1
2006 Managing Large-Scale Scientific Workflows in Distributed Environments: Experiences and Challenges
abstract
In this paper we discuss several challenges associated scientific workflow design and management in distributed, heterogeneous environments. Based on our prior work with a number of scientific applications, we describe the workflow lifecycle and examine our experiences and the challenges ahead as they pertain to the user experience, planning the workflow execution and managing the execution itself.
Ewa Deelman, Yolanda Gil
e-Science1
2006 Automating Climate Science: Large Ensemble Simulations on the TeraGrid with the GriPhyN Virtual Data System
abstract
Ensemble simulations are a promising technique for identifying the signal of atmospheric response to extra-tropical sea surface temperature variability with high statistical significance. The basic idea is to perform multiple simulations from slightly different initial conditions and then to study the average signal of the ensemble. A significant obstacle to performing such ensemble simulations is the bookkeeping required to prepare, execute, and track the progress of hundreds of different computations. We describe an ensemble simulation experiment in which the Fast Ocean Atmosphere Model was run on the U.S. TeraGrid. In this experiment, we used the GriPhyN Virtual Data System to manage our ensemble simulations and their execution on distributed resources, achieving dramatic (order-of-magnitude) reductions in turnaround time relative to previous manual experiments.
Veronika Nefedova, Robert L. Jacob, Ian T. Foster, Yun Liu 0032, Ewa Deelman, Gaurang Mehta, Mei-Hui Su, Karan Vahi
e-Science6
2006 Application-Level Resource Provisioning on the Grid
abstract
In this paper, we present algorithms for Grid resource provisioning that employ agreement-based resource management. These algorithms allow userlevel resource allocation and scheduling of applications that are structured as a precedenceconstrained set of tasks. We present a provisioning model where the resource availability in the Grid can be enumerated as a set of slots. A slot is defined as a number of processors available from a certain start time for a certain duration at a certain cost. Using a cost model that combines the cost of resource allocation and the expected application runtime, we evaluate the performance of the Min-Min and of the Genetic algorithm (GA)-based heuristics for a range of synthetic applications. We show that the GA paired with a list scheduling algorithm can obtain significantly better solutions than the Min-Min heuristic alone.
Carl Kesselman, Ewa Deelman
e-Science3
2006 What makes workflows work in an opportunistic environment?
abstract
Abstract In this paper, we examine the issues of workflow mapping and execution in opportunistic environments such as the Grid. As applications become ever more complex, the process of choosing the appropriate resources and successfully executing the application components becomes ever more difficult. This may include extension or reduction of the initial workflow mapping as necessary for the actual execution. In this paper, we focus on the interplay between a workflow‐mapping component that plans the high‐level resource assignments and the workflow executor that oversees the component execution. We concentrate particularly on issues of data management and we draw from the experiences with mapping and execution systems: Pegasus, DAGMan and Stork. Copyright © 2005 John Wiley & Sons, Ltd.
Ewa Deelman, Tevfik Kosar, Carl Kesselman, Miron Livny
Concurr. Comput. Pract. Exp.1
2005 Task scheduling strategies for workflow-based applications in grids
abstract
Grid applications require allocating a large number of heterogeneous tasks to distributed resources. A good allocation is critical for efficient execution. However, many existing grid toolkits use matchmaking strategies that do not consider overall efficiency for the set of tasks to be run. We identify two families of resource allocation algorithms: task-based algorithms, that greedily allocate tasks to resources, and workflow-based algorithms, that search for an efficient allocation for the entire workflow. We compare the behavior of workflow-based algorithms and task-based algorithms, using simulations of workflows drawn from a real application and with varying ratios of computation cost to data transfer cost. We observe that workflow-based approaches have a potential to work better for data-intensive applications even when estimates about future tasks are inaccurate.
Jim Blythe, Ewa Deelman, Yolanda Gil, Karan Vahi, Anirban Mandal, Ken Kennedy
CCGRID3
2005 Preface
Ewa Deelman, Ian J. Taylor
J. Grid Comput.1
2005 Optimizing Grid-Based Workflow Execution
Carl Kesselman, Ewa Deelman
J. Grid Comput.3
2004 Artemis: Integrating Scientific Data on the Grid
Rattapoom Tuchinda, Snehal Thakkar, Yolanda Gil, Ewa Deelman
AAAI4
2004 The Grid2003 Production Grid: Principles and Practice
Ian T. Foster, Jerry Gieraltowski, Scott Gose, Natalia Maltsev, Edward N. May, Alexis A. Rodriguez, Dinanath Sulakhe, A. Vaniachine, Jim Shank, Saul Youssef, David Adams, Richard Baker 0003, Wensheng Deng, Dantong Yu, Iosif Legrand, Conrad Steenberg, M. Anzar Afaq, Eileen Berman, James Annis, L. A. T. Bauerdick, Michael Ernst, Ian Fisk, Lisa Giacchetti, Gregory E. Graham, Anne Heavey, Joseph Kaiser, Nickolai Kuropatkin, Ruth Pordes, Vijay Sekhri, John Weigand, Yujun Wu, Keith Baker, Lawrence Sorrillo, John Huth, Matthew Allen, Leigh Grundhoefer, John Hicks, Fred Luehring, Steve Peck, Robert Quick, Stephen C. Simms, George Fekete, Jan vandenBerg, Kihyeon Cho, Kihwan Kwon, Dongchul Son, Hyoungwoo Park, Shane Canon, Keith R. Jackson, David E. Konerding, Jason Lee 0001, Doug Olson, Iwona Sakrejda, Brian Tierney, Mark Green 0001, Russ Miller, James Letts, Terrence Martin, David Bury, Catalin Dumitrescu, Daniel Engh, Robert W. Gardner, Marco Mambelli, Yuri Smirnov, Jens-S. Vöckler, Michael Wilde, Yong Zhao 0009, Paul Avery, Richard Cavanaugh, Bockjoo Kim, Craig Prescott, Jorge Rodríguez 0002, Andrew Zahn, Shawn McKee, Christopher T. Jordan, James E. Prewett, Timothy L. Thomas, Horst Severini, Ben Clifford, Ewa Deelman, Larry Flon, Carl Kesselman, Gaurang Mehta, Nosa Olomu, Karan Vahi, Kaushik De, Patrick McGuigan, Mark Sosebee, Dan Bradley, Peter Couvares, Alan DeSmet, Carey Kireyev, Erik Paulson 0001, Alain J. Roy, Scott Koranda, Brian Moe, Bobby Brown, Paul Sheldon
HPDC84
2004 Grid-Based Metadata Services
Ewa Deelman, Malcolm P. Atkinson 0001, Ann L. Chervenak, Neil P. Chue Hong, Carl Kesselman, Sonal Patil, Laura Pearlman, Mei-Hui Su
SSDBM1
2003 Transparent Grid Computing: A Knowledge-Based Approach
Jim Blythe, Ewa Deelman, Yolanda Gil, Carl Kesselman
IAAI2
2003 Grid-Based Galaxy Morphology Analysis for the National Virtual Observatory
abstract
As part of the development of the National Virtual Observatory (NVO), a Data Grid for astronomy, we have developed a prototype science application to explore the dynamical history of galaxy clusters by analyzing the galaxies' morphologies. The purpose of the prototype is to investigate how Grid-based technologies can be used to provide specialized computational services within the NVO environment. In this paper we focus on the key enabling technology components, particularly Chimera and Pegasus which are used to create and manage the computational workflow that must be present to deal with the challenging application requirements. We illustrate how the components interplay with each other and can be driven from a special purpose application portal.
Ewa Deelman, Raymond Plante, Carl Kesselman, Mei-Hui Su, Gretchen Greene, Robert J. Hanisch, Niall Gaffney, Antonio Volpicelli, James Annis, Vijay Sekhri, Tamás Budavári, María A. Nieto-Santisteban, William O'Mullane, David Bohlender, Tom McGlynn, Arnold H. Rots, Olga Pevunova
SC1
2003 A Metadata Catalog Service for Data Intensive Applications
abstract
Advances in computational, storage and network technologies as well as middle ware such as the Globus Toolkit allow scientists to expand the sophistication and scope of data-intensive applications. These applications produce and analyze terabytes and petabytes of data that are distributed in millions of files or objects. To manage these large data sets efficiently, metadata or descriptive information about the data needs to be managed. There are various types of metadata, and it is likely that a range of metadata services will exist in Grid environments that are specialized for particular types of metadata cataloguing and discovery. In this paper, we present the design of a Metadata Catalog Service (MCS) that provides a mechanism for storing and accessing descriptive metadata and allows users to query for data items based on desired attributes. We describe our experience in using the MCS with several applications and present a scalability study of the service.
Shishir Bharathi, Ann L. Chervenak, Ewa Deelman, Carl Kesselman, Mary Manohar, Sonal Patil, Laura Pearlman
SC4
2003 Multi-wavelength image space: another Grid-enabled science
abstract
Abstract We describe how the Grid enables new research possibilities in astronomy through multi‐wavelength images. To see sky images in the same pixel space, they must be projected to that space, a computer‐intensive process. There is thus a virtual data space induced that is defined by an image and the applied projection. This virtual data can be created and replicated with Planners and Replica catalog technology developed under the GriPhyN project. We plan to deploy our system (MONTAGE) on the U.S. Teragrid. Grid computing is also needed for ingesting data—computing background correction on each image—which forms a separate virtual data space. Multi‐wavelength images can be used for pushing source detection and statistics by an order of magnitude from current techniques; for optimization of multi‐wavelength image registration for detection and characterization of extended sources; and for detection of new classes of essentially multi‐wavelength astronomical phenomena. The paper discusses both the Grid architecture and the scientific goals. Copyright © 2003 John Wiley & Sons, Ltd.
Roy Williams, G. Bruce Berriman, Ewa Deelman, John Good, Joseph C. Jacob, Carl Kesselman, Carol Lonsdale, Seb Oliver, Thomas A. Prince
Concurr. Comput. Pract. Exp.3
2003 Mapping Abstract Complex Workflows onto Grid Environments
Ewa Deelman, Jim Blythe, Yolanda Gil, Carl Kesselman, Gaurang Mehta, Karan Vahi, Kent Blackburn, Albert Lazzarini, Adam Arbree, Richard Cavanaugh, Scott Koranda
J. Grid Comput.1
2003 High-performance remote access to climate simulation data: a challenge problem for data grid technologies
Ann L. Chervenak, Ewa Deelman, Carl Kesselman, William E. Allcock, Ian T. Foster, Veronika Nefedova, Jason Lee 0001, Alex Sim, Arie Shoshani, Bob Drach, Dean N. Williams, Don Middleton
Parallel Comput.2
2002 GriPhyN and LIGO, Building a Virtual Data Grid for Gravitational Wave Scientists
abstract
Many Physics experiments today generate large volumes of data. That data is then processed in a variety of ways in order to achieve the understanding of fundamental physical phenomena. The goal of the NSF-funded GriPhyN project (Grid Physics Network) is to enable scientists to seamlessly access data whether it is raw experimental data or a data product which is a result of further processing. GriPhyN provides a new degree of transparency in how data-handling and processing capabilities are integrated to deliver data products to end-users or applications, so that requests for such products are easily mapped into computation and/or data access at multiple locations. GriPhyN refers to the set of all data products available to the user as virtual data. Among the physics applications participating in the project is the Laser Interferometer Gravitational-wave Observatory (LIGO), which is being built to observe the gravitational waves predicted by general relativity. We describe our initial design and prototype of a virtual data Grid for LIGO.
Ewa Deelman, Carl Kesselman, Gaurang Mehta, Leila Meshkat, Laura Pearlman, Kent Blackburn, Phil Ehrens, Albert Lazzarini, Roy Williams, Scott Koranda
HPDC1
2002 Giggle: a framework for constructing scalable replica location services
abstract
In wide area computing systems, it is often desirable to create remote read-only copies (replicas) of files. Replication can be used to reduce access latency, improve data locality, and/or increase robustness, scalability and performance for distributed applications. We define a replica location service (RLS) as a system that maintains and provides access to information about the physical locations of copies. An RLS typically functions as one component of a data grid architecture. This paper makes the following contributions. First, we characterize RLS requirements. Next, we describe a parameterized architectural framework, which we name Giggle (for GIGa-scale Global Location Engine), within which a wide range of RLSs can be defined. We define several concrete instantiations of this framework with different performance characteristics. Finally, we present initial performance results for an RLS prototype, demonstrating that RLS systems can be constructed that meet performance goals.
Ann L. Chervenak, Ewa Deelman, Ian T. Foster, Leanne Guy, Wolfgang Hoschek, Adriana Iamnitchi, Carl Kesselman, Peter Z. Kunszt, Matei Ripeanu, Robert Schwartzkopf, Heinz Stockinger, Kurt Stockinger, Brian Tierney
SC2
2002 Compiler-Optimized Simulation of Large-Scale Applications on High Performance Architectures
Vikram S. Adve, Rajive L. Bagrodia, Ewa Deelman, Rizos Sakellariou
J. Parallel Distributed Comput.3
2002 Simulating Spatially Explicit Problems on High Performance Architectures
Ewa Deelman, Boleslaw K. Szymanski
J. Parallel Distributed Comput.1
2001 High-performance remote access to climate simulation data: a challenge problem for data grid technologies
abstract
In numerous scientific disciplines, terabyte and soon petabyte-scale data collections are emerging as critical community resources. A new class of Data Grid infrastructure is required to support management, transport, distributed access to, and analysis of these datasets by potentially thousands of users. Researchers who face this challenge include the Climate Modeling community, which performs long-duration computations accompanied by frequent output of very large files that must be further analyzed. We describe the Earth System Grid prototype, which brings together advanced analysis, replica management, data transfer, request management, and other technologies to support high-performance, interactive analysis of replicated data. We present performance results that demonstrate our ability to manage the location and movement of large datasets from the user's desktop. We report on experiments conducted over SciNET at SC'2000, where we achieved peak performance of 1.55Gb/s and sustained performance of 512.9Mb/s for data transfers between Texas and California.
William E. Allcock, Ian T. Foster, Veronika Nefedova, Ann L. Chervenak, Ewa Deelman, Carl Kesselman, Jason Lee 0001, Alex Sim, Arie Shoshani, Bob Drach, Dean N. Williams
SC5
2000 POEMS: End-to-End Performance Design of Large Parallel Adaptive Computational Systems
abstract
The POEMS project is creating an environment for end-to-end performance modeling of complex parallel and distributed systems, spanning the domains of application software, runtime and operating system software, and hardware architecture. Toward this end, the POEMS framework supports composition of component models from these different domains into an end-to-end system model. This composition can be specified using a generalized graph model of a parallel system, together with interface specifications that carry information about component behaviors and evaluation methods. The POEMS Specification Language compiler will generate an end-to-end system model automatically from such a specification. The components of the target system may be modeled using different modeling paradigms and at various levels of detail. Therefore, evaluation of a POEMS end-to-end system model may require a variety of evaluation tools including specialized equation solvers, queuing network solvers, and discrete event simulators. A single application representation based on static and dynamic task graphs serves as a common workload representation for all these modeling approaches. Sophisticated parallelizing compiler techniques allow this representation to be generated automatically for a given parallel program. POEMS includes a library of predefined analytical and simulation component models of the different domains and a knowledge base that describes performance properties of widely used algorithms. The paper provides an overview of the POEMS methodology and illustrates several of its key components. The modeling capabilities are demonstrated by predicting the performance of alternative configurations of Sweep3D, a benchmark for evaluating wavefront application technologies and high-performance, parallel architectures.
Vikram S. Adve, Rajive L. Bagrodia, James C. Browne, Ewa Deelman, Aditya Dube, Elias N. Houstis, John R. Rice, Rizos Sakellariou, David Sundaram-Stukel, Patricia J. Teller, Mary K. Vernon
IEEE Trans. Software Eng.4
2000 Asynchronous Parallel Simulation of Parallel Programs
abstract
Parallel simulation of parallel programs for large datasets has been shown to offer significant reduction in the execution time of many discrete event models. The paper describes the design and implementation of MPI-SIM, a library for the execution driven parallel simulation of task and data parallel programs. MPI-SIM can be used to predict the performance of existing programs written using MPI for message passing, or written in UC, a data parallel language, compiled to use message passing. The simulation models can be executed sequentially or in parallel. Parallel execution of the models are synchronized using a set of asynchronous conservative protocols. The paper demonstrates how protocol performance is improved by the use of application-level, runtime analysis. The analysis targets the communication patterns of the application. We show the application-level analysis for message passing and data parallel languages. We present the validation and performance results for the simulator for a set of applications that include the NAS Parallel Benchmark suite. The application-level optimization described in the paper yielded significant performance improvements in the simulation of parallel programs, and in some cases completely eliminated the synchronizations in the parallel execution of the simulation model.
Sundeep Prakash, Ewa Deelman, Rajive L. Bagrodia
IEEE Trans. Software Eng.2
1999 Performance Prediction of Large Parallel Applications using Parallel Simulations
abstract
Accurate simulation of large parallel applications can be facilitated with the use of direct execution and parallel discrete event simulation. This paper describes the use of COMPASS, a direct execution-driven, parallel simulator for performance prediction of programs that include both communication and I/O intensive applications. The simulator has been used to predict the performance of such applications on both distributed memory machines like the IBM SP and shared-memory machines like the SGI Origin 2000. The paper illustrates the usefulness of COMPASS as a versatile performance prediction tool. We use both real-world applications and synthetic benchmarks to study application scalability, sensitivity to communication latency, and the interplay between factors like communication pattern and parallel file system caching on application performance. We also show that the simulator is accurate in its predictions and that it is also efficient in its ability to use parallel simulation to reduce its own execution time which, in some cases, has yielded a nearlinear speedup.
Rajive L. Bagrodia, Ewa Deelman, Steven Docy, Thomas Phan
PPoPP2
1999 Compiler-Supported Simulation of Highly Scalable Parallel Applications
abstract
In this paper, we propose and evaluate practical, automatic techniques that exploit compiler analysis to facilitate simulation of very large message-passing systems.We use a compilersynthesized static task graph model to identify the control-flow and the subset of the computations that determine the parallelism, communication and synchronization of the code, and to generate symbolic estimates of sequential task execution times.This information allows us to avoid executing or simulating large portions of the computational code during the simulation.We have used these techniques to integrate the MPI-Sim parallel simulator at UCLA with the Rice dHPF compiler infrastructure.The integrated system can simulate unmodified High Performance Fortran (HPF) programs compiled to the Message-Passing Interface standard (MPI) by the dHPF compiler, and we expect to simulate MPI programs as well.We evaluate the accuracy and benefits of these techniques for three standard benchmarks on a wide range of problem and system sizes.Our results show that the optimized simulator has errors of less than 17% compared with direct program measurement in all the cases we studied, and typically much smaller errors.Furthermore, it requires factors of 5 to 2000 less memory and up to a factor of 10 less time to execute than the original simulator.These dramatic savings allow us to simulate systems and problem sizes 10 to 100 times larger than is possible with the original simulator.
Vikram S. Adve, Rajive L. Bagrodia, Ewa Deelman, Thomas Phan, Rizos Sakellariou
SC3