Kyle Chard

dblp:10/6661 · DBLP profile ↗
← Back
122ranked-venue papers
16as first author
66since 2021 · last 2026
0000-0002-7370-4805ORCID · verified

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

Systems, architecture and hardware · 62 · 6 first-author · 41 since 2021Applied, interdisciplinary, general and emerging computing · 46 · 8 first-author · 20 since 2021Software engineering, systems software and programming languages · 41 · 7 first-author · 19 since 2021Databases, data management, data science and information retrieval · 8 · 1 first-author · 3 since 2021Artificial intelligence and machine learning · 6 · 2 first-author · 2 since 2021Human-computer interaction and ubiquitous computing · 3 · 1 first-author · 2 since 2021
YearPublicationVenuePosition
2026 Towards Transparent Checkpointing with AI-driven Code Generation
abstract
Adding reliable checkpoint/restart support to an MPI scientific application is a time-consuming expert effort that requires deep knowledge of both the application and resilience. We ask whether a frontier large language model can perform this work end-to-end without human intervention. We assemble a benchmark suite of MPI applications spanning diverse domains and computation patterns, and drive an iterative code-generation loop for each application using Anthropic’s Claude Opus 4.7 invoked through the OpenCode CLI. Across six scientific applications, the LLM generates working checkpoint/restart code in 50 minutes on average while consuming 3.4 M tokens per application. The generated code adds negligible overhead during normal failure-free execution on five of six applications and recovers from injected process failures with efficiency comparable to human-engineered checkpoint/restart implementations. These results suggest that automated end-to-end LLM-driven resilience engineering is technically viable today for a meaningful fraction of HPC applications.
Hai Nguyen 0005, Tekin Bicer, Kyle Chard, Ian T. Foster, Bogdan Nicolae
HPDC3
2026 StreamGuard: Low-Overhead Resilience for Real-time HPC Data Streams
abstract
Real-time scientific workflows operate on continuous data streams and must produce timely, high-quality results despite executing on complex, failure-prone infrastructure. Hardware faults, network disruptions, and performance anomalies caused by resource contention or system heterogeneity can severely degrade performance and violate real-time constraints. We focus on strengthening the resilience of the producer–consumer streaming pattern, a fundamental building block of scientific streaming workflows. We present two complementary techniques: (i) a dynamic, asynchronous, non-blocking checkpointing mechanism that preserves progress without interrupting computation, and (ii) a progress-aware load redistribution strategy that detects slow workers and proactively rebalances tasks. Together, these mechanisms maintain forward progress and balanced execution even in highly error-prone environments. Experimental results show that our approach reduces the impact of failures and performance anomalies by up to 6 ×, while introducing less than 1% overhead in failure-free execution.
Hai Nguyen 0005, Bogdan Nicolae, Tekin Bicer, Amal Gueroudji, Matthieu Dorier, Kyle Chard, Ian T. Foster
ICS6
2026 Empowering Scientific Workflows with Federated Agents
abstract
Agentic systems, in which diverse agents cooperate to tackle challenging problems, are exploding in popularity in the AI community. However, existing agentic frameworks take a relatively narrow view of agents, apply a centralized model, and target conversational, cloud-native applications (e.g., LLM-based AI chatbots). In contrast, scientific applications require myriad agents be deployed and managed across diverse cyberinfrastructure. Here we introduce Academy, a modular and extensible middleware designed to deploy autonomous agents across the federated research ecosystem, including HPC systems, experimental facilities, and data repositories. To meet the demands of scientific computing, Academy supports asynchronous execution, heterogeneous resources, high-throughput data flows, and dynamic resource availability. It provides abstractions for expressing stateful agents, managing inter-agent coordination, and integrating computation with experimental control. We present microbenchmark results that demonstrate high performance and scalability in HPC environments. To explore the breadth of applications that can be supported by agentic workflow designs, we also present case studies in materials discovery, astronomy, decentralized learning, and information extraction in which agents are deployed across diverse HPC systems.
Alok Kamatar, J. Gregory Pauloski, Yadu N. Babuji, Ryan Chard, Mansi Sakarvadia, Daniel Babnigg, Ian T. Foster, Kyle Chard
IPDPS8
2026 Flight: A FaaS-based framework for complex and Hierarchical Federated Learning
Nathaniel Hudson 0001, Valérie Hayot-Sasson, Yadu N. Babuji, Matt Baughman, J. Gregory Pauloski, Ryan Chard, Ian T. Foster, Kyle Chard
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.6
2025 Dynostore: A Wide-Area Distribution System for the Management of Data Over Heterogeneous Storage
abstract
Data distribution across different facilities offers benefits such as enhanced resource utilization, increased resilience through replication, and improved performance by processing data near its source. However, managing such data is challenging due to heterogeneous access protocols, disparate authentication models, and the lack of a unified coordination framework. This paper presents DynoStore, a system that manages data across heterogeneous storage systems. At the core of DynoStore are data containers, an abstraction that provides standardized interfaces for seamless data management, irrespective of the underlying storage systems. Multiple data container connections create a cohesive wide-area storage network, ensuring resilience using erasure coding policies. Furthermore, a load-balancing algorithm ensures equitable and efficient utilization of storage resources. We evaluate DynoStore using benchmarks and realworld case studies, including the management of medical and satellite data across geographically distributed environments. Our results demonstrate a 10 % performance improvement compared to centralized cloud-hosted systems while maintaining competitive performance with state-of-the-art solutions such as Redis and IPFS. DynoStore also exhibits superior fault tolerance, withstanding more failures than traditional systems.
Dante D. Sánchez-Gallegos, José Luis González 0002, Maxime Gonthier, Valérie Hayot-Sasson, J. Gregory Pauloski, Haochen Pan, Kyle Chard, Jesús Carretero 0001, Ian T. Foster
CCGrid7
2025 Wrath: Workload Resilience Across Task Hierarchies in Task-Based Parallel Programming Frameworks
abstract
Failures in Task-based Parallel Programming (TBPP) can severely degrade performance and result in incomplete or incorrect outcomes. Existing failure-handling approaches, including reactive, proactive, and resilient methods such as retry and checkpointing mechanisms, often apply uniform retry mechanisms regardless of the root cause of failures, failing to account for the unique characteristics of TBPP frameworks such as heterogeneous resource availability and task-level failures. To address these limitations, we propose Wrath, a novel systematic approach that categorizes failures based on the unique layered structure of TBPP frameworks and defines specific responses to address failures at different layers. Wrath combines a distributed monitoring system and a resilient module to collaboratively address different types of failures in real time. The monitoring system captures execution and resource information, reports failures, and profiles tasks across different layers of TBPP frameworks. The resilient module then categorizes failures and responds with appropriate actions, such as hierarchically retrying failed tasks on suitable resources. Evaluations demonstrate that Wrath significantly improves TBPP robustness, tripling the task success rate and maintaining an application success rate of over 90 % for resolvable failures. Additionally, Wrath can reduce the time to failure by$20 \%-50 \%$, allowing tasks that are destined to fail to be identified and fail more quickly.
Zhuozhao Li, Valérie Hayot-Sasson, Haochen Pan, Maxime Gonthier, J. Gregory Pauloski, Ryan Chard, Kyle Chard, Ian T. Foster
CCGrid8
2025 ControlA: Agentic Workflow Control Mechanisms for Reliable Science
abstract
AI-driven scientific discovery has emerged as a transformative fifth paradigm in research, with agentic AI playing an increasingly prominent role across scientific domains. Agentic AI can enable collaborative AI-human or even fully autonomous decision-making, but it also introduces significant reliability challenges due to the dynamic and evolutionary nature of the AI agents. Specifically, foundation model-powered agents are prone to generating hallucinated, misleading, or adversarial outputs that can propagate silently through workflows and corrupt downstream results. In this paper we present a conceptual framework for a unified approach that integrates agentic workflow-level instrumentation and agent-level safeguards to enhance the reliability of the wider system, particularly critical in science. Embedding these mechanisms into a provenance-augmented infrastructure enables early detection, containment, and recovery from erroneous behavior, ultimately enhancing reliability and reproducibility in AI-assisted scientific workflows.
Amal Gueroudji, Tanwi Mallick, Renan Souza 0001, Rafael Ferreira da Silva, Robert B. Ross, Matthieu Dorier, Philip H. Carns, Kyle Chard, Ian T. Foster
eScience8
2025 Slice-Aware Attention for Quality Control of Breast MRI Segmentations
abstract
Quality control (QC) is essential for ensuring the reliability of automated tumor segmentations in breast magnetic resonance imaging (MRI), which are increasingly used in both clinical and research workflows. Manual QC by experts is subjective, time-consuming, and unscalable, highlighting the need for automated solutions. In this work, we propose a 2D convolutional neural network (CNN) with attention-based slice aggregation to classify the quality of tumor segmentations in breast MRI. Leveraging the large-scale, expert-annotated MAMA-MIA dataset, our method processes 2D slices as 2-channel inputs (image and segmentation) and learns to weight their contribution using an attention mechanism for volume-level classification. Our best-performing configuration (attention + augmentation) achieved 64% accuracy, 0.58 F1, and 0.67 AUC, outperforming other settings. These findings show that augmentation and attention are complementary—augmentation increases slice-level diversity, while attention helps identify informative slices among mostly uninformative ones. While performance was moderate, the results establish a baseline for scalable QC in breast MRI, with considerable potential for improvement through enhanced architectures, richer augmentations, and integration with 3D context.
Edwin Ma, Rachel Gordon, Anna Woodard, Ian T. Foster, Kyle Chard
eScience5
2025 SMURF: Federated Multimodal Retrieval for Scientific Data via Embedding Alignment
abstract
Scientific data spans heterogeneous modalities and is distributed across diverse storage systems. Vector databases have emerged as powerful tools to index high-dimensional embeddings for semantic search, but most existing systems are centralized and focus on indexing a single modality. We introduce SMURF, a system for federated multimodal indexing and retrieval. SMURF aligns pre-computed embeddings from different scientific modalities into a unified embedding space. Our method extends semi-supervised learning techniques to connect disparate embeddings by combining geometry-preserving alignment with weak supervision using a small set of semantically matched data. We evaluate SMURF on cross-modal retrieval tasks using scientific datasets and find that aligned embeddings substantially improve performance—on average, top-5 accuracy nearly doubles (19.18%→38.82%), and nDCG@10, a ranking quality metric, increases by 78% (0.18→0.32) across various retrieval settings. These results highlight the promise of bridging modality gaps to support more effective scientific discovery in federated environments.
Song Young Oh, Arham Khan, Ian T. Foster, Kyle Chard
eScience4
2025 Diamond: Harnessing GPU Resources for Scientific Deep Learning
abstract
Modern research computing cyberinfrastructure, such as ACCESS-CI and NAIRR Pilot, offers GPU resources across geographically distributed clusters to accommodate the increasing needs of scientific deep learning (DL) workloads. Even for high-performance computing (HPC) experts, configuring environments and managing DL workloads across supercomputers remain significant barriers. To address these obstacles, we present Diamond, an open-source platform to simplify and streamline the DL lifecycle on HPC. Diamond provides an intuitive graphical interface that abstracts system-level complexity, enabling users to develop, debug, and deploy DL models with minimal overhead. We identify several challenges in building such a platform, including portability, security, and usability, and propose effective architectural solutions to each. Notably, Diamond enables users to share and reuse DL workload environments across systems and collaborators, reducing redundant setup efforts. Experimental results demonstrate that Diamond reduces the time to first successful deployment by an average of 68%, compared to manual configuration with command lines. The Diamond service is available at https://diamondhpc.ai.
Haotian Xie, Rohan Marwaha, Minu Mathew, Song Bian 0002, Gengcong Yang, Minghao Yan, Yadu N. Babuji, Owen Price, Yinzhi Wang, Volodymyr V. Kindratenko, Shivaram Venkataraman, Kyle Chard, Ian T. Foster, Zhao Zhang 0007
eScience12
2025 Mitigating Memorization in Language Models
abstract
Language models (LMs) can “memorize” information, i.e., encode training data in their weights in such a way that inference-time queries can lead to verbatim regurgitation of that data. This ability to extract training data can be problematic, for example, when data are private or sensitive. In this work, we investigate methods to mitigate memorization: three regularizer-based, three fine-tuning-based, and eleven machine unlearning-based methods, with five of the latter being new methods that we introduce. We also introduce TinyMem, a suite of small, computationally-efficient LMs for the rapid development and evaluation of memorization-mitigation methods. We demonstrate that the mitigation methods that we develop using TinyMem can successfully be applied to production-grade LMs, and we determine via experiment that: regularizer-based mitigation methods are slow and ineffective at curbing memorization; fine-tuning-based methods are effective at curbing memorization, but overly expensive, especially for retaining higher accuracies; and unlearning-based methods are faster and more effective, allowing for the precise localization and removal of memorized information from LM weights prior to inference. We show, in particular, that our proposed unlearning method BalancedSubnet outperforms other mitigation methods at removing memorized information while preserving performance on target tasks.
Mansi Sakarvadia, Aswathy Ajith, Arham Khan, Nathaniel Hudson 0001, Caleb Geniesse, Kyle Chard, Yaoqing Yang 0002, Ian T. Foster, Michael W. Mahoney
ICLR6
2025 Deadline-Aware Scheduling of Mixed-Criticality Tasks
abstract
High-performance computing centers and cloud providers host a wide variety of workloads, ranging from routine calibration tasks with no strict timing requirements to urgent real-time computations that must be completed within hard deadlines. Traditional approaches reserve resources for high-criticality tasks or preempt and kill lower-criticality tasks when necessary, resulting in wasted compute time and longer turnaround times for lower-criticality tasks. We suggest that a better solution is to interleave the execution of critical and non-critical tasks. We formulate a bi-objective optimization problem: guarantee that all critical tasks meet their deadlines, and minimize the maximum flow, defined as the time a task spends in the system, of non-critical tasks. We introduce a formal model, derive an approximation algorithm and a lower bound, and develop several heuristics based on the approximation framework. Through extensive simulations, based on synthetic and real-world workloads, we show that one of our heuristics reduces the maximum flow of non-critical tasks by up to 14% compared to static resource partitioning.
Maxime Gonthier, Kyle Chard, Ian T. Foster, Loris Marchal, Frédéric Vivien
ICPP2
2025 D-Rex: Heterogeneity-Aware Reliability Framework and Adaptive Algorithms for Distributed Storage
abstract
The exponential growth of data necessitates distributed storage models, such as peer-to-peer systems and data federations.While distributed storage can reduce costs and increase reliability, the heterogeneity in storage capacity, I/O performance, and failure rates of storage resources makes their efficient use a challenge.Further, node failures are common and can lead to data unavailability and even data loss.
Maxime Gonthier, Dante D. Sánchez-Gallegos, Haochen Pan, Bogdan Nicolae, Hai Nguyen 0005, Valérie Hayot-Sasson, J. Gregory Pauloski, Jesús Carretero 0001, Kyle Chard, Ian T. Foster
ICS10
2025 Optimizing Fine-Grained Parallelism Through Dynamic Load Balancing on Multi-Socket Many-Core Systems
abstract
Achieving efficient task parallelism on many-core architectures is an important challenge. The widely used GNU OpenMP implementation of the popular OpenMP parallel programming model incurs high overhead for fine-grained, shortrunning tasks due to time spent on runtime synchronization. In this work, we introduce and analyze three key advances that collectively achieve significant performance gains. First, we introduce XQueue, a lock-less concurrent queue implementation to replace GNU's priority task queue and remove the global task lock. Second, we develop a scalable, efficient, and hybrid lock-free/lock-less distributed tree barrier to address the high hardware synchronization overhead from GNU's centralized barrier. Third, we develop two lock-less and NUMA-aware load balancing strategies. We evaluate our implementation using Barcelona OpenMP Task Suite (BOTS) benchmarks. We show that the use of XQueue and the distributed tree barrier can improve performance by up to$1522.8 \times$compared to the original GNU OpenMP. We further show that lock-less load balancing can improve performance by up to$4 \times$compared to GNU OpenMP using XQueue.
Maxime Gonthier, Poornima Nookala, Haochen Pan, Ian T. Foster, Ioan Raicu, Kyle Chard
IPDPS7
2025 What to Support When You're Compressing: The State of Practice Gaps and Opportunities for Scientific Data Compression
abstract
Over the last nearly 20 years, lossy compression has become an essential aspect of HPC applications’ data pipelines, allowing them to overcome limitations in storage capacity and bandwidth and, in some cases, increase computational throughput and capacity. However, with the adoption of lossy compression comes the requirement to assess and control the impact lossy compression has on scientific outcomes. In this work, we take a major step forward in describing the state of practice and by characterizing workloads. We examine applications’ needs and compressors’ capabilities across 9 different supercomputing application domains. We present 24 takeaways that provide best practices for applications, operational impacts for facilities achieving compressed data, and gaps in application needs not addressed by production compressors that point towards opportunities for future compression research.
Franck Cappello, Robert Underwood, Yuri Alexeev, Allison H. Baker, Ebru Bozdag, Martin Burtscher, Kyle Chard, Sheng Di, Kyle Gerard Felker, Paul Christopher O'Grady, Hanqi Guo 0001, Yafan Huang, Peng Jiang 0004, Sian Jin, Petter Johansson, Shaomeng Li, Xin Liang 0001, Erik Lindahl, Peter Lindstrom 0001, Zarija Lukic, Magnus Lundborg, Danylo Lykov, Masaru Nagaso, Kento Sato, Amarjit Singh, Seung Woo Son 0001, Shihui Song, William Tang 0002, Dingwen Tao, Jiannan Tian, Kazutomo Yoshii, Kai Zhao 0008
SC7
2025 XaaS Containers: Performance-Portable Representation With Source and IR Containers
abstract
High-performance computing (HPC) systems and cloud data centers are converging, and containers are becoming the default method of portable software deployment. Yet, while containers simplify software management, they face significant performance challenges in HPC environments as they must sacrifice hardware-specific optimizations to achieve portability. Although HPC containers can use runtime hooks to access optimized MPI libraries and GPU devices, they are limited by application binary interface (ABI) compatibility and cannot overcome the effects of early-stage compilation decisions. Acceleration as a Service (XaaS) proposes a vision of performance-portable containers, where a containerized application should achieve peak performance across all HPC systems. We present a practical realization of this vision through Source and Intermediate Representation (IR) containers, where we delay performance-critical decisions until the target system specification is known. We analyze specialization mechanisms in HPC software and propose a new LLM-assisted method for automatic discovery of specializations. By examining the compilation pipeline, we develop a methodology to build containers optimized for target architectures at deployment time. Our prototype demonstrates that new XaaS containers combine the convenience of containerization with the performance benefits of system-specialized builds.
Marcin Copik, Eiman Alnuaimi, Alok Kamatar, Valérie Hayot-Sasson, Alberto Madonna, Todd Gamblin, Kyle Chard, Ian T. Foster, Torsten Hoefler
SC7
2025 Addressing Reproducibility Challenges in HPC with Continuous Integration
abstract
The high-performance computing (HPC) community has adopted incentive structures to motivate reproducible research, with major conferences awarding badges to papers that meet reproducibility requirements. Yet, many papers do not meet such requirements. The uniqueness of HPC infrastructure and software, coupled with strict access requirements, may limit opportunities for reproducibility. In the absence of resource access, we believe that regular documented testing, through continuous integration (CI), coupled with complete provenance information, can be used as a substitute. Here, we argue that better HPC-compliant CI solutions will improve reproducibility of applications. We present a survey of reproducibility initiatives and describe the barriers to reproducibility in HPC. To address existing limitations, we present a GitHub Action, CORRECT, that enables secure execution of tests on remote HPC resources. We evaluate CORRECT’s usability across three different types of HPC applications, demonstrating the effectiveness of using CORRECT for automating and documenting reproducibility evaluations.
Valérie Hayot-Sasson, Nathaniel Hudson 0001, André Bauer 0001, Maxime Gonthier, Ian T. Foster, Kyle Chard
SC6
2025 Core Hours and Carbon Credits: Incentivizing Sustainability in HPC
abstract
Efforts to reduce the environmental impact of HPC often focus on resource providers, but choices made by users, e.g., concerning where to run, can be equally consequential. Here we present evidence that new accounting methods that charge users for energy used can incentivize significantly more efficient behavior. We first survey 300 HPC users and find that fewer than 30% are aware of their energy consumption, and that energy efficiency is a low priority concern. We then propose two new multi-resource accounting methods that charge for computations based on their energy consumption or carbon footprint, respectively. Finally, we conduct both simulation studies and a user study to evaluate the impact of these two methods on user behavior. We find that while only providing users feedback on their energy use had no impact on their behavior, associating energy with cost incentivized users to select more efficient resources, and use 40% less energy.
Alok Kamatar, Maxime Gonthier, Valérie Hayot-Sasson, André Bauer 0001, Marcin Copik, Raul Castro Fernandez, Torsten Hoefler, Kyle Chard, Ian T. Foster
SC8
2025 Ocelot: An Interactive, Efficient Distributed Compression-As-a-Service Platform With Optimized Data Compression Techniques
abstract
Large volumes of data generated by scientific simulations, genome sequencing, and other applications need to be moved among clusters for data collection/analysis. Data compression techniques have effectively reduced data storage and transfer costs. However, users' requirements on interactively controlling both data quality and compression ratios are non-trivial to fulfill. We propose a novel Compression-as-a-Service (CaaS) platform called Ocelot with four important contributions: (1) It offers real-time visualization, interactive compression, and transfer of scientific datasets. (2) It incorporates new strategies for compressing diverse types of datasets more effectively than traditional methods. (3) It provides an effective method for estimating the compression ratio and execution time of compression tasks. (4) Experiments on multiple real-world datasets on geographically distributed computers show that Ocelot can significantly improve data transfer efficiency with a performance gain of more than 10x in computing clusters with relatively slow networks.
Yuanjian Liu, Sheng Di, Jiajun Huang 0001, Kyle Chard, Ian T. Foster
IEEE Trans. Parallel Distributed Syst.5
2025 Object Proxy Patterns for Accelerating Distributed Applications
abstract
Workflow and serverless frameworks have empowered new approaches to distributed application design by abstracting compute resources. However, their typically limited or one-size-fits-all support for advanced data flow patterns leaves optimization to the application programmer—optimization that becomes more difficult as data become larger. The transparent object proxy, which provides wide-area references that can resolve to data regardless of location, has been demonstrated as an effective low-level building block in such situations. Here we propose three high-level proxy-based programming patterns—distributed futures, streaming, and ownership—that make the power of the proxy pattern usable for more complex and dynamic distributed program structures. We motivate these patterns via careful review of application requirements and describe implementations of each pattern. We evaluate our implementations through a suite of benchmarks and by applying them in three meaningful scientific applications, in which we demonstrate substantial improvements in runtime, throughput, and memory usage.
J. Gregory Pauloski, Valérie Hayot-Sasson, Logan T. Ward, Alex Brace, André Bauer 0001, Kyle Chard, Ian T. Foster
IEEE Trans. Parallel Distributed Syst.6
2024 An Empirical Investigation of Container Building Strategies and Warm Times to Reduce Cold Starts in Scientific Computing Serverless Functions
abstract
Serverless computing has revolutionized application development and deployment by abstracting infrastructure management, allowing developers to focus on writing code. To do so, serverless platforms dynamically create execution environments, often using containers. The cost to create and deploy these environments is known as "cold start" latency, and this cost can be particularly detrimental to scientific computing workloads characterized by sporadic and dynamic demands. We investigate methods to mitigate cold start issues in scientific computing applications by pre-installing Python packages in container images. Using data from Globus Compute and Binder, we empirically analyze cold start behavior and evaluate four strategies for building containers, including fully pre-built environments and dynamic, on-demand installations. Our results show that pre-installing all packages reduces initial cold start time but requires significant storage. Conversely, dynamic installation offers lower storage requirements but incurs repetitive delays. Additionally, we implemented a simulator and assessed the impact of different warm times, finding that moderate warm times significantly reduce cold starts without the excessive overhead of maintaining always-hot states.
André Bauer 0001, Maxime Gonthier, Haochen Pan, Ryan Chard, Daniel Grzenda, Martin Sträßer, J. Gregory Pauloski, Alok Kamatar, Matt Baughman, Nathaniel Hudson 0001, Ian T. Foster, Kyle Chard
e-Science12
2024 Diaspora: Resilience-Enabling Services for Real-Time Distributed Workflows
abstract
The need for real-time processing to enable automated decision making and experimental steering has driven a shift from high-performance computing workflows on a centralized system to a distributed approach that integrates remote data sources, edge devices, and diverse compute facilities. Under this paradigm, data can be processed close to the source where it is generated, thus reducing latency and bandwidth usage. System resilience is thus a key challenge, requiring distributed workflows to survive component failures and to meet stringent quality-of-service requirements, which results in the need to mitigate anomalies such as congestion and low availability of resources. To address these challenges, we propose Diaspora, a unified resilience framework that is inspired by event-driven communication patterns used in public clouds. Specifically, we propose an event fabric that extends across sites, facilities, and computations to provide timely, reliable, and accurate information about data, application, and resource status. On top of the event fabric, we build resilience-enabling services that combine QoS-aware data streaming, resilient data views, resilient compute and data resources, and anomaly detection and prediction, all of which collectively enhance workflow resilience for these scientific cases.
Bogdan Nicolae, Justin M. Wozniak, Tekin Bicer, Hai Nguyen 0005, Haochen Pan, Amal Gueroudji, Maxime Gonthier, Valérie Hayot-Sasson, Eliu A. Huerta, Kyle Chard, Ryan Chard, Matthieu Dorier, Nageswara S. V. Rao, Anees Al-Najjar, Alessandra Corsi, Ian T. Foster
e-Science11
2024 TaPS: A Performance Evaluation Suite for Task-based Execution Frameworks
abstract
Task-based execution frameworks, such as parallel programming libraries, computational workflow systems, and function-as-a-service platforms, enable the composition of distinct tasks into a single, unified application designed to achieve a computational goal and abstract the parallel and distributed execution of those tasks on arbitrary hardware. Research into these task executors has accelerated as computational sciences increasingly need to take advantage of parallel compute and/or heterogeneous hardware. However, the lack of evaluation standards makes it challenging to compare and contrast novel systems against existing implementations. Here, we introduce TaPS, the Task Performance Suite, to support continued research in distributed task executor frameworks. TaPS provides (1) a unified, modular interface for writing and evaluating applications using arbitrary execution frameworks and data management systems and (2) an initial set of reference synthetic and real-world science applications. We discuss how the design of TaPS supports the reliable evaluation of frameworks and demonstrate TaPS through a survey of benchmarks using the provided reference applications.
J. Gregory Pauloski, Valérie Hayot-Sasson, Maxime Gonthier, Nathaniel Hudson 0001, Haochen Pan, Ian T. Foster, Kyle Chard
e-Science8
2024 Accelerating Function-Centric Applications by Discovering, Distributing, and Retaining Reusable Context in Workflow Systems
abstract
Workflow systems provide a convenient way for users to write large-scale applications by composing independent tasks into large graphs that can be executed concurrently on high-performance clusters. In many newer workflow systems, tasks are often expressed as a combination of function invocations in a high-level language. Because necessary code and data are not statically known prior to execution, they must be moved into the cluster at runtime. An obvious way of doing this is to translate function invocations into self-contained executable programs and run them as usual, but this brings a hefty performance penalty: a function invocation now needs to piggyback its context with extra code and data to a remote node, and the remote node needs to take extra time to reconstruct the invocation's context before executing it, both detrimental to lightweight short-running functions.
Thanh Son Phung, Colin Thomas, Logan T. Ward, Kyle Chard, Douglas Thain
HPDC4
2024 UniFaaS: Programming across Distributed Cyberinfrastructure with Federated Function Serving
abstract
Modern scientific applications are increasingly decomposable into individual functions that may be deployed across distributed and diverse cyberinfrastructure such as supercomputers, clouds, and accelerators. Such applications call for new approaches to programming, distributed execution, and function-level management. We present UniFaaS, a parallel programming framework that relies on a federated function-as-a-service (FaaS) model to enable composition of distributed, scalable, and high-performance scientific workflows, and to support fine-grained function-level management. UniFaaS provides a unified programming interface to compose dynamic task graphs with transparent wide-area data management. UniFaaS exploits an observe-predict-decide approach to efficiently map workflow tasks to target heterogeneous and dynamic resources. We propose a dynamic heterogeneity-aware scheduling algorithm that employs a delay mechanism and a re-scheduling mechanism to accommodate dynamic resource capacity. Our experiments show that UniFaaS can efficiently execute workflows across computing resources with minimal scheduling overhead. We show that UniFaaS can improve the performance of a real-world drug screening workflow by as much as 22.99% when employing an additional 19.48% of resources and a montage workflow by 54.41% when employing an additional 47.83% of resources across multiple distributed clusters, in contrast to using a single cluster.
Ryan Chard, Yadu N. Babuji, Kyle Chard, Ian T. Foster, Zhuozhao Li
IPDPS4
2024 The globus compute dataset: An open function-as-a-service dataset from the edge to the cloud
André Bauer 0001, Haochen Pan, Ryan Chard, Yadu N. Babuji, Josh Bryan, Devesh Tiwari, Ian T. Foster, Kyle Chard
Future Gener. Comput. Syst.8
2024 QoS-aware edge AI placement and scheduling with multiple implementations in FaaS-based edge computing
Nathaniel Hudson 0001, Hana Khamfroush, Matt Baughman, Daniel Enrique Lucani, Kyle Chard, Ian T. Foster
Future Gener. Comput. Syst.5
2024 X-OpenMP - eXtreme fine-grained tasking using lock-less work stealing
Poornima Nookala, Kyle Chard, Ioan Raicu
Future Gener. Comput. Syst.2
2024 SCIPIS: Scalable and concurrent persistent indexing and search in high-end computing systems
Alexandru Iulian Orhean, Anna Giannakou, Lavanya Ramakrishnan, Kyle Chard, Boris Glavic, Ioan Raicu
J. Parallel Distributed Comput.4
2023 Trillion Parameter AI Serving Infrastructure for Scientific Discovery: A Survey and Vision
abstract
Deep learning methods are transforming research, enabling new techniques, and ultimately leading to new discoveries. As the demand for more capable AI models continues to grow, we are now entering an era of Trillion Parameter Models (TPM), or models with more than a trillion parameters---such as Huawei's PanGu-Σ. We describe a vision for the ecosystem of TPM users and providers that caters to the specific needs of the scientific community. We then outline the significant technical challenges and open problems in system design for serving TPMs to enable scientific research and discovery. Specifically, we describe the requirements of a comprehensive software stack and interfaces to support the diverse and flexible requirements of researchers.
Nathaniel Hudson 0001, J. Gregory Pauloski, Matt Baughman, Alok Kamatar, Mansi Sakarvadia, Logan T. Ward, Ryan Chard, André Bauer 0001, Maksim Levental, Will Engler, Owen Price Skelly, Ben Blaiszik, Rick L. Stevens, Kyle Chard, Ian T. Foster
BDCAT15
2023 An Empirical Study of Container Image Configurations and Their Impact on Start Times
abstract
A core selling point of application containers is their fast start times compared to other virtualization approaches like virtual machines. Predictable and fast container start times are crucial for improving and guaranteeing the performance of containerized cloud, serverless, and edge applications. While previous work has investigated container starts, there remains a lack of understanding of how start times may vary across container configurations. We address this shortcoming by presenting and analyzing a dataset of approximately 200,000 open-source Docker Hub images featuring different image configurations (e.g., image size and exposed ports). Leveraging this dataset, we investigate the start times of containers in two environments and identify the most influential features. Our experiments show that container start times can vary between hundreds of milliseconds and tens of seconds in the same environment. Moreover, we conclude that no single dominant configuration feature determines a container's start time, and hardware and software parameters must be considered together for an accurate assessment.
Martin Sträßer, André Bauer 0001, Robert Leppich, Nikolas Herbst, Kyle Chard, Ian T. Foster, Samuel Kounev
CCGrid5
2023 PSI/J: A Portable Interface for Submitting, Monitoring, and Managing Jobs
abstract
It is generally desirable for high-performance computing (HPC) applications to be portable between HPC systems, for example to make use of more performant hardware, make effective use of allocations, and to co-locate compute jobs with large datasets. Unfortunately, moving scientific applications between HPC systems is challenging for various reasons, most notably that HPC systems have different HPC schedulers. We introduce PSI/J, a job management abstraction API intended to simplify the construction of software components and applications that are portable over various HPC scheduler implementations. We argue that such a system is both necessary and that no viable alternative currently exists. We analyze similar notable APIs and attempt to determine the factors that influenced their evolution and adoption by the HPC community. We base the design of PSI/J on that analysis. We describe how PSI/J has been integrated in three workflow systems and one application, and also show via experiments that PSI/J imposes minimal overhead.
Mihael Hategan, André Merzky, Nicholson T. Collier, Ketan Maheshwari, Jonathan Ozik, Matteo Turilli, Andreas Wilke, Justin M. Wozniak, Kyle Chard, Ian T. Foster, Rafael Ferreira da Silva, Shantenu Jha, Daniel E. Laney
e-Science9
2023 Lazy Python Dependency Management in Large-Scale Systems
abstract
Python has become the language of choice for managing many scientific applications. However, when distributing a Python application, it is necessary that all application dependencies be distributed and available in the target execution environment. A specific consequence is that Python workflows suffer from slow scale out due to the time required to import dependencies. We describe ProxyImports, a method to package and distribute Python dependencies in a lazy fashion while remaining transparent and easy to use. Using ProxyImports, Python packages are loaded only once (e.g., by a workflow head node) and are transferred asynchronously to compute nodes. We evaluate our implementation on the Perlmutter and Theta supercomputers and in an HPC cloud-bursting scenario. Our experiments show that ProxyImports significantly reduces the average time to import large modules across an HPC system and demonstrate that this method can be used easily to distribute user-packages to cloud resources. We conclude that ProxyImports improves application runtime, reduces contention on metadata servers and facilitates runtime portability of Python applications.
Alok Kamatar, Mansi Sakarvadia, Valérie Hayot-Sasson, Kyle Chard, Ian T. Foster
e-Science4
2023 Can Automated Metadata Extraction Make Scientific Data More Navigable?
abstract
FAIR principles require that scientific data be findable, discoverable, and reusable by users. To enable FAIRness, practioners of a science repository will often construct a rich, searchable index of metadata derived from the data. Unfortunately, manual metadata annotation methods do not scale to the many data files generated by many projects; and instead automated extraction systems are needed to scalably parse these files—often with nonstandard schema requiring specialized parsing strategies—and deposit representative metadata into a search index. In this work, we evaluate whether, and the extent to which, automatically extracted metadata make research repositories more navigable. We present a two-part user study conducted with scientists at two U.S. national laboratories from projects spanning spectroscopy and battery modeling. We constructed research indexes automatically by using the Xtract metadata extraction system. In the first part of our study, we learned about each user's role and identified key navigation concerns for scientists. We found that participants wished to navigate for purposes of discovery, retrieval, and organization. In the second part, participants completed simulated research data navigation tasks crafted to reflect real-world navigability concerns. We found that regardless of the interface used, participants consistently solved navigation tasks with high degrees of confidence and correctness, and significantly ($1.2\mathrm{X}-50\times$) faster than via their alternative methods (e.g., manual directory scans or designing a customized navigational tool).
Tyler J. Skluzacek, Kyle Chard, Ian T. Foster
e-Science2
2023 SECRE: Surrogate-Based Error-Controlled Lossy Compression Ratio Estimation Framework
abstract
Error-controlled lossy compression has been effective in reducing data storage/transfer costs while preserving reconstructed data fidelity based on user-defined error bounds. State-of-the-art error-controlled lossy compressors primarily fo-cus on error control rather than compression size, and thus, compression ratios are unknown until the compression operation is fully completed. Many use cases, however, require knowledge of compression ratios a priori, for example, pre-allocating appropri-ate memory for the compressed data at runtime. In this paper, we propose a novel, efficient Surrogate-based Error-controlled Lossy Compression Ratio Estimation Framework (SECRE), which includes three key features/contributions. (1) We carefully design the SECRE framework, which, in principle, can be applied to different error-bounded lossy compressors. (2) We implement a compression ratio estimation method for four state-of-the-art error-controlled lossy compressors-SZx, SZ3, ZFP, and SPERR-by devising a corresponding lightweight compression surrogate for each. (3) We evaluate the performance and accuracy of SECRE using four real-world scientific simulation datasets. Experiments show that SECREcan obtain highly accurate com-pression ratio estimates (e.g., ~ 1 % estimation errors for SZx) with low execution overhead (e.g., ~ 2 % estimation cost for SZx).
Arham Khan, Sheng Di, Kai Zhao 0008, Jinyang Liu 0003, Kyle Chard, Ian T. Foster, Franck Cappello
HiPC5
2023 Optimizing Scientific Data Transfer on Globus with Error-Bounded Lossy Compression
abstract
The increasing volume and velocity of science data necessitate the frequent movement of enormous data volumes as part of routine research activities. As a result, limited wide-area bandwidth often leads to bottlenecks in research progress. However, in many cases, consuming applications (e.g., for analysis, visualization, and machine learning) can achieve acceptable performance on reduced-precision data, and thus researchers may wish to compromise on data precision to reduce transfer and storage costs. Error-bounded lossy compression presents a promising approach as it can significantly reduce data volumes while preserving data integrity based on user-specified error bounds. In this paper, we propose a novel data transfer framework called Ocelot that integrates error-bounded lossy compression into the Globus data transfer infrastructure. We note four key contributions: (1) Ocelot is the first integration of lossy compression in Globus to significantly improve scientific data transfer performance over wide area network (WAN). (2) We propose an effective machine-learning based lossy compression quality estimation model that can predict the quality of error-bounded lossy compressors, which is fundamental to ensure that transferred data are acceptable to users. (3) We develop optimized strategies to reduce the compression time overhead, counter the compute-node waiting time, and improve transfer speed for compressed files. (4) We perform evaluations using many real-world scientific applications across different domains and distributed Globus endpoints. Our experiments show that Ocelot can improve dataset transfer performance substantially, and the quality of lossy compression (time, ratio and data distortion) can be predicted accurately for the purpose of quality assurance.
Yuanjian Liu, Sheng Di, Kyle Chard, Ian T. Foster, Franck Cappello
ICDCS3
2023 Accelerating Communications in Federated Applications with Transparent Object Proxies
abstract
Advances in networks, accelerators, and cloud services encourage programmers to reconsider where to compute---such as when fast networks make it cost-effective to compute on remote accelerators despite added latency. Workflow and cloud-hosted serverless computing frameworks can manage multi-step computations spanning federated collections of cloud, high-performance computing (HPC), and edge systems, but passing data among computational steps via cloud storage can incur high costs. Here, we overcome this obstacle with a new programming paradigm that decouples control flow from data flow by extending the pass-by-reference model to distributed applications. We describe ProxyStore, a system that implements this paradigm by providing object proxies that act as wide-area object references with just-in-time resolution. This proxy model enables data producers to communicate data unilaterally, transparently, and efficiently to both local and remote consumers. We demonstrate the benefits of this model with synthetic benchmarks and real-world scientific applications, running across various computing platforms.
J. Gregory Pauloski, Valérie Hayot-Sasson, Logan T. Ward, Nathaniel Hudson 0001, Charlie Sabino, Matt Baughman, Kyle Chard, Ian T. Foster
SC7
2023 Globus automation services: Research process automation across the space-time continuum
abstract
Research process automation–the reliable, efficient, and reproducible execution of linked sets of actions on scientific instruments, computers, data stores, and other resources–has emerged as an essential element of modern science. We report here on new services within the Globus research data management platform that enable the specification of diverse research processes as reusable sets of actions, flows , and the execution of such flows in heterogeneous research environments. To support flows with broad spatial extent (e.g., from scientific instrument to remote data center) and temporal extent (from seconds to weeks), these Globus automation services feature: (1) cloud hosting for reliable execution of even long-lived flows despite sporadic failures; (2) a simple specification and extensible asynchronous action provider API, for defining and executing a wide variety of actions and flows involving heterogeneous resources ; (3) an event-driven execution model for automating execution of flows in response to arbitrary events; and (4) a rich security model enabling authorization delegation mechanisms for secure execution of long-running actions across distributed resources . These services permit researchers to outsource and automate the management of a broad range of research tasks to a reliable, scalable, and secure cloud platform. We present use cases for Globus automation services, describe their design and implementation, present microbenchmark studies, and review experiences applying the services in a range of applications.
Ryan Chard, Jim Pruyne, Kurt McKee, Josh Bryan, Brigitte Raumann, Rachana Ananthakrishnan, Kyle Chard, Ian T. Foster
Future Gener. Comput. Syst.7
2023 Landlord: Coordinating Dynamic Software Environments to Reduce Container Sprawl
abstract
Containers provide customizable software environments that are independent from the system on which they are deployed. Online services for task execution must often generate containers on the fly to meet user-generated requests. However, as the number of users grows and container environments are changed and updated over time, there is an explosion in the number of containers that must be managed, despite the fact that there is significant overlap among many of the containers in use. We analyze a trace of container launches on the public Binder service and demonstrate the performance and resource usage issues associated with container sprawl. We presentLandlord, an algorithm that coalesces related container environments, and show that it can improve container reuse and reduce the number of container builds required in the Binder trace by 40%. We perform a sensitivity analysis ofLandlordusing randomized synthetic workloads on a high-energy physics (HEP) software repository and demonstrate thatLandlordshows benefits for container management across a wide range of usage patterns. Finally, we compareLandlordto offline clustering, and observe that the continuous churn in software necessitates an online approach.
Timothy Shaffer, Thanh Son Phung, Kyle Chard, Douglas Thain
IEEE Trans. Parallel Distributed Syst.3
2022 SCANNS: Towards Scalable and Concurrent Data Indexing and Searching in High-End Computing System
Alexandru Iulian Orhean, Anna Giannakou, Lavanya Ramakrishnan, Kyle Chard, Ioan Raicu
CCGRID4
2022 Exploring Tradeoffs in Federated Learning on Serverless Computing Architectures
abstract
Federated learning is driving the development of new techniques to efficiently and securely use data across multiple sites while using diverse resources. One of these techniques is the use of the serverless computing paradigm to abstract away resource specific configurations, allowing federated learning across heterogeneous environments. However, deploying federated learning across edge resources, the cloud, and traditional HPC sites will require specialized approaches in order to best account for the weaknesses and strengths of each resource. In this work, we explore the new tradeoffs presented by managing a federated learning task across heterogeneous resources and demonstrate these tradeoffs with experiments using a serverless federated learning framework.
Matt Baughman, Ian T. Foster, Kyle Chard
e-Science3
2022 FLoX: Federated Learning with FaaS at the Edge
abstract
Federated learning (FL) is a technique for distributed machine learning that enables the use of siloed and distributed data. With FL, individual machine learning models are trained separately and then only model parameters (e.g., weights in a neural network) are shared and aggregated to create a global model, allowing data to remain in its original environment. While many applications can benefit from FL, existing frameworks are incomplete, cumbersome, and environment-dependent. To address these issues, we present FLoX, an FL framework built on the funcX federated serverless computing platform. FLoX decouples FL model training/inference from infrastructure management and thus enables users to easily deploy FL models on one or more remote computers with a single line of Python code. We evaluate FLoX using three benchmark datasets deployed on ten heterogeneous and distributed compute endpoints. We show that FLoX incurs minimal overhead, especially with respect to the large communication overheads between endpoints for data transfer. We show how balancing the number of samples and epochs with respect to the capacities of participating endpoints can significantly reduce training time with minimal reduction in accuracy. Finally, we show that global models consistently outperform any single model on average by 8%.
Nikita Kotsehub, Matt Baughman, Ryan Chard, Nathaniel Hudson 0001, Panos Patros, Omer F. Rana, Ian T. Foster, Kyle Chard
e-Science8
2022 Automated metadata extraction: challenges and opportunities
abstract
Proper application of the FAIR data principles is what separates a vibrant data ecosystem, in which research data are frequently shared and reused, from a lifeless data graveyard. Automated metadata extraction systems have been proposed as a means of bolstering the findability, interoperability, and reusability of data repositories with little or no human intervention. These extraction systems mine metadata by crawling a repository and applying lightweight extractors that, for various types of file (e.g., image, CSV file), extract or synthesize relevant attributes. In practice, however, the automated creation of generally useful metadata is fraught with challenges. Data consumers may have different perspectives as to what metadata representations are useful, the standards for recording metadata tend to change over time, and the software model for processing updates can introduce unnecessary human and computational effort. Thus, generalizing extraction for a broad audience of data consumers is a difficult and relatively unsolved problem. In this work, we explore these challenges faced by extraction systems in the context of constructing our own extraction system for science data. We first define the metadata extraction problem and provide context to the issues faced in generalizing metadata. Additionally, we identify potential research directions to help alleviate many of these challenges for all automated extraction systems. Ultimately, this work represents a first step in designing Ubiquitous metadata extraction systems that can maximize the value of research data while minimizing the human efforts required in doing so.
Tyler J. Skluzacek, Kyle Chard, Ian T. Foster
e-Science2
2022 HiPS2022: The 2nd Workshop on High Performance Serverless Computing
abstract
Serverless computing presents an attractive model for general distributed computing as it focuses on abstracting the infrastructure required to execute an application. This workshop investigates the intersection between high performance computing and serverless computing, looking both at the high performance and distributed systems used to deliver serverless platforms and also at the use of serverless models for high performance and distributed systems.
Yadu N. Babuji, Kyle Chard, Ian T. Foster, Zhuozhao Li
HPDC2
2022 Summarizing Sets of Related ML-Driven Recommendations for Improving File Management in Cloud Storage
abstract
Personal cloud storage systems increasingly offer recommendations to help users retrieve or manage files of interest. For example, Google Drive’s Quick Access predicts and surfaces files likely to be accessed. However, when multiple, related recommendations are made, interfaces typically present recommended files and any accompanying explanations individually, burdening users. To improve the usability of ML-driven personal information management systems, we propose a new method for summarizing related file-management recommendations. We generate succinct summaries of groups of related files being recommended. Summaries reference the files’ shared characteristics. Through a within-subjects online study in which participants received recommendations for groups of files in their own Google Drive, we compare our summaries to baselines like visualizing a decision tree model or simply listing the files in a group. Compared to the baselines, participants expressed greater understanding and confidence in accepting recommendations when shown our novel recommendation summaries.
Will Brackenbury, Kyle Chard, Aaron J. Elmore, Blase Ur
UIST2
2022 Greedy Nominator Heuristic: Virtual function placement on fog resources
abstract
Abstract Fog computing is an intermediate infrastructure between edge devices (e.g., Internet of Things) and cloud systems that is used to reduce latency in real‐time applications. An application can be composed of a collection of virtual functions, between which dependency constraints can be captured in a service function chain (SFC). Virtual functions within an SFC can be executed at different geo‐distributed locations. However, virtual functions are prone to failure and often do not complete within a deadline. This results in function reallocation to other nodes within the infrastructure; causing delays, potential data loss during function migration, and increased costs. We proposed Greedy Nominator Heuristic (GNH) to address these issues. GNH is based on redundant deployment and failure tracking of virtual functions. GNH places replicas of each function at multiple locations—taking account of expected completion time, failure risk, and cost. We make use of a MapReduce‐based mechanism, where Mappers find suitable locations in parallel, and a Reducer then ranks these locations. Our results show that GNH reduces latency by up to 68%, and is more cost effective than other approaches which rely on state‐of‐the‐art optimization algorithms to allocate replicas.
Osama Almurshed, Omer F. Rana, Kyle Chard
Concurr. Comput. Pract. Exp.3
2022 Data Station: Delegated, Trustworthy, and Auditable Computation to Enable Data-Sharing Consortia with a Data Escrow
abstract
Pooling and sharing data increases and distributes its value. But since data cannot be revoked once shared, scenarios that require controlled release of data for regulatory, privacy, and legal reasons default to not sharing. Because selectively controlling what data to release is difficult, the few data-sharing consortia that exist are often built around data-sharing agreements resulting from long and tedious one-off negotiations. We introduce Data Station, a data escrow designed to enable the formation of data-sharing consortia. Data owners share data with the escrow knowing it will not be released without their consent. Data users delegate their computation to the escrow. The data escrow relies on delegated computation to execute queries without releasing the data first. Data Station leverages hardware enclaves to generate trust among participants, and exploits the centralization of data and computation to generate an audit log. We evaluate Data Station on machine learning and data-sharing applications while running on an untrusted intermediary. In addition to important qualitative advantages, we show that Data Station: i) outperforms federated learning baselines in accuracy and runtime for the machine learning application; ii) is orders of magnitude faster than alternative secure data-sharing frameworks; and iii) introduces small overhead on the critical path.
Siyuan Xia, Zhiru Zhu, Chris Zhu, Kyle Chard, Aaron J. Elmore, Ian T. Foster, Michael J. Franklin, Sanjay Krishnan, Raul Castro Fernandez
Proc. VLDB Endow.5
2022 $f$funcX: Federated Function as a Service for Science
abstract
ƒuncX is a distributed function as a service (FaaS) platform that enables flexible, scalable, and high performance remote function execution. Unlike centralized FaaS systems, ƒuncX decouples the cloud-hosted management functionality from the edge-hosted execution functionality. ƒuncX's endpoint software can be deployed, by users or administrators, on arbitrary laptops, clouds, clusters, and supercomputers, in effect turning them into function serving systems. ƒuncX's cloud-hosted service provides a single location for registering, sharing, and managing both functions and endpoints. It allows for transparent, secure, and reliable function execution across the federated ecosystem of endpoints—enabling users to route functions to endpoints based on specific needs. ƒuncX uses containers (e.g., Docker, Singularity, and Shifter) to provide common execution environments across endpoints. ƒuncX implements various container management strategies to execute functions with high performance and efficiency on diverse ƒuncX endpoints. ƒuncX also integrates with an in-memory data store and Globus for managing data that may span endpoints. We motivate the need for ƒuncX, present our prototype design and implementation, and demonstrate, via experiments on two supercomputers, that ƒuncX can scale to more than 130000 concurrent workers. We show that ƒuncX's container warming-aware routing algorithm can reduce the completion time for 3,000 functions by up to 61% compared to a randomized algorithm and the in-memory data store can speed up data transfers by up to 3x compared to a shared file system.
Zhuozhao Li, Ryan Chard, Yadu N. Babuji, Ben Galewsky, Tyler J. Skluzacek, Kirill Nagaitsev, Anna Woodard, Ben Blaiszik, Josh Bryan, Daniel S. Katz, Ian T. Foster, Kyle Chard
IEEE Trans. Parallel Distributed Syst.12
2022 Optimizing Error-Bounded Lossy Compression for Scientific Data With Diverse Constraints
abstract
Vast volumes of data are produced by today's scientific simulations and advanced instruments. These data cannot be stored and transferred efficiently because of limited I/O bandwidth, network speed, and storage capacity. Error-bounded lossy compression can be an effective method for addressing these issues: not only can it significantly reduce data size, but it can also control the data distortion based on user-defined error bounds. In practice, many scientific applications have specific requirements or constraints for lossy compression, in order to guarantee that the reconstructed data are valid for post hoc analysis. For example, some datasets contain irrelevant data that should be isolated in particular and users often have intuition regarding value ranges, geospatial regions, and other data subsets that are crucial for subsequent analysis. Existing state-of-the-art error-bounded lossy compressors, however, do not consider these constraints during compression, resulting in inferior compression ratios with respect to user's post hoc analysis, due to the fact that the data itself provides little or no value for post hoc analysis. In this work we address this issue by proposing an optimized framework that can preserve diverse constraints during the error-bounded lossy compression, e.g., cleaning the irrelevant data, efficiently preserving different precision for multiple value intervals, and allowing users to set diverse precision over both regular and irregular regions. We perform our evaluation on a supercomputer with up to 2,100 cores. Experiments with six real-world applications show that our proposed diverse constraints based error-bounded lossy compressor can obtain a higher visual quality or data fidelity on reconstructed data with the same or even higher compression ratios compared with the traditional state-of-the-art compressor SZ. Our experiments also demonstrate very good scalability in compression performance compared with the I/O throughput of the parallel file system.
Yuanjian Liu, Sheng Di, Kai Zhao 0008, Sian Jin, Kyle Chard, Dingwen Tao, Ian T. Foster, Franck Cappello
IEEE Trans. Parallel Distributed Syst.6
2022 Deep Neural Network Training With Distributed K-FAC
abstract
Scaling deep neural network training to more processors and larger batch sizes is key to reducing end-to-end training time; yet, maintaining comparable convergence and hardware utilization at larger scales is challenging. Increases in training scales have enabled natural gradient optimization methods as a reasonable alternative to stochastic gradient descent and variants thereof. Kronecker-factored Approximate Curvature (K-FAC), a natural gradient method, preconditions gradients with an efficient approximation of the Fisher Information Matrix to improve per-iteration progress when optimizing an objective function. Here we propose a scalable K-FAC algorithm and investigate K-FAC’s applicability in large-scale deep neural network training. Specifically, we explore layer-wise distribution strategies, inverse-free second-order gradient evaluation, and dynamic K-FAC update decoupling, with the goal of preserving convergence while minimizing training time. We evaluate the convergence and scaling properties of our K-FAC gradient preconditioner, for image classification, object detection, and language modeling applications. In all applications, our implementation converges to baseline performance targets in 9–25% less time than the standard first-order optimizers on GPU clusters across a variety of scales.
J. Gregory Pauloski, Lei Huang 0019, Weijia Xu, Kyle Chard, Ian T. Foster, Zhao Zhang 0007
IEEE Trans. Parallel Distributed Syst.4
2021 Federated Function as a Service for eScience
abstract
The function as a service paradigm aims to abstract the complexities of managing computing infrastructure for users. While adoption in industry has been swift, we have yet to see widespread adoption in academia. This is in part due to barriers such as the need to access large research data, diverse hardware requirements, monolithic code bases, and existing systems available to researchers. We describe funcX, a federated functionas-a-service platform that addresses important requirements for use of FaaS in research computing. We outline how funcX has been used in early science deployments.
Yadu N. Babuji, Josh Bryan, Ryan Chard, Kyle Chard, Ian T. Foster, Ben Galewsky, Daniel S. Katz, Zhuozhao Li
e-Science4
2021 Enhancing Automated FaaS with Cost-aware Provisioning of Cloud Resources
abstract
Compute resources are becoming more diverse, more specialized, and more physically distributed every day. To properly use these resources, the new paradigm of serverless computing aims to abstract away the complexity associated with resource configuration and workload deployment. Building on this serverless architecture are efforts for automated resource selection and automated task distribution and coordination. As computing resources advance, so too does application-specific optimizations, as is the case of deep learning on GPUs or even the more specialized TPUs. To better navigate the efficient and effective use of these diverse resources, we present the addition of automated, cost-aware provisioning of cloud resources to a state-of-the-art automated serverless framework. By automating the selection, provisioning, and configuration of cloud resources, our framework, which we call DELTA+, will enable truly cost-aware usage of the cloud by navigating complex computational tradeoffs. In this proposal, we will outline how we plan to introduce burstable cloud computing to automated serverless infrastructure and to present our initial findings with respect to system design and performance.
Matt Baughman, Ian T. Foster, Kyle Chard
e-Science3
2021 Gold Panning: Automatic Extraction of Scientific Information from Publications
abstract
The millions of scientific papers published every year have made it well beyond any human’s ability to read and curate. As a result, valuable data and potentially ground-breaking insights remain buried deep in the mountain of publications. Unfortunately, while modern Information Extraction techniques show great promise, they are rarely developed for scientific text and often perform poorly in scientific scenarios. We present a Relation Extraction model that automates identification of relationships between entity pairs. We incorporate dependency parse information, relation-term attention, and specialized word embeddings in our model. Application to a benchmark dataset results in an F1 score of 0.755, outperforming the 0.682 by a more traditional model.
Zhi Hong, Kyle Chard, Ian T. Foster
e-Science2
2021 Ultrafast Focus Detection for Automated Microscopy
abstract
Recent advances in scientific instrumentation technology have led to an increase in the volumes and velocities of data being generated in typical laboratories. Such advanced instruments now necessitate commensurately advanced computing resources and techniques to realize their full potential. In particular, such as in the case of quality control for automated microscopes collecting large volumes of images. We propose a faster than real-time “out-of-focus” detection algorithm for electron microscopy images. Our technique, Multi-scale Histologic Feature Detection (MHFD), adapts classical computer vision techniques and is based on detecting various fine-grained histologic features. We also study the inherent parallelism in the technique in order to employ GPU accelerators. We present preliminary tests that demonstrate the efficacy of MHFD and discuss applying our solution to images produced by a scanning electron microscope in an active connectomics laboratory.
Maksim Levental, Ryan Chard, Kyle Chard, Ian T. Foster, Gregg A. Wildenberg
e-Science3
2021 An Empirical Study of Package Dependencies and Lifetimes in Binder Python Containers
abstract
Containers are widely used in scientific applications as they provide greater precision and flexibility in controlling nearly every aspect of the software environment. They can also be easily shared, enabling researchers to run with the same environment on different host systems—an important requirement for scientific reproducibility. In this work, we studied logs of container launches from Binder, a publicly accessible online service for executing Git repositories. Binder dynamically builds and deploys containers following a recipe stored in the repository. These logs capture usage over several years, and include nearly 14 million container launches of around 70,000 unique repositories. To gain more insight about the types of containers and software environments in use for container-based scientific computing services, we downloaded the user-provided recipe repositories referenced in the logs and captured the software specifications and repository metadata. We discovered a number of interesting trends that may be of interest to site administrators provisioning container-based infrastructure for scientific computing in Python. Based on this analysis, we found that automatically generated containers present unique management challenges that are not well handled by straightforward caching. We used historical metadata on package releases from Pip and Conda to quantify difficulties in keeping previously built containers up to date with changing external dependencies. Finally, we proposed several management strategies for reducing infrastructure costs and improving user experience when managing a large container-based service, and back-tested these strategies against Binder launch activity and historical package metadata to demonstrate the value of dependency-oriented container management.
Timothy Shaffer, Kyle Chard, Douglas Thain
e-Science2
2021 Extreme Scale Survey Simulation with Python Workflows
abstract
The Vera C. Rubin Observatory Legacy Survey of Space and Time (LSST) will soon carry out an unprecedented wide, fast, and deep survey of the sky in multiple optical bands. The data from LSST will open up a new discovery space in astronomy and cosmology, simultaneously providing clues toward addressing burning issues of the day, such as the origin of dark energy and and the nature of dark matter, while at the same time yielding data that will, in turn, pose fresh new questions. To prepare for the imminent arrival of this remarkable data set, it is crucial that the associated scientific communities be able to develop the software needed to analyze it. Computational power now available allows us to generate synthetic data sets that can be used as a realistic training ground for such an effort. This effort raises its own challenges—the need to generate very large simulations of the night sky, scaling up simulation campaigns to large numbers of compute nodes across multiple computing centers with different architectures, and optimizing the complex workload around memory requirements and widely varying wall clock times. We describe here a large-scale workflow that melds together Python code to steer the workflow, Parsl to manage the large-scale distributed execution of workflow components, and containers to carry out the image simulation campaign across multiple sites. Taking advantage of these tools, we developed an extreme-scale computational framework and used it to simulate five years of observations for 300 square degrees of sky area. We describe our experiences and lessons learned in developing this workflow capability, and highlight how the scalability and portability of our approach enabled us to efficiently execute it on up to 4000 compute nodes on two supercomputers.
A. S. Villarreal, Yadu N. Babuji, Thomas D. Uram, Daniel S. Katz, Kyle Chard, Katrin Heitmann
e-Science5
2021 Optimizing Multi-Range based Error-Bounded Lossy Compression for Scientific Datasets
abstract
Vast volumes of scientific data cannot be stored and transferred efficiently because of limited I/O bandwidth, network bandwidth, and storage capacity. Error-bounded lossy compression can be an effective method for resolving these big data issues, since not only can it significantly reduce the data size but it can also control the data distortion based on user-defined error bounds. In practice, many scientific applications have specific data fidelity requirements across different value ranges/intervals of the dataset for the lossy compression, in order to guarantee that the reconstructed data are valid for post hoc analysis. Existing state-of-the-art error-bounded lossy compressors, however, do not support multi-range based error-bounds in the lossy compression, leaving a critical gap that hampers their effective use in practice. In this work, we address this issue by proposing a multi-range based error-bounded lossy compressor based on the state-of-the-art SZ lossy compressor. Our approach allows users to set different error bounds in different value ranges for a compressoin task. We evaluate our approach on several real-world datasets and show that it can obtain a higher visual quality or data fidelity on reconstructed data with the same or even higher compression ratios achieved by SZ.
Yuanjian Liu, Sheng Di, Kai Zhao 0008, Sian Jin, Kyle Chard, Dingwen Tao, Ian T. Foster, Franck Cappello
HiPC6
2021 A Serverless Framework for Distributed Bulk Metadata Extraction
abstract
We introduce Xtract, an automated and scalable system for bulk metadata extraction from large, distributed research data repositories. Xtract orchestrates the application of metadata extractors to groups of files, determining which extractors to apply to each file and, for each extractor and file, where to execute. A hybrid computing model, built on the funcX federated FaaS platform, enables Xtract to balance tradeoffs between extraction time and data transfer costs by dispatching each extraction task to the most appropriate location. Experiments on a range of clouds and supercomputers show that Xtract can efficiently process multi-million-file repositories by orchestrating the concurrent execution of container-based extractors on thousands of nodes. We highlight the flexibility of Xtract by applying it to a large, semi-curated scientific data repository and to an uncurated scientific Google Drive repository. We show that by remotely orchestrating metadata extraction across decentralized storage and compute nodes, Xtract can process large repositories in 50% of the time it takes just to transfer the same data to a machine within the same computing facility. We also show that when transferring data is necessary (e.g., no local compute is available), Xtract can scale to process files as fast as they are received, even over a multi-GB/s network.
Tyler J. Skluzacek, Ryan Wong 0002, Zhuozhao Li, Ryan Chard, Kyle Chard, Ian T. Foster
HPDC5
2021 IMPECCABLE: Integrated Modeling PipelinE for COVID Cure by Assessing Better LEads
abstract
The drug discovery process currently employed in the pharmaceutical industry typically requires about 10 years and $2–3 billion to deliver one new drug. This is both too expensive and too slow, especially in emergencies like the COVID-19 pandemic. In silico methodologies need to be improved both to select better lead compounds, so as to improve the efficiency of later stages in the drug discovery protocol, and to identify those lead compounds more quickly. No known methodological approach can deliver this combination of higher quality and speed. Here, we describe an Integrated Modeling PipEline for COVID Cure by Assessing Better LEads (IMPECCABLE) that employs multiple methodological innovations to overcome this fundamental limitation. We also describe the computational framework that we have developed to support these innovations at scale, and characterize the performance of this framework in terms of throughput, peak performance, and scientific results. We show that individual workflow components deliver 100 × to 1000 × improvement over traditional methods, and that the integration of methods, supported by scalable infrastructure, speeds up drug discovery by orders of magnitudes. IMPECCABLE has screened ∼ 1011 ligands and has been used to discover a promising drug candidate. These capabilities have been used by the US DOE National Virtual Biotechnology Laboratory and the EU Centre of Excellence in Computational Biomedicine.
Aymen Alsaadi, Dario Alfè, Yadu N. Babuji, Agastya Bhati, Ben Blaiszik, Alex Brace, Thomas S. Brettin, Kyle Chard, Ryan Chard, Austin Clyde, Peter V. Coveney, Ian T. Foster, Tom Gibbs, Shantenu Jha, Kristopher Keipert, Dieter Kranzlmüller, Thorsten Kurth, Hyungro Lee, Zhuozhao Li, Gerald Mathias, André Merzky, Alexander Partin, Arvind Ramanathan, Ashka Shah, Abraham C. Stern, Rick L. Stevens, Mikhail Titov, Anda Trifan, Aristeidis Tsaris, Matteo Turilli, Huub J. J. Van Dam, Shunzhou Wan, David Wifling, Junqi Yin
ICPP8
2021 Lightweight Function Monitors for Fine-Grained Management in Large Scale Python Applications
abstract
Python has become a widely used programming language for research, not only for small one-off analyses, but also for complex application pipelines running at supercomputer-scale. Modern parallel programming frameworks for Python present users with a more granular unit of management than traditional Unix processes and batch submissions: the Python function. We review the challenges involved in running native Python functions at scale, and present techniques for dynamically determining a minimal set of dependencies and for assembling a lightweight function monitor (LFM) that captures the software environment and manages resources at the granularity of single functions. We evaluate these techniques in a range of environments, from campus cluster to supercomputer, and show that our advanced dependency management planning and dynamic resource management methods provide superior performance and utilization relative to coarser-grained management approaches, achieving several-fold decrease in execution time for several large Python applications.
Timothy Shaffer, Zhuozhao Li, Benjamín Tovar, Yadu N. Babuji, T. J. Dasso, Zoe Surma, Kyle Chard, Ian T. Foster, Douglas Thain
IPDPS7
2021 Enabling Extremely Fine-grained Parallelism via Scalable Concurrent Queues on Modern Many-core Architectures
abstract
Enabling efficient fine-grained task parallelism is a significant challenge for hardware platforms with increasingly many cores. Existing techniques do not scale to hundreds of threads due to the high cost of synchronization in concurrent data structures. To overcome these limitations we present XQueue, a novel lock-less concurrent queuing system with relaxed ordering semantics that is geared towards realizing scalability up to hundreds of concurrent threads. We demonstrate the scalability of XQueue using microbenchmarks and show that XQueue can deliver concurrent operations with latencies as low as 110 cycles at scales of up to 192 cores (up to 6900× improvement compared to traditional synchronization mechanisms) across our diverse hardware, including x86, ARM, and Power9. The reduced latency allows XQueue to provide orders of magnitude (3300×) better throughput that existing techniques. To evaluate the real-world benefits of XQueue, we integrated XQueue with LLVM OpenMP and evaluated five unmodified benchmarks from the Barcelona OpenMP Task Suite (BOTS) as well as a graph traversal benchmark from the GAP benchmark suite. We compared the XQueue-enabled LLVM OpenMP implementation with the native LLVM and GNU OpenMP versions. Using fine-grained task workloads, XQueue can deliver 4× to 6× speedup compared to native GNU OpenMP and LLVM OpenMP in many cases, with speedups as high as 116× in some cases.
Poornima Nookala, Peter A. Dinda, Kyle C. Hale, Kyle Chard, Ioan Raicu
MASCOTS4
2021 KAISA: an adaptive second-order optimizer framework for deep neural networks
abstract
Kronecker-factored Approximate Curvature (K-FAC) has recently been shown to converge faster in deep neural network (DNN) training than stochastic gradient descent (SGD); however, K-FAC's larger memory footprint hinders its applicability to large models. We present KAISA, a K-FAC-enabled, Adaptable, Improved, and ScAlable second-order optimizer framework that adapts the memory footprint, communication, and computation given specific models and hardware to improve performance and increase scalability. We quantify the tradeoffs between memory and communication cost and evaluate KAISA on large models, including ResNet-50, Mask R-CNN, U-Net, and BERT, on up to 128 NVIDIA A100 GPUs. Compared to the original optimizers, KAISA converges 18.1--36.3% faster across applications with the same global batch size. Under a fixed memory budget, KAISA converges 32.5% and 41.6% faster in ResNet-50 and BERT-Large, respectively. KAISA can balance memory and communication to achieve scaling efficiency equal to or better than the baseline optimizers.
J. Gregory Pauloski, Lei Huang 0019, Shivaram Venkataraman, Kyle Chard, Ian T. Foster, Zhao Zhang 0007
SC5
2021 Files of a Feather Flock Together? Measuring and Modeling How Users Perceive File Similarity in Cloud Storage
abstract
Prior work suggests that users conceptualize the organization of personal collections of digital files through the lens of similarity. However, it is unclear to what degree similar files are actually located near one another (e.g., in the same directory) in actual file collections, or whether leveraging file similarity can improve information retrieval and organization for disorganized collections of files. To this end, we conducted an online study combining automated analysis of 50 Google Drive and Dropbox users' cloud accounts with a survey asking about pairs of files from those accounts. We found that many files located in different parts of file hierarchies were similar in how they were perceived by participants, as well as in their algorithmically extractable features. Participants often wished to co-manage similar files (e.g., deleting one file implied deleting the other file) even if they were far apart in the file hierarchy. To further understand this relationship, we built regression models, finding several algorithmically extractable file features to be predictive of human perceptions of file similarity and desired file co-management. Our findings pave the way for leveraging file similarity to automatically recommend access, move, or delete operations based on users' prior interactions with similar files.
Will Brackenbury, Galen Harrison, Kyle Chard, Aaron J. Elmore, Blase Ur
SIGIR3
2021 KondoCloud: Improving Information Management in Cloud Storage via Recommendations Based on File Similarity
abstract
Users face many challenges in keeping their personal file collections organized. While current file-management interfaces help users retrieve files in disorganized repositories, they do not aid in organization. Pertinent files can be difficult to find, and files that should have been deleted may remain. To help, we designed KondoCloud, a file-browser interface for personal cloud storage. KondoCloud makes machine learning-based recommendations of files users may want to retrieve, move, or delete. These recommendations leverage the intuition that similar files should be managed similarly.
Will Brackenbury, Andrew M. McNutt, Kyle Chard, Aaron J. Elmore, Blase Ur
UIST3
2021 DLHub: Simplifying publication, discovery, and use of machine learning models in science
Zhuozhao Li, Ryan Chard, Logan T. Ward, Kyle Chard, Tyler J. Skluzacek, Yadu N. Babuji, Anna Woodard, Steven Tuecke, Ben Blaiszik, Michael J. Franklin, Ian T. Foster
J. Parallel Distributed Comput.4
2020 funcX: A Federated Function Serving Fabric for Science
abstract
Exploding data volumes and velocities, new computational methods and platforms, and ubiquitous connectivity demand new approaches to computation in the sciences. These new approaches must enable computation to be mobile, so that, for example, it can occur near data, be triggered by events (e.g., arrival of new data), be offloaded to specialized accelerators, or run remotely where resources are available. They also require new design approaches in which monolithic applications can be decomposed into smaller components, that may in turn be executed separately and on the most suitable resources. To address these needs we present funcX---a distributed function as a service (FaaS) platform that enables flexible, scalable, and high performance remote function execution. funcX's endpoint software can transform existing clouds, clusters, and supercomputers into function serving systems, while funcX's cloud-hosted service provides transparent, secure, and reliable function execution across a federated ecosystem of endpoints. We motivate the need for funcX with several scientific case studies, present our prototype design and implementation, show optimizations that deliver throughput in excess of 1 million functions per second, and demonstrate, via experiments on two supercomputers, that funcX can scale to more than more than 130 000 concurrent workers.
Ryan Chard, Yadu N. Babuji, Zhuozhao Li, Tyler J. Skluzacek, Anna Woodard, Ben Blaiszik, Ian T. Foster, Kyle Chard
HPDC8
2019 Measuring, Quantifying, and Predicting the Cost-Accuracy Tradeoff
abstract
Exponentially increasing data volumes, coupled with new modes of analysis have created significant new opportunities for data scientists. However, the stochastic nature of many data science techniques results in tradeoffs between costs and accuracy. For example, machine learning algorithms can be trained iteratively and indefinitely with diminishing returns in terms of accuracy. In this paper we explore the cost-accuracy tradeoff through three representative examples: we vary the number of models in an ensemble, the number of epochs used to train a machine learning model, and the amount of data used to train a machine learning model. We highlight the feasibility and benefits of being able to measure, quantify, and predict cost accuracy tradeoffs by demonstrating the presence and usability of these tradeoffs in two different case studies.
Matt Baughman, Nifesh Chakubaji, Hong Linh Truong 0001, Krists Kreics, Kyle Chard, Ian T. Foster
IEEE BigData5
2019 ParaOpt: Automated Application Parameterization and Optimization for the Cloud
abstract
The variety of instance types available on cloud platforms offers enormous flexibility to match the requirements of applications with available resources. However, selecting the most suitable instance type and configuring an application to optimally execute on that instance type can be complicated and time-consuming. For example, application parallelism flags must match available cores and problem sizes must be tuned to match available memory. As the search space of application configurations can be enormous, we propose an automated approach, called ParaOpt, to automatically explore and tune application configurations on arbitrary cloud instances. ParaOpt supports arbitrary applications, enables use of custom optimization methods, and can be configured with different optimization targets such as runtime and cost. We evaluate ParaOpt by optimizing genomics, molecular dynamics, and machine learning applications with four types of optimizers. We show with as few as 15 parameterized executions of an application, representing between 1.2%-26.7% of the search space, that ParaOpt is able to identify the optimal configuration in 32.7% of experiments and a near-optimal configuration in 83.2% of cases. As a result of using near-optimal configurations, ParaOpt reduces overall execution time by up to 85.8% when compared with using the default configuration.
Ian T. Foster, Ted Summer, Zhuozhao Li, Anna Woodard, Ryan Chard, Matt Baughman, Yadu N. Babuji, Kyle Chard, Jason Pitt
CloudCom9
2019 FSMonitor: Scalable File System Monitoring for Arbitrary Storage Systems
abstract
Data automation, monitoring, and management tools are reliant on being able to detect, report, and respond to file system events. Various data event reporting tools exist for specific operating systems and storage devices, such as inotify for Linux, kqueue for BSD, and FSEvents for macOS. However, these tools are not designed to monitor distributed file systems. Indeed, many cannot scale to monitor many thousands of directories, or simply cannot be applied to distributed file systems. Moreover, each tool implements a custom API and event representation, making the development of generalized and portable event-based applications challenging. As file systems grow in size and become increasingly diverse, there is a need for scalable monitoring solutions that can be applied to a wide range of both distributed and local systems. We present here a generic and scalable file system monitor and event reporting tool, FSMonitor, that provides a file-system-independent event representation and event capture interface. FSMonitor uses a modular Data Storage Interface (DSI) architecture to enable the selection and application of appropriate event monitoring tools to detect and report events from a target file system, and implements efficient and fault-tolerant mechanisms that can detect and report events even on large file systems. We describe and evaluate DSIs for common UNIX, macOS, and Windows storage systems, and for the Lustre distributed file system. Our experiments on a 897 TB Lustre file system show that FSMonitor can capture and process almost 38 000 events per second.
Arnab Kumar Paul, Ryan Chard, Kyle Chard, Steven Tuecke, Ali Raza Butt, Ian T. Foster
CLUSTER3
2019 Serverless Science for Simple, Scalable, and Shareable Scholarship
abstract
The adoption of computation- and data-intensive science, or eScience, makes research progress increasingly dependent on the availability, management, and use of sophisticated cyberinfrastructure. An unfortunate consequence is that researchers face increasingly burdensome demands for managing and maintaining cyberinfrastructure. The advent of virtualization and cloud computing has helped, by allowing outsourcing of some such tasks to reliable and scalable cloud providers. But much more progress is needed before we can create a research cyberinfrastructure that allows researchers to focus on creative thought rather than systems management. We examine here how the emerging paradigm of serverless computing, in which arbitrary functions can be dispatched seamlessly to scalable, secure, and reliable service providers, can move us in that direction. To demonstrate how serverless computing can transform scientific computing, we describe three serverless computing models: service-oriented computing, research automation, and function as a service, presenting illustrative case studies for each.
Kyle Chard, Ian T. Foster
eScience1
2019 Application of BagIt-Serialized Research Object Bundles for Packaging and Re-Execution of Computational Analyses
abstract
In this paper we describe our experience adopting the Research Object Bundle (RO-Bundle) format with BagIt serialization (BagIt-RO) for the design and implementation of "tales" in the Whole Tale platform. A tale is an executable research object intended for the dissemination of computational scientific findings that captures information needed to facilitate understanding, transparency, and re-execution for review and computational reproducibility at the time of publication. We describe the Whole Tale platform and requirements that led to our adoption of BagIt-RO, specifics of our implementation, and discuss migrating to the emerging Research Object Crate (RO-Crate) standard.
Kyle Chard, Thomas Thelen, Matthew J. Turk, Craig Willis, Niall Gaffney, Matthew B. Jones, Kacper Kowalik, Bertram Ludäscher, Timothy M. McPhillips, Jarek Nabrzyski, Victoria Stodden, Ian J. Taylor
eScience1
2019 Active Learning Yields Better Training Data for Scientific Named Entity Recognition
abstract
Despite significant progress in natural language processing, machine learning models require substantial expertannotated training data to perform well in tasks such as named entity recognition (NER) and entity relations extraction. Furthermore, NER is often more complicated when working with scientific text. For example, in polymer science, chemical structure may be encoded using nonstandard naming conventions, the same concept can be expressed using many different terms (synonymy), and authors may refer to polymers with ad-hoc labels. These challenges, which are not unique to polymer science, make it difficult to generate training data, as specialized skills are needed to label text correctly. We have previously designed polyNER, a semi-automated system for efficient identification of scientific entities in text. PolyNER applies word embedding models to generate entity-rich corpora for productive expert labeling, and then uses the resulting labeled data to bootstrap a context-based classifier. PolyNER facilitates a labeling process that is otherwise tedious and expensive. Here, we use active learning to efficiently obtain more annotations from experts and improve performance. Our approach requires just five hours of expert time to achieve discrimination capacity comparable to that of a state-of-the-art chemical NER toolkit.
Roselyne Tchoua, Aswathy Ajith, Zhi Hong, Logan T. Ward, Kyle Chard, Debra Audus, Shrayesh Patel, Juan de Pablo, Ian T. Foster
eScience5
2019 Parsl: Pervasive Parallel Programming in Python
abstract
High-level programming languages such as Python are increasingly used to provide intuitive interfaces to libraries written in lower-level languages and for assembling applications from various components. This migration towards orchestration rather than implementation, coupled with the growing need for parallel computing (e.g., due to big data and the end of Moore's law), necessitates rethinking how parallelism is expressed in programs. Here, we present Parsl, a parallel scripting library that augments Python with simple, scalable, and flexible constructs for encoding parallelism. These constructs allow Parsl to construct a dynamic dependency graph of components that it can then execute efficiently on one or many processors. Parsl is designed for scalability, with an extensible set of executors tailored to different use cases, such as low-latency, high-throughput, or extreme-scale execution. We show, via experiments on the Blue Waters supercomputer, that Parsl executors can allow Python scripts to execute components with as little as 5 ms of overhead, scale to more than 250000 workers across more than 8000 nodes, and process upward of 1200 tasks per second. Other Parsl features simplify the construction and execution of composite programs by supporting elastic provisioning and scaling of infrastructure, fault-tolerant execution, and integrated wide-area data management. We show that these capabilities satisfy the needs of many-task, interactive, online, and machine learning applications in fields such as biology, cosmology, and materials science.
Yadu N. Babuji, Anna Woodard, Zhuozhao Li, Daniel S. Katz, Ben Clifford, Lukasz Lacinski, Ryan Chard, Justin M. Wozniak, Ian T. Foster, Michael Wilde, Kyle Chard
HPDC12
2019 DLHub: Model and Data Serving for Science
abstract
While the Machine Learning (ML) landscape is evolving rapidly, there has been a relative lag in the development of the “learning systems” needed to enable broad adoption. Furthermore, few such systems are designed to support the specialized requirements of scientific ML. Here we present the Data and Learning Hub for science (DLHub), a multi-tenant system that provides both model repository and serving capabilities with a focus on science applications. DLHub addresses two significant shortcomings in current systems. First, its self-service model repository allows users to share, publish, verify, reproduce, and reuse models, and addresses concerns related to model reproducibility by packaging and distributing models and all constituent components. Second, it implements scalable and low-latency serving capabilities that can leverage parallel and distributed computing resources to democratize access to published models through a simple web interface. Unlike other model serving frameworks, DLHub can store and serve any Python 3-compatible model or processing function, plus multiple-function pipelines. We show that relative to other model serving systems including TensorFlow Serving, SageMaker, and Clipper, DLHub provides greater capabilities, comparable performance without memoization and batching, and significantly better performance when the latter two techniques can be employed. We also describe early uses of DLHub for scientific applications.
Ryan Chard, Zhuozhao Li, Kyle Chard, Logan T. Ward, Yadu N. Babuji, Anna Woodard, Steven Tuecke, Ben Blaiszik, Michael J. Franklin, Ian T. Foster
IPDPS3
2019 Computing environments for reproducibility: Capturing the "Whole Tale"
abstract
The act of sharing scientific knowledge is rapidly evolving away from traditional articles and presentations to the delivery of executable objects that integrate the data and computational details (e.g., scripts and workflows) upon which the findings rely. This envisioned coupling of data and process is essential to advancing science but faces technical and institutional barriers. The Whole Tale project aims to address these barriers by connecting computational, data-intensive research efforts with the larger research process—transforming the knowledge discovery and dissemination process into one where data products are united with research articles to create “living publications” or tales. The Whole Tale focuses on the full spectrum of science, empowering users in the long tail of science, and power users with demands for access to big data and compute resources. We report here on the design, architecture, and implementation of the Whole Tale environment.
Adam Brinckman, Kyle Chard, Niall Gaffney, Mihael Hategan, Matthew B. Jones, Kacper Kowalik, Sivakumar Kulasekaran, Bertram Ludäscher, Bryce D. Mecum, Jarek Nabrzyski, Victoria Stodden, Ian J. Taylor, Matthew J. Turk, Kandace Turner
Future Gener. Comput. Syst.2
2019 Co-Operative Resource Allocation: Building an Open Cloud Market Using Shared Infrastructure
abstract
In this paper we present DRIVE, a distributed service-based system designed to facilitate an open economic market for federating Cloud providers. To address the challenges associated with market ownership and operation we propose the use of a co-operative (co-op) infrastructure in which the services that make up DRIVE are hosted across participants' resources. To prevent malicious behavior we use cryptographic, secure and privacy preserving allocation protocols as a means of establishing trust in the allocation infrastructure. We investigate through simulation the effect of different strategies, pricing functions, and penalty models on allocation performance and revenue, and show that the overhead of running DRIVE's services on commodity infrastructure is modest.
Kyle Chard, Kris Bubendorfer
IEEE Trans. Cloud Comput.1
2018 Skluma: An Extensible Metadata Extraction Pipeline for Disorganized Data
abstract
To mitigate the effects of high-velocity data expansion and to automate the organization of filesystems and data repositories, we have developed Skluma-a system that automatically processes a target filesystem or repository, extracts content-and context-based metadata, and organizes extracted metadata for subsequent use. Skluma is able to extract diverse metadata, including aggregate values derived from embedded structured data; named entities and latent topics buried within free-text documents; and content encoded in images. Skluma implements an overarching probabilistic pipeline to extract increasingly specific metadata from files. It applies machine learning methods to determine file types, dynamically prioritizes and then executes a suite of metadata extractors, and explores contextual metadata based on relationships among files. The derived metadata, represented in JSON, describes probabilistic knowledge of each file that may be subsequently used for discovery or organization. Skluma's architecture enables it to be deployed both locally and used as an on-demand, cloud-hosted service to create and execute dynamic extraction workflows on massive numbers of files. It is modular and extensible-allowing users to contribute their own specialized metadata extractors. Thus far we have tested Skluma on local filesystems, remote FTP-accessible servers, and publicly-accessible Globus endpoints. We have demonstrated its efficacy by applying it to a scientific environmental data repository of more than 500,000 files. We show that we can extract metadata from those files with modest cloud costs in a few hours.
Tyler J. Skluzacek, Ryan Chard, Galen Harrison, Paul G. Beckman, Kyle Chard, Ian T. Foster
eScience6
2018 Scalable pCT Image Reconstruction Delivered as a Cloud Service
abstract
We describe a cloud-based medical image reconstruction service designed to meet a real-time and daily demand to reconstruct thousands of images from proton cancer treatment facilities worldwide. Rapid reconstruction of a three-dimensional Proton Computed Tomography (pCT) image can require the transfer of 100 GB of data and use of approximately 120 GPU-enabled compute nodes. The nature of proton therapy means that demand for such a service is sporadic and comes from potentially hundreds of clients worldwide. We thus explore the use of a commercial cloud as a scalable and cost-efficient platform for pCT reconstruction. To address the high performance requirements of this application we leverage Amazon Web Services' GPU-enabled cluster resources that are provisioned with high performance networks between nodes. To support episodic demand, we develop an on-demand multi-user provisioning service that can dynamically provision and resize clusters based on image reconstruction requirements, priorities, and wait times. We compare the performance of our pCT reconstruction service running on commercial cloud resources with that of the same application on dedicated local high performance computing resources. We show that we can achieve scalable and on-demand reconstruction of large scale pCT images for simultaneous multi-client requests, processing images in less than 10 minutes for less than $10 per image.
Ryan Chard, Ravi K. Madduri, Nicholas T. Karonis, Kyle Chard, Kirk L. Duffin, Caesar E. Ordoñez, Thomas D. Uram, Justin Fleischauer, Ian T. Foster, Michael E. Papka, John Winans
IEEE Trans. Cloud Comput.4
2017 Safe Collections and Stewardship on Cloud Kotta
abstract
The increasing collection and use of sensitive datasets in science, coupled with the need for inter-institutional collaboration, poses new challenges for infrastructure and administrative models. While enclaves, such as CLOUD KOTTA, provide for the secure management and analysis of data, they do not yet support the administrative models needed by today's researcher practices. Current approaches often rely on a single data administrator to be responsible for the research activities of multiple analysts. However, this approach is not scalable and is prone to errors. To address these challenges we introduce two new abstractions in CLOUD KOTTA: 'safe collections' and 'stewardship'. Safe-collections define a novel abstraction that refines the scope of a data-store and the policies tied to it. Stewards are a new class of privileged user who own and manage safe-collections. By introducing these abstractions, our aim is to relieve the tension between limiting access to, and promoting research on, sensitive data.
Yadu N. Babuji, Kyle Chard, Eamon Duede, Ian T. Foster
eScience2
2017 Software Defined Cyberinfrastructure for Data Management
abstract
Scientific research is data-centric, relying on the acquisition, management, movement, analysis, and sharing of data. Proficiently managing the end-to-end lifecycle of scientific data is non-trivial and comprises many time consuming and mundane tasks. While individual tasks are not prohibitive, when done repeatedly and frequently they represent a significant strain on researchers. We posit that a better approach is to automate these tasks through a Software Defined Cyberinfrastructure. We have developed R IPPLE to provide such capabilities by automating research data management activities via a programmable and event-based cyber-environment. Users specify high-level management policies, such as data movement and metadata extraction, using intuitive If-Trigger-Then-Action rules. These rules are then autonomously, and reliably, executed and managed by Ripple.
Ryan Chard, Kyle Chard, Steven Tuecke, Ian T. Foster
eScience2
2017 Towards a Hybrid Human-Computer Scientific Information Extraction Pipeline
abstract
The emerging field of materials informatics has the potential to greatly reduce time-to-market and development costs for new materials. The success of such efforts hinges on access to large, high-quality databases of material properties. However, many such data are only to be found encoded in text within esoteric scientific articles, a situation that makes automated extraction difficult and manual extraction time-consuming and error-prone. To address this challenge, we present a hybrid Information Extraction (IE) pipeline to improve the machine-human partnership with respect to extraction quality and person-hours, through a combination of rule-based, machine learning, and crowdsourcing approaches. Our goal is to leverage computer and human strengths to alleviate the burden on human curators by automating initial extraction tasks before prioritizing and assigning specialized curation tasks to humans with different levels of training: using non-experts for straightforward tasks such as validation of higher accuracy results (e.g., completing partial facts) and domain experts for low-certainty results (e.g., reviewing specialized compound labels). To validate our approaches, we focus on the task of extracting the glass transition temperature of polymers from published articles. Applying our approaches to 6 090 articles, we have so far extracted 259 refined data values. We project that this number will grow considerably as we tune our methods and process more articles, to exceed that found in standard, expert-curated polymer data handbooks while also being easier to keep up-to-date. The freely available data can be found on our Polymer Properties Predictor and Database website at http://pppdb.uchicago.edu.
Roselyne Tchoua, Kyle Chard, Debra Audus, Logan T. Ward, Joshua Lequieu, Juan de Pablo, Ian T. Foster
eScience2
2017 Software Defined Cyberinfrastructure
abstract
Within and across thousands of science labs, researchers and students struggle to manage data produced in experiments, simulations, and analyses. Largely manual research data lifecycle management processes mean that much time is wasted, research results are often irreproducible, and data sharing and reuse remain rare. In response, we propose a new approach to data lifecycle management in which researchers are empowered to define the actions to be performed at individual storage systems when data are created or modified: actions such as analysis, transformation, copying, and publication. We term this approach software-defined cyberinfrastructure because users can implement powerful data management policies by deploying rules to local storage systems, much as software-defined networking allows users to configure networks by deploying rules to switches.We argue that this approach can enable a new class of responsive distributed storage infrastructure that will accelerate research innovation by allowing any researcher to associate data workflows with data sources, whether local or remote, for such purposes as data ingest, characterization, indexing, and sharing. We report on early experiments with this approach in the context of experimental science, in which a simple if-trigger-then-action (IFTA) notation is used to define rules.
Ian T. Foster, Ben Blaiszik, Kyle Chard, Ryan Chard
ICDCS3
2017 Probabilistic guarantees of execution duration for Amazon spot instances
abstract
In this paper we propose DrAFTS - a methodology for implementing probabilistic guarantees of instance reliability in the Amazon Spot tier. Amazon offers "unreliable" virtual machine instances (ones that may be terminated at any time) at a potentially large discount relative to "reliable" On-demand and Reserved instances. Our method predicts the "bid values" that users can specify to provision Spot instances which ensure at least a fixed duration of execution with a given probability. We illustrate the method and test its validity using Spot pricing data post facto, both randomly and using real-world workload traces. We also test the efficacy of the method experimentally by using it to launch Spot instances and then observing the instance termination rate. Our results indicate that it is possible to obtain the same level of reliability from unreliable instances that the Amazon service level agreement guarantees for reliable instances with a greatly reduced cost.
Richard Wolski, John Brevik, Ryan Chard, Kyle Chard
SC4
2017 Skluma: A Statistical Learning Pipeline for Taming Unkempt Data Repositories
abstract
Scientists' capacity to make use of existing data is predicated on their ability to find and understand those data. While significant progress has been made with respect to data publication, and indeed one can point to a number of well organized and highly utilized data repositories, there remain many such repositories in which archived data are poorly described and thus impossible to use. We present Skluma---an automated system designed to process vast amounts of data and extract deeply embedded metadata, latent topics, relationships between data, and contextual metadata derived from related documents. We show that Skluma can be used to organize and index a large climate data collection that totals more than 500GB of data in over a half-million files.
Paul G. Beckman, Tyler J. Skluzacek, Kyle Chard, Ian T. Foster
SSDBM3
2017 A social content delivery network for e-Science
abstract
Summary We are in the midst of a scientific data explosion in which the rate of data growth is rapidly increasing. While large‐scale research projects have developed sophisticated data distribution networks to share their data with researchers globally, there is no such support for the many millions of research projects generating data of interest to much smaller audiences (as exemplified by the long tail scientist). In data‐oriented research, every aspect of the research process is influenced by data access. However, sharing and accessing data efficiently as well as lowering access barriers are difficult. In the absence of dedicated large‐scale storage, many have noted that there is an enormous storage capacity available via connected peers, none more so than the storage resources of many research groups. With widespread usage of the content delivery network model for disseminating web content, we believe a similar model can be applied to distributing, sharing, and accessing long tail research data in an e‐Science context. We describe the vision and architecture of a social content delivery network – a model that leverages the social networks of researchers to automatically share and replicate data on peers' resources based upon shared interests and trust. Using this model, we describe a simulator and investigate how aspects such as user activity, geographic distribution, trust, and replica selection algorithms affect data access and storage performance. From these results, we show that socially informed replication strategies are comparable with more general strategies in terms of availability and outperform them in terms of spatial efficiency. Copyright © 2016 John Wiley & Sons, Ltd.
Kyle Chard, Simon Caton, Kai Kugler 0002, Omer F. Rana, Daniel S. Katz
Concurr. Comput. Pract. Exp.1
2016 Cloud Kotta: Enabling secure and scalable data analytics in the cloud
abstract
Distributed communities of researchers rely increasingly on valuable, proprietary, or sensitive datasets. Given the growth of such data, especially in fields new to data-driven research like the social sciences and humanities, coupled with what are often strict and complex data-use agreements, many research communities now require methods that allow secure, scalable and cost-effective storage and analysis. Here we present Cloud Kotta: a cloud-based data management and analytics framework. Cloud Kotta delivers an end-to-end solution for coordinating secure access to large datasets, and an execution model that provides both automated infrastructure scaling and support for executing analytics near to the data. Cloud Kotta implements a fine-grained security model ensuring that only authorized users may access, analyze, and download protected data. It also implements automated methods for acquiring and configuring low-cost storage and compute resources as they are needed. We present the architecture and implementation of Cloud Kotta and demonstrate the advantages it provides in terms of increased performance and flexibility. We show that Cloud Kotta's elastic provisioning model can reduce costs by up to 16x when compared with statically provisioned models.
Yadu N. Babuji, Kyle Chard, Aaron Gerow, Eamon Duede
IEEE BigData2
2016 I'll take that to go: Big data bags and minimal identifiers for exchange of large, complex datasets
abstract
Big data workflows often require the assembly and exchange of complex, multi-element datasets. For example, in biomedical applications, the input to an analytic pipeline can be a dataset consisting thousands of images and genome sequences assembled from diverse repositories, requiring a description of the contents of the dataset in a concise and unambiguous form. Typical approaches to creating datasets for big data workflows assume that all data reside in a single location, requiring costly data marshaling and permitting errors of omission and commission because dataset members are not explicitly specified. We address these issues by proposing simple methods and tools for assembling, sharing, and analyzing large and complex datasets that scientists can easily integrate into their daily workflows. These tools combine a simple and robust method for describing data collections (BDBags), data descriptions (Research Objects), and simple persistent identifiers (Minids) to create a powerful ecosystem of tools and services for big data analysis and sharing. We present these tools and use biomedical case studies to illustrate their use for the rapid assembly, sharing, and analysis of large datasets.
Kyle Chard, Mike D'Arcy, Benjamin D. Heavner, Ian T. Foster, Carl Kesselman, Ravi K. Madduri, Alexis A. Rodriguez, Stian Soiland-Reyes, Carole A. Goble, Kristi Clark, Eric W. Deutsch, Ivo D. Dinov, Nathan D. Price 0001, Arthur W. Toga
IEEE BigData1
2016 An Automated Tool Profiling Service for the Cloud
abstract
Cloud providers offer a diverse set of instance types with varying resource capacities, designed to meet the needs of a broad range of user requirements. While this flexibility is a major benefit of the cloud computing model, it also creates challenges when selecting the most suitable instance type for a given application. Sub-optimal instance selection can result in poor performance and/or increased cost, with significant impacts when applications are executed repeatedly. Yet selecting an optimal instance type is challenging, as each instance type can be configured differently, application performance is dependent on input data and configuration, and instance types and applications are frequently updated. We present a service that supports automatic profiling of application performance on different instance types to create rich application profiles that can be used for comparison, provisioning, and scheduling. This service can dynamically provision cloud instances, automatically deploy and contextualize applications, transfer input datasets, monitor execution performance, and create a composite profile with fine grained resource usage information. We use real usage data from four production genomics gateways and estimate the use of profiles in autonomic provisioning systems can decrease execution time by up to 15.7% and cost by up to 86.6%.
Ryan Chard, Kyle Chard, Bryan C. K. Ng, Kris Bubendorfer, Alexis A. Rodriguez, Ravi K. Madduri, Ian T. Foster
CCGrid2
2016 A secure data enclave and analytics platform for social scientists
abstract
Data-driven research is increasingly ubiquitous and data itself is a defining asset for researchers, particularly in the computational social sciences and humanities. Entire careers and research communities are built around valuable, proprietary or sensitive datasets. However, many existing computation resources fail to support secure and cost-effective storage of data while also enabling secure and flexible analysis of the data. To address these needs we present CLOUD KOTTA, a cloud-based architecture for the secure management and analysis of social science data. CLOUD KOTTA leverages reliable, secure, and scalable cloud resources to deliver capabilities to users, and removes the need for users to manage complicated infrastructure. CLOUD KOTTA implements automated, cost-aware models for efficiently provisioning tiered storage and automatically scaled compute resources. CLOUD KOTTA has been used in production for several months and currently manages approximately 10TB of data and has been used to process more than 5TB of data with over 75,000 CPU hours. It has been used for a broad variety of text analysis workflows, matrix factorization, and various machine learning algorithms, and more broadly, it supports fast, secure and cost-effective research.
Yadu N. Babuji, Kyle Chard, Aaron Gerow, Eamon Duede
eScience2
2016 Globus auth: A research identity and access management platform
abstract
Globus Auth is a foundational identity and access management platform service designed to address unique needs of the science and engineering community. It serves to broker authentication and authorization interactions between end-users, identity providers, resource servers (services), and clients (including web, mobile, desktop, and command line applications, and other services). Globus Auth thus makes it easy, for example, for a researcher to authenticate with one credential, connect to a specific remote storage resource with another identity, and share data with colleagues based on another identity. By eliminating friction associated with the frequent need for multiple accounts, identities, credentials, and groups when using distributed cyberinfrastructure, Globus Auth streamlines the creation, integration, and use of advanced research applications and services. Globus Auth builds upon the OAuth 2 and OpenID Connect specifications to enable standards-compliant integration using existing client libraries. It supports identity federation models that enable diverse identities to be linked together, while also providing delegated access tokens via which client services can obtain short term delegated tokens to access other services. We describe the design and implementation of Globus Auth, and report on experiences integrating it with a range of research resources and services, including the JetStream cloud, XSEDE, NCAR's Research Data Archive, and FaceBase.
Steven Tuecke, Rachana Ananthakrishnan, Kyle Chard, Mattias Lidman, Brendan McCollam, Stephen Rosen, Ian T. Foster
eScience3
2016 The Discovery Cloud: Accelerating and Democratizing Research on a Global Scale
abstract
Modern science and engineering require increasingly sophisticated information technology (IT) for data analysis, simulation, and related tasks. Yet the small to medium laboratories (SMLs) in which the majority of research advances occur increasingly lack the human and financial capital needed to acquire and operate such IT. New methods are needed to provide all researchers with access to state-of-the-art scientific capabilities, regardless of their location and budget. Industry has demonstrated the value of cloud-hosted software-and platform-as-a-service approaches, small businesses that outsource their IT to third-party providers slash costs and accelerate innovation. However, few business cloud services are transferable to science. We thus propose the Discovery Cloud, an ecosystem of new, community-produced services to which SMLs can outsource common activities, from data management and analysis to collaboration and experiment automation. We explain the need for a Discovery Platform to streamline the creation and operation of new and interoperable services, and a Discovery Exchange to facilitate the use and sustainability of Discovery Cloud services. We report on our experiences building early elements of the Discovery Platform in the form of Globus services, and on the experiences of those who have applied those services in innovative applications.
Ian T. Foster, Kyle Chard, Steven Tuecke
IC2E2
2016 Globus Nexus: A Platform-as-a-Service provider of research identity, profile, and group management
Kyle Chard, Mattias Lidman, Brendan McCollam, Josh Bryan, Rachana Ananthakrishnan, Steven Tuecke, Ian T. Foster
Future Gener. Comput. Syst.1
2016 Guest Editors Introduction: Special Issue on Scientific Cloud Computing
abstract
The papers in this special section contribute important advances towards leveraging clouds for scientific applications. The contributions focus on a broad range of topics, including: performance modeling and optimization, data management, resource allocation and scheduling, elasticity, reconfiguration, cost prediction and optimization. Most papers revolve around general techniques and approaches that are agnostic of the applications, while two contributions demonstrate how domain specific scientific applications can be migrated to the cloud.
Kate Keahey, Ioan Raicu, Kyle Chard, Bogdan Nicolae
IEEE Trans. Cloud Comput.3
2015 Cost-Aware Elastic Cloud Provisioning for Scientific Workloads
abstract
Cloud computing provides an efficient model to host and scale scientific applications. While cloud-based approaches can reduce costs as users pay only for the resources used, it is often challenging to scale execution both efficiently and cost-effectively. We describe here a cost-aware elastic cloud provisioner designed to elastically provision cloud infrastructure to execute analyses cost-effectively. The provisioner considers real-time spot instance prices across availability zones, leverages application profiles to optimize instance type selection, over-provisions resources to alleviate bottlenecks caused by oversubscribed instance types, and is capable of reverting to on-demand instances when spot prices exceed thresholds. We evaluate the usage of our cost-aware provisioner using four production scientific gateways and show that it can produce cost savings of up to 97.2% when compared to naive provisioning approaches.
Ryan Chard, Kyle Chard, Kris Bubendorfer, Lukasz Lacinski, Ravi K. Madduri, Ian T. Foster
CLOUD2
2015 Cost-Aware Cloud Provisioning
abstract
Cloud computing is often suggested as a low-cost and scalable model for executing and scaling scientific analyses. However, while the benefits of cloud computing are frequently touted, there are inherent technical challenges associated with scaling execution efficiently and cost-effectively. We describe here a cost-aware elastic provisioner designed to dynamically and cost-effectively provision cloud infrastructure based on the requirements of user-submitted scientific workflows. Our provisioner is used in the Globus Galaxies platform -- a Software-as-a-Service provider of scientific analysis capabilities using commercial cloud infrastructure. Using workloads from production usage of this platform we investigate the performance of our provisioner in terms of cost, spot instance termination rate, and execution time. We demonstrate cost savings across six production gateways of up to 95% and 12% improvement in total execution time when compared to a worst case scenario using a single instance type in a single availability zone.
Ryan Chard, Kyle Chard, Kris Bubendorfer, Lukasz Lacinski, Ravi K. Madduri, Ian T. Foster
e-Science2
2015 Globus Data Publication as a Service: Lowering Barriers to Reproducible Science
abstract
Broad access to the data on which scientific results are based is essential for verification, reproducibility, and extension. Scholarly publication has long been the means to this end. But as data volumes grow, new methods beyond traditional publications are needed for communicating, discovering, and accessing scientific data. We describe data publication capabilities within the Globus research data management service, which supports publication of large datasets, with customizable policies for different institutions and researchers, the ability to publish data directly from both locally owned storage and cloud storage, extensible metadata that can be customized to describe specific attributes of different research domains, flexible publication and curation workflows that can be easily tailored to meet institutional requirements, and public and restricted collections that give complete control over who may access published data. We describe the architecture and implementation of these new capabilities and review early results from pilot projects involving nine research communities that span a range of data sizes, data types, disciplines, and publication policies.
Kyle Chard, Jim Pruyne, Ben Blaiszik, Rachana Ananthakrishnan, Steven Tuecke, Ian T. Foster
e-Science1
2015 Using Active Data to Provide Smart Data Surveillance to E-Science Users
abstract
Modern scientific experiments often involve multiple storage and computing platforms, software tools, and analysis scripts. The resulting heterogeneous environments make data management operations challenging, the significant number of events and the absence of data integration makes it difficult to track data provenance, manage sophisticated analysis processes, and recover from unexpected situations. Current approaches often require costly human intervention and are inherently error prone. The difficulties inherent in managing and manipulating such large and highly distributed datasets also limits automated sharing and collaboration. We study a real world e-Science application involving terabytes of data, using three different analysis and storage platforms, and a number of applications and analysis processes. We demonstrate that using a specialized data life cycle and programming model -- Active Data -- we can easily implement global progress monitoring, and sharing, recover from unexpected events, and automate a range of tasks.
Anthony Simonet, Kyle Chard, Gilles Fedak, Ian T. Foster
PDP2
2015 Globus platform-as-a-service for collaborative science applications
abstract
Globus, developed as Software-as-a-Service (SaaS) for research data management, also provides APIs that constitute a flexible and powerful Platform-as-a-Service (PaaS) to which developers can outsource data management activities such as transfer and sharing, as well as identity, profile and group management. By providing these frequently important but always challenging capabilities as a service, accessible over the network, Globus PaaS streamlines web application development and makes it easy for individuals, teams, and institutions to create collaborative applications such as science gateways for science communities. We introduce the capabilities of this platform and review representative applications.
Rachana Ananthakrishnan, Kyle Chard, Ian T. Foster, Steven Tuecke
Concurr. Comput. Pract. Exp.2
2015 The Globus Galaxies platform: delivering science gateways as a service
abstract
Summary The use of public cloud computers to host sophisticated scientific data and software is transforming scientific practice by enabling broad access to capabilities previously available only to the few. The primary obstacle to more widespread use of public clouds to host scientific software (‘cloud‐based science gateways’) has thus far been the considerable gap between the specialized needs of science applications and the capabilities provided by cloud infrastructures. We describe here a domain‐independent, cloud‐based science gateway platform, the Globus Galaxies platform, which overcomes this gap by providing a set of hosted services that directly address the needs of science gateway developers. The design and implementation of this platform leverages our several years of experience with Globus Genomics, a cloud‐based science gateway that has served more than 200 genomics researchers across 30 institutions. Building on that foundation, we have implemented a platform that leverages the popular Galaxy system for application hosting and workflow execution; Globus services for data transfer, user and group management, and authentication; and a cost‐aware elastic provisioning model specialized for public cloud resources. We describe here the capabilities and architecture of this platform, present six scientific domains in which we have successfully applied it, report on user experiences, and analyze the economics of our deployments. Published 2015. This article is a U.S. Government work and is in the public domain in the USA.
Ravi K. Madduri, Kyle Chard, Ryan Chard, Lukasz Lacinski, Alexis A. Rodriguez, Dinanath Sulakhe, David Kelly, Utpal J. Dave, Ian T. Foster
Concurr. Comput. Pract. Exp.2
2015 Big biomedical data as the key resource for discovery science
abstract
Modern biomedical data collection is generating exponentially more data in a multitude of formats. This flood of complex data poses significant opportunities to discover and understand the critical interplay among such diverse domains as genomics, proteomics, metabolomics, and phenomics, including imaging, biometrics, and clinical data. The Big Data for Discovery Science Center is taking an "-ome to home" approach to discover linkages between these disparate data sources by mining existing databases of proteomic and genomic data, brain images, and clinical assessments. In support of this work, the authors developed new technological capabilities that make it easy for researchers to manage, aggregate, manipulate, integrate, and model large amounts of distributed data. Guided by biological domain expertise, the Center's computational resources and software will reveal relationships and patterns, aiding researchers in identifying biomarkers for the most confounding conditions and diseases, such as Parkinson's and Alzheimer's.
Arthur W. Toga, Ian T. Foster, Carl Kesselman, Ravi K. Madduri, Kyle Chard, Eric W. Deutsch, Nathan D. Price 0001, Gwênlyn Glusman, Benjamin D. Heavner, Ivo D. Dinov, Joseph Ames, John D. Van Horn, Roger Kramer, Leroy E. Hood
J. Am. Medical Informatics Assoc.5
2014 Globus Nexus: Research Identity, Profile, and Group Management as a Service
abstract
Collaborative e-Science applications often need to manage large numbers of user identities, profiles, and groups. However, developing and maintaining such capabilities is often challenging given the plethora of security protocols available and requirements for scalable, robust, and highly available implementations. Globus Nexus is a professionally hosted Platform-as-a-Service that provides these capabilities for collaborative e-Science applications, with a particular focus on the needs of scientific communities. It provides features such as identity provisioning, identity federation, profile management, user-oriented group management, and branded web interfaces that are important to many e-Science applications. Globus Nexus implements best practices approaches for each of these features for example using delegated security protocols such as OAuth, provides sophisticated workflows for actions such as email validation, and implements complex user-defined policies regarding permissible actions. We present here Globus Nexus' capabilities, motivate design choices, and present results that characterize the scalability, reliability, and availability of its implementation and deployment.
Kyle Chard, Mattias Lidman, Josh Bryan, Tom Howe, Brendan McCollam, Rachana Ananthakrishnan, Steven Tuecke, Ian T. Foster
eScience1
2014 On Replica Placement in a Social CDN for e-Science
abstract
Research data is experiencing a seemingly endless increase in both volume and production rate. At the same time, efficiently transferring, storing, and analyzing large scale research data have become major research foci. In this paper, we expand on our approach to sharing data for e-Science: a Social Content Delivery Network (S-CDN). A S-CDN leverages the social networks of researchers to automatically share data and place replicas on peers' resources based upon the premises of trust and interest in shared data. We denote a consumer of shared data as a data follower, similar to the notion of Twitter followers, except we add the element of bilateral authorization to capture a notion of trust. We describe a prototypical implementation for a S-CDN that captures an efficient asynchronous transfer mechanism for data management and replication. In addition, we study via simulation the interplay of user behavior with different replication strategies that capture social as well as more general premises for data sharing. Our results illustrate the opportunities and pitfalls of various replication and data access management strategies. Specifically, we show that socially-informed replication strategies are competitive with more general strategies in terms of availability, and outperform them in terms of spatial efficiency.
Kai Kugler 0002, Simon Caton, Kyle Chard, Daniel S. Katz
eScience3
2014 Experiences building Globus Genomics: a next-generation sequencing analysis service using Galaxy, Globus, and Amazon Web Services
abstract
We describe Globus Genomics, a system that we have developed for rapid analysis of large quantities of next-generation sequencing (NGS) genomic data. This system achieves a high degree of end-to-end automation that encompasses every stage of data analysis including initial data retrieval from remote sequencing centers or storage (via the Globus file transfer system); specification, configuration, and reuse of multi-step processing pipelines (via the Galaxy workflow system); creation of custom Amazon Machine Images and on-demand resource acquisition via a specialized elastic provisioner (on Amazon EC2); and efficient scheduling of these pipelines over many processors (via the HTCondor scheduler). The system allows biomedical researchers to perform rapid analysis of large NGS datasets in a fully automated manner, without software installation or a need for any local computing infrastructure. We report performance and cost results for some representative workloads.
Ravi K. Madduri, Dinanath Sulakhe, Lukasz Lacinski, Bo Liu 0010, Alexis A. Rodriguez, Kyle Chard, Utpal J. Dave, Ian T. Foster
Concurr. Comput. Pract. Exp.6
2014 Cloud-based bioinformatics workflow platform for large-scale next-generation sequencing analyses
Bo Liu 0010, Ravi K. Madduri, Borja Sotomayor, Kyle Chard, Lukasz Lacinski, Utpal J. Dave, Jianqiang Li 0002, Ian T. Foster
J. Biomed. Informatics4
2014 A Social Compute Cloud: Allocating and Sharing Infrastructure Resources via Social Networks
abstract
Social network platforms have rapidly changed the way that people communicate and interact. They have enabled the establishment of, and participation in, digital communities as well as the representation, documentation and exploration of social relationships. We believe that as `apps' become more sophisticated, it will become easier for users to share their own services, resources and data via social networks. To substantiate this, we present a social compute cloud where the provisioning of cloud infrastructure occurs through “friend” relationships. In a social compute cloud, resource owners offer virtualized containers on their personal computer(s) or smart device(s) to their social network. However, as users may have complex preference structures concerning with whom they do or do not wish to share their resources, we investigate, via simulation, how resources can be effectively allocated within a social community offering resources on a best effort basis. In the assessment of social resource allocation, we consider welfare, allocation fairness, and algorithmic runtime. The key findings of this work illustrate how social networks can be leveraged in the construction of cloud computing infrastructures and how resources can be allocated in the presence of user sharing preferences.
Simon Caton, Christian Haas 0003, Kyle Chard, Kris Bubendorfer, Omer F. Rana
IEEE Trans. Serv. Comput.3
2013 Globus Nexus: An identity, profile, and group management platform for science gateways and other collaborative science applications
abstract
Globus Nexus is a flexible and powerful Platform-as-a-Service to which developers can outsource identity, group, and profile management needs. By providing these frequently important but always challenging capabilities as a service, accessible over the network, Globus Nexus streamlines web application development and makes it easy for individuals, teams, and institutions to create collaborative web applications such as science gateways for the science community. We introduce the capabilities of this platform and review representative applications.
Rachana Ananthakrishnan, Josh Bryan, Kyle Chard, Ian T. Foster, Tom Howe, Mattias Lidman, Steven Tuecke
CLUSTER3
2013 Constructing a Social Content Delivery Network for eScience
abstract
Increases in the size of research data and the move towards citizen science, in which everyday users contribute data and analyses, have resulted in a research data deluge. Researchers must now carefully determine how to store, transfer and analyze "Big Data" in collaborative environments. This task is even more complicated when considering budget and locality constraints on data storage and access. In this paper we investigate the potential to construct a Social Content Delivery Network (S-CDN) based upon the social networks that exist between researchers. The S-CDN model builds upon the incentives of collaborative researchers within a given scientific community to address their data challenges collaboratively and in proven trusted settings. In this paper we present a prototype implementation of a S-CDN and investigate the performance of the data transfer mechanisms (using Glob us Online) and the potential cost advantages of this approach.
Kai Kugler 0002, Kyle Chard, Simon Caton, Omer F. Rana, Daniel S. Katz
e-Science2
2013 eScience in the Social Cloud
Kris Bubendorfer, Kyle Chard, John Koshy, Ashfag M. Thaufeeg
Future Gener. Comput. Syst.2
2013 High Performance Resource Allocation Strategies for Computational Economies
abstract
Utility computing models have long been the focus of academic research, and with the recent success of commercial cloud providers, computation and storage is finally being realized as the fifth utility. Computational economies are often proposed as an efficient means of resource allocation, however adoption has been limited due to a lack of performance and high overheads. In this paper, we address the performance limitations of existing economic allocation models by defining strategies to reduce the failure and reallocation rate, increase occupancy and thereby increase the obtainable utilization of the system. The high-performance resource utilization strategies presented can be used by market participants without requiring dramatic changes to the allocation protocol. The strategies considered include overbooking, advanced reservation, just-in-time bidding, and using substitute providers for service delivery. The proposed strategies have been implemented in a distributed metascheduler and evaluated with respect to Grid and cloud deployments. Several diverse synthetic workloads have been used to quantity both the performance benefits and economic implications of these strategies.
Kyle Chard, Kris Bubendorfer
IEEE Trans. Parallel Distributed Syst.1
2012 Experiences in the design and implementation of a Social Cloud for Volunteer Computing
abstract
Volunteer computing provides an alternative computing paradigm for establishing the resources required to support large scale scientific computing. The model is particularly well suited for projects that have high popularity and little available computing infrastructure. The premise of volunteer computing platforms is the contribution of computing resources by individuals for little to no gain. It is therefore difficult to attract and retain contributors to projects. The Social Cloud for Volunteer Computing aims to exploit social engineering principles and the ubiquity of social networks to increase the outreach of volunteer computing, by providing an integrated volunteer computing application and creating gamification algorithms based on social principles to encourage contribution. In this paper we present the development of a production SoCVC, detailing the architecture, implementation and performance of the SoCVC Facebook application and show that the approach proposed could have a high impact on volunteer computing projects.
Ryan Chard, Kris Bubendorfer, Kyle Chard
eScience3
2012 Social Cloud Computing: A Vision for Socially Motivated Resource Sharing
abstract
Online relationships in social networks are often based on real world relationships and can therefore be used to infer a level of trust between users. We propose leveraging these relationships to form a dynamic "Social Cloud,” thereby enabling users to share heterogeneous resources within the context of a social network. In addition, the inherent socially corrective mechanisms (incentives, disincentives) can be used to enable a cloud-based framework for long term sharing with lower privacy concerns and security overheads than are present in traditional cloud environments. Due to the unique nature of the Social Cloud, a social market place is proposed as a means of regulating sharing. The social market is novel, as it uses both social and economic protocols to facilitate trading. This paper defines Social Cloud computing, outlining various aspects of Social Clouds, and demonstrates the approach using a social storage cloud implementation in Facebook.
Kyle Chard, Kris Bubendorfer, Simon Caton, Omer F. Rana
IEEE Trans. Serv. Comput.1
2011 Scalability and cost of a cloud-based approach to medical NLP
abstract
Natural Language Processing (NLP) in the medical field has the potential to dramatically influence the way in which everyday clinical care and medical research is conducted. NLP systems provide access to structured content embedded in raw medical texts, therefore enabling automated processing. There are however, several barriers prohibiting wide spread adoption of NLP technology primarily driven by the complexity and cost. This paper describes an approach and implementation which leverages cloud-based deployment and service-based interfaces to extract, process, synthesize, mine, compare/contrast, explore, and manage medical text data in a flexibly secure and scalable architecture. Through a virtual appliance architecture users are able to discover, deploy and utilize NLP engines on demand without requiring knowledge of the underlying, potentially complex, NLP engine. As highlighted in this paper, the system architecture can scale in several configurations: by increasing the number of instances deployed, the number of NLP engines, and the number of databases.
Kyle Chard, Michael Russell, Yves A. Lussier, Eneida A. Mendonça, Jonathan C. Silverstein
CBMS1
2011 A Social Cloud for Public eResearch
abstract
Scientific researchers faced with extremely large computations or the requirement of storing vast quantities of data have come to rely on distributed computational models like cloud computing. However, distributed computation is typically complex and expensive. The Social Cloud for Public eResearch aims to provide researchers with a platform to exploit social networks to reach out to users who would otherwise be unlikely to donate computational time for scientific and other research oriented projects. In this paper we explore the motivations of users to contribute computational time and examine the various ways these motivations can be catered to through established social networks. We specifically look at integrating Face book and BOINC, and discuss the architecture of the functional system and the novel social engineering algorithms that power it.
John Koshy, Kris Bubendorfer, Kyle Chard
eScience3
2011 Collaborative eResearch in a Social Cloud
abstract
Social networks provide a useful basis for enabling collaboration among groups of individuals. This is applicable not only to social communities but also to the scientific community. Already scientists are leveraging social networking concepts in projects to form groups, share information and communicate with their peers. For scientific projects which require large computing resources, one useful aspect of collaboration is the sharing of computing resources among project members. A social network provides an ideal platform to share these resources. This paper introduces a framework for Social Cloud computing with a view towards collaboration and resource sharing within a scientific community. The architecture of a Social Cloud, where individuals or institutions contribute the capacity of their computing resources by means of Virtual Machines leased through the social network, is outlined. Members of the Social Cloud can contribute, request, and use Virtual Machines from other members, as well as form Virtual Organizations among groups of members.
Ashfag M. Thaufeeg, Kris Bubendorfer, Kyle Chard
eScience3
2010 Social Cloud: Cloud Computing in Social Networks
abstract
With the increasingly ubiquitous nature of Social networks and Cloud computing, users are starting to explore new ways to interact with, and exploit these developing paradigms. Social networks are used to reflect real world relationships that allow users to share information and form connections between one another, essentially creating dynamic Virtual Organizations. We propose leveraging the pre-established trust formed through friend relationships within a Social network to form a dynamic “Social Cloud”, enabling friends to share resources within the context of a Social network. We believe that combining trust relationships with suitable incentive mechanisms (through financial payments or bartering) could provide much more sustainable resource sharing mechanisms. This paper outlines our vision of, and experiences with, creating a Social Storage Cloud, looking specifically at possible market mechanisms that could be used to create a dynamic Cloud infrastructure in a Social network environment.
Kyle Chard, Simon Caton, Omer F. Rana, Kris Bubendorfer
IEEE CLOUD1
2010 High occupancy resource allocation for grid and cloud systems, a study with DRIVE
abstract
Economic models have long been advocated as a means of efficient resource allocation, however they are often criticized due to a lack of performance and high overheads. The widespread adoption of utility computing models as seen in commercial Cloud providers has re-motivated the need for economic allocation mechanisms. The aim of this work is to address some of the performance limitations of existing economic allocation models, by reducing the failure/reallocation rate, increasing occupancy and thereby increasing the obtainable utilization of the system. This paper is a study of high performance resource utilization strategies that can be employed in Grid and Cloud systems. In particular we have implemented and quantified the results for strategies including overbooking, advanced reservation, justin-time bidding and using substitute providers for service delivery. These strategies are analyzed in a meta-scheduling context using synthetic workloads derived from a production Grid trace to quantify the performance benefits obtained.
Kyle Chard, Kris Bubendorfer, Peter Komisarczuk
HPDC1
2009 Wrap Scientific Applications as WSRF Grid Services Using gRAVI
abstract
Web service models are increasingly being used in the Grid community as way to create distributed applications exposing data and/or applications through self describing interfaces. Scientific research is one key field in which the benefits are apparent as individual services can be orchestrated into experimental workflows that model the research process and facilitate verification and extension. However, many applications are not web enabled and the task of creating services from scratch is cumbersome in part due to the range of complex technologies, tools, standards and languages involved. In this paper we present gRAVI, a WSRF Web service wrapping tool that allows scientists to rapidly expose applications, scripts and workflows as Web services. gRAVI generated services include GSI security, Grid scheduling, state notifications, persistence and data staging. All service code, scripts and definition files are created automatically without any developer input. gRAVI services are created in standard Grid Archive files and are able to be moved and deployed to any compliant container with no requirement for any gRAVI or Grid infrastructure on the target machine. gRAVI supports deployment to the open science cloud Nimbus, whilst also being able to parse Taverna workflow definition files to create strongly typed services.
Kyle Chard, Wei Tan 0001, Joshua Boverhof, Ravi K. Madduri, Ian T. Foster
ICWS1
2009 Scientific Workflows as Services in caGrid: A Taverna and gRAVI Approach
abstract
In scientific collaboration platforms such as caGrid, workflow-as-a-service is a useful concept for various reasons, such as easy reuse of workflows, access to remote resources, security concerns, and improved execution performance. We propose a solution for facilitating workflow-as-a-service based on Taverna as the workflow engine and gRAVI as a service wrapping tool. We provide both a generic service to execute all Taverna workflows, and an easy-to-use tool (gRAVI-t) for users to wrap their workflows as workflow-specific services, without developing service code. The signature of the specific service is identical to the corresponding workflow's input/output definition and is therefore more self-explained to workflow users. These two categories of services are useful in different scenarios, respectively. We use a tumor analysis workflow as an example to demonstrate how the workflow-as-a-service approach benefits the execution performance. Finally a conclusion is drawn and future research opportunities are discussed.
Wei Tan 0001, Kyle Chard, Dinanath Sulakhe, Ravi K. Madduri, Ian T. Foster, Stian Soiland-Reyes, Carole A. Goble
ICWS2
2008 A Distributed Economic Meta-scheduler for the Grid
abstract
In this paper we present DRIVE, a novel architecture for a virtual organization (VO) based distributed economic meta-scheduler in which members of the VO collaboratively allocate grid resources. Resource providers joining the VO contribute obligation services to the VO. These contributed services are in effect membership 'dues' and are used in the running of the VO's operations - allocation, advertising, management, etc. We use an auction plug-in mechanism to support arbitrary auction protocols which allows users to choose a protocol based on specific requirements and infrastructural availability. For instance, within a single organization, where internal trust exists, users can achieve maximum allocation performance by choosing a conventional sealed bid auction plug-in. In a global utility Grid no such trust exists. The same meta-scheduler architecture can be used with a more expensive secure auction protocol plug- in which ensures the allocation is carried out fairly in the absence of trust The DRIVE prototype has been implemented as a collection of standard WSRF Web services and includes a distributed implementation of a secure combinatorial Garbled Circuit protocol.
Kyle Chard, Kris Bubendorfer
CCGRID1
2008 Build Grid Enabled Scientific Workflows Using gRAVI and Taverna
abstract
Scientific communities are increasingly exposing information and tools as online services in an effort to abstract complex scientific processes and large data sets. Clients are then able to access services without knowledge of their internal workings therefore simplifying the process of replicating scientific research. Taking a service-oriented approach to science (SOS) facilitates reuse, extension, and scalability of components, whist also making information and tools available to a wider audience. Scientific workflows play a key role in realizing SOS by orchestrating services into well formed logical pipelines created to model the requirements of complex scientific experiments. The task of developing such service-oriented infrastructures is not trivial as developers must create and deploy Web services and then coordinate multiple services into workflows. This paper presents an end-to-end approach for developing SOS-based workflows with the aim of simplifying development, deployment, and execution. In particular, we use gRAVI to wrap applications as WSRF Web services and Taverna to compose and execute workflows. The process is validated through the creation of a real world bioinformatics workflow involving multiple services and complex execution paths.
Kyle Chard, Cem Onyuksel, Wei Tan 0001, Dinanath Sulakhe, Ravi K. Madduri, Ian T. Foster
eScience1
2005 Efficient dynamic resource specifications
abstract
In the effort to reach beyond 3G, researchers have been actively looking at utilizing new models for network based services. Small mobile, pervasive and ubiquitous devices will benefit from networked services and computation provided by utility computing providers and the virtual organizations that lease resources from them. As an additional factor, we believe that it is critical that the mobile, pervasive or ubiquitous devices be able to dynamically manipulate their resource specifications when obtaining services and resources from the utility computing and communication network. This requires a simple, manipulatable, and preferably modular resource specification structure. This paper presents the Resource Description Graph (RDG). The RDG is used to represent available and required resources for hosts and applications in a directed acyclic graph. The RDG has many desirable properties including inherent security, expressiveness, modularity, and composition. We show that the computational time to match RDG resource specifications, thirty resource types and constraints, is less than 1ms --- demonstrating that the RDG is a practical approach to resource specification with a low computational overhead.
Kris Bubendorfer, Peter Komisarczuk, Kyle Chard
Mobile Data Management3