EDBT 2026 Demo / reviewers in the wild / expert
Douglas Thain
dblp:83/1871
· DBLP profile ↗
100ranked-venue papers
13as first author
15since 2021 · last 2026
0000-0001-5218-1956ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 72 · 12 first-author · 11 since 2021Software engineering, systems software and programming languages · 18 · 1 first-author · 4 since 2021Applied, interdisciplinary, general and emerging computing · 14 · 1 first-author · 4 since 2021Security and privacy · 2Databases, data management, data science and information retrieval · 2Artificial intelligence and machine learning · 1Computer networks · 1Graphics, computer vision, multimedia, augmented reality and games · 1
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | Efficiently Reproducing Distributed Workflows in Notebook-based Systems
Talha Azaz, Raza Ahmad, Douglas Thain, Tanu Malik |
CCGrid | 4 |
| 2026 | SciWIND: Effectively Exploiting Node-Local Storage for Data-Intensive High-Energy Physics Workflows
Colin Thomas, Barry Sly-Delgado, Connor Moore, Benjamín Tovar, Kevin Lannon, Douglas Thain |
IPDPS | 7 |
| 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. | 23 |
| 2026 | X-Bucket: A Family of Algorithms for Adaptive Resource Allocation in Dynamic Workflows
Thanh Son Phung, Douglas Thain |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2025 | Backpacks for Notebooks: Enabling Containerized Notebook Workflows in Distributed EnvironmentsabstractNotebooks have become widely adopted in the scientific community due to their interactive interface and ease of sharing. However, using notebooks to execute large-scale scientific workflows remains challenging. Scientific workflows are typically distributed and require resource provisioning and data management prior to execution. Because notebooks do not natively embed workflow specifications, users often resort to inserting custom configuration steps directly within notebook cells to enable provisioning. This practice undermines reproducibility, as the same notebook may not run consistently across different cluster environments. In this paper, we introduce the concept of a notebook backpack—a companion specification that captures the embedded workflow along with all relevant configuration elements. We describe how notebook tracing can be leveraged to automatically populate the backpack. We then describe an integrated tool that provisions a backpack on distributed resources. Using real-world case studies, we demonstrate that the backpack abstraction enables minimal modification of the notebook, portable execution, and cross-site reproducibility of notebook-based workflows on HPC clusters without significantly increasing notebook execution time. Talha Azaz, Raza Ahmad, A. S. M. Shahadat Hossain, Furqan Baig, Shaowen Wang 0001, Kevin Lannon, Tanu Malik, Douglas Thain |
eScience | 9 |
| 2025 | Liberating the Data Aware Scheduler to Achieve Locality in Layered Scientific Workflow SystemsabstractLarge scale scientific workflows are typically run using a task-based distributed workflow system. Many contemporary workflow systems are organized in a layered or modular fashion, with one component managing the DAG or workflow graph structure, and another component acting as the scheduler or executor. These components communicate and regulate the flow of tasks from the application to the remote execution site. Scientific workflows perform significant amounts of I/O, typically using an HPC parallel filesystem. Node-Local storage technologies are often used as an effective supplement to parallel filesystems, relieving them of the large amount of bandwidth consumed by intermediate data reads and writes which would otherwise cause I/O bottlenecks. Node-Local storage techniques depend upon effective scheduling to place tasks close to their necessary data, thus benefiting from data locality. Effective locality based scheduling is a challenge however. The conventional layered architecture results in the scheduler considering tasks on an individual basis with a narrow view of the greater DAG. We present a modified architecture of a workflow system which allows the efficient construction of dependency-based task groups which are passed through the DAG manager, scheduler, and finally to the remote worker node where improved use of data locality can be achieved. These modifications were implemented using the Parsl parallel library and TaskVine execution engine. We evaluate this implementation with a benchmark application and a Montage workflow. We compare the results between a conventional data-aware scheduler, our task grouping implementation, and a workflow system using only shared storage. We find that task grouping achieves a significantly greater degree of data locality, thereby improving performance and reducing total data movement in the cluster when compared to the other two methods. Colin Thomas, Douglas Thain |
eScience | 2 |
| 2024 | Accelerating Function-Centric Applications by Discovering, Distributing, and Retaining Reusable Context in Workflow SystemsabstractWorkflow 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 |
HPDC | 5 |
| 2024 | Adaptive Task-Oriented Resource Allocation for Large Dynamic Workflows on Opportunistic ResourcesabstractDynamic workflow management systems offer a solution to the problem of distributing a local application by packaging individual computations and their dependencies on-the-fly into tasks executable on remote workers. Such independent task execution allows workers to be launched in an opportunistic manner to maximize the current pool of resources at any given time, either through opportunistic systems (e.g., HTCondor, AWS Spot Instances), or conventional systems (e.g., SLURM, SGE) with backfilling enabled, as opposed to monolithic or message-passing applications requiring a fixed block of non-preemptible workers. However, the dynamic nature of task generation presents a significant challenge in terms of resource management as tasks must be allocated with some unknown amount of resources pre-execution but are only observable at runtime. This in turn results in potentially huge resource waste per task as (1) users lack direct knowledge about the relationship between tasks and resources, and thus cannot correctly specify the amount of resources a task needs in advance, and (2) workflows and tasks may exhibit stochastic behaviors at runtime, which complicates the process of resource management.In this paper, we (1) argue for the need of an adaptive resource allocator capable of allocating tasks at runtime and adjusting to random fluctuations and abrupt changes in a dynamic workflow without requiring any prior knowledge, and (2) introduce Greedy Bucketing and Exhaustive Bucketing: two robust, online, general-purpose, and prior-free allocation algorithms capable of producing quality estimates of a task’s resource consumption as the workflow runs. Our results show that a resource allocator equipped with either algorithm consistently outperforms 5 alternative allocation algorithms on 7 diverse workflows and incurs at most 1.6 ms overhead per allocation in the steady state. Thanh Son Phung, Douglas Thain |
IPDPS | 2 |
| 2024 | Reshaping High Energy Physics Applications for Near-Interactive Execution Using TaskVineabstractHigh energy physics experiments produce petabytes of data annually that must be reduced to gain insight into the laws of nature. Early-stage reduction executes long-running, high-throughput workflows across thousands of nodes spanning multiple facilities to produce shared datasets. Later stages are typically written by individuals or small groups and must be refined and re-run many times for correctness. Reducing iteration times of later stages is key to accelerating discovery. We demonstrate our experience reshaping late-stage analysis applications on thousands of nodes. It is not enough merely to increase scale: it is necessary to make changes throughout the stack, including storage systems, data management, task scheduling, and application design. We demonstrate these changes when applied to two analysis applications built on open source data analysis frameworks (Coffea, Dask, TaskVine). We evaluate the performance of the applications on opportunistic campus clusters, showing effective scaling up to 7200 cores, thus producing significant speedup. Barry Sly-Delgado, Benjamín Tovar, Douglas Thain |
SC | 4 |
| 2023 | Mixed Modality Workflows in TaskVineabstractModern scientific workflows desire to mix several different computing modalities: self-contained computational tasks, data-intensive transformations, and serverless function calls. To date, these modalities have required distinct system architectures with different scheduling objectives and constraints. In this paper, we describe how TaskVine, a new workflow execution platform, combines these modalities into an execution platform with shared abstractions. We demonstrate results of the system executing a machine learning workflow with combined standalone tasks and serverless functions. David Simonetti, Benjamín Tovar, Douglas Thain |
HPDC | 3 |
| 2023 | Landlord: Coordinating Dynamic Software Environments to Reduce Container SprawlabstractContainers 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. | 4 |
| 2022 | Robust Meta-Workflow Management with MufasaabstractWorkflow management systems (WMS) are widely used to describe and execute large computational or data intensive applications. However, when a large ensemble of workflows is run on a cluster, new resource management problems occur. Each WMS itself consumes otherwise unmanaged resources, such as the shared head node where the WMS coordinator runs, the shared filesystem where intermediate data is stored, and the shared batch queue itself. We introduce Mufasa, a meta-workflow management system, which is designed to control the concurrency of multiple workflows in an ensemble, by observing and controlling the resources required by each WMS. We show some initial results demonstrating that Mufasa correctly handles the overcommitment of different resource types by starting, pausing, and cancelling workflows with unexpected behavior. Ben Lyons, Douglas Thain |
e-Science | 2 |
| 2022 | Dynamic Task Shaping for High Throughput Data Analysis Applications in High Energy PhysicsabstractDistributed data analysis frameworks are widely used for processing large datasets generated by instruments in scientific fields such as astronomy, genomics, and particle physics. Such frameworks partition petabyte-size datasets into chunks and execute many parallel tasks to search for common patterns, locate unusual signals, or compute aggregate properties. When well-configured, such frameworks make it easy to churn through large quantities of data on large clusters. However, configuring frameworks presents a challenge for end users, who must select a variety of parameters such as the blocking of the input data, the number of tasks, the resources allocated to each task, and the size of nodes on which they run. If poorly configured, the result may perform many orders of magnitude worse than optimal, or the application may even fail to make progress at all. Even if a good configuration is found through painstaking observations, the performance may change drastically when the input data or analysis kernel changes. This paper considers the problem of automatically configuring a data analysis application for high energy physics (TopEFT) built upon standard frameworks for physics analysis (Coffea) and distributed tasking (Work Queue). We observe the inherent variability within the application, demonstrate the problems of poor configuration, and then develop several techniques for automatically sizing tasks to meet goals of resource consumption, and overall application completion. Benjamín Tovar, Ben Lyons, Kelci Mohrman, Barry Sly-Delgado, Kevin Lannon, Douglas Thain |
IPDPS | 6 |
| 2021 | An Empirical Study of Package Dependencies and Lifetimes in Binder Python ContainersabstractContainers 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-Science | 3 |
| 2021 | Lightweight Function Monitors for Fine-Grained Management in Large Scale Python ApplicationsabstractPython 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 |
IPDPS | 9 |
| 2020 | Autoscaling High-Throughput Workloads on Container OrchestratorsabstractHigh-throughput computing (HTC) workloads seek to complete as many jobs as possible over a long period of time. Such workloads require efficient execution of many parallel jobs and can occupy a large number of resources for a long time. As a result, full utilization is the normal state of an HTC facility. The widespread use of container orchestrators eases the deployment of HTC frameworks across different platforms, which also provides an opportunity to scale up HTC workloads with almost infinite resources on the public cloud. However, the autoscaling mechanisms of container orchestrators are primarily designed to support latency-sensitive microservices, and result in unexpected behavior when presented with HTC workloads. In this paper, we design a feedback autoscaler, High Throughput Autoscaler (HTA), that leverages the unique characteristics of the HTC workload to autoscales the resource pools used by HTC workloads on container orchestrators. HTA takes into account a reference input, the real-time status of the jobs' queue, as well as two feedback inputs, resource consumption of jobs, and the resource initialization time of the container orchestrator. We implement HTA using the Makeflow workload manager, Work Queue job scheduler, and the Kubernetes cluster manager. We evaluate its performance on both CPU-bound and IO-bound workloads. The evaluation results show that, by using HTA, we improve resource utilization by 5.6× with a slight increase in execution time (about 15%) for a CPU-bound workload, and shorten the workload execution time by up to 3.65× for an IO-bound workload. Chao Zheng 0002, Nathaniel Kremer-Herman, Timothy Shaffer, Douglas Thain |
CLUSTER | 4 |
| 2020 | Solving the Container Explosion Problem for Distributed High Throughput ComputingabstractContainer technologies are seeing wider use at advanced computing facilities for managing highly complex applications that must execute at multiple sites. However, in a distributed high throughput computing setting, the unrestricted use of containers can result in the container explosion problem. If a new container image is generated for each variation of a job dispatched to a site, shared storage is soon exceeded. On the other hand, if a single large container image is used to meet multiple needs, the size of that container may become a problem for storage and transport. To address this problem, we observe that many containers have an internal structure generated by a structured package manager, and this information could be used to strategically combine and share container images. We develop Landlord to exploit this property and evaluate its performance through a combination of simulation studies and empirical measurement of high energy physics applications. Timothy Shaffer, Nicholas L. Hazekamp, Jakob Blomer, Douglas Thain |
IPDPS | 4 |
| 2019 | Dynamic Sizing of Continuously Divisible Jobs for Heterogeneous ResourcesabstractMany scientific applications operate on large datasets that can be partitioned and operated on concurrently. The existing approaches for concurrent execution generally rely on statically partitioned data. This static partitioning can lock performance in a sub-optimal configuration, leading to higher execution time and an inability to respond to dynamic resources. We present the Continuously Divisible Job abstraction which allows statically defined applications to have their component tasks dynamically sized responding to system behavior. The Continuously Divisible Job abstraction defines a simple interface that dictates how work can be recursively divided, executed, and merged. Implementing this abstraction allows scientific applications to leverage dynamic job coordinators for execution. We also propose the Virtual File abstraction which allows read-only subsets of large files to be treated as separate files. In exploring the Continuously Divisible Job abstraction, two applications were implemented using the Continuously Divisible Job interface: a bioinformatics application and a high-energy physics event analysis. These were tested using an abstract job interface and several job coordinators. Comparing these against a previous static partitioning implementation we show comparable or better performance without having to make static decisions or implement complex dynamic application handling. Nicholas L. Hazekamp, Benjamín Tovar, Douglas Thain |
eScience | 3 |
| 2018 | Wharf: Sharing Docker Images in a Distributed File SystemabstractContainer management frameworks, such as Docker, package diverse applications and their complex dependencies in self-contained images, which facilitates application deployment, distribution, and sharing. Currently, Docker employs a shared-nothing storage architecture, i.e. every Docker-enabled host requires its own copy of an image on local storage to create and run containers. This greatly inflates storage utilization, network load, and job completion times in the cluster. In this paper, we investigate the option of storing container images in and serving them from a distributed file system. By sharing images in a distributed storage layer, storage utilization can be reduced and redundant image retrievals from a Docker registry become unnecessary. We introduce Wharf, a middleware to transparently add distributed storage support to Docker. Wharf partitions Docker's runtime state into local and global parts and efficiently synchronizes accesses to the global state. By exploiting the layered structure of Docker images, Wharf minimizes the synchronization overhead. Our experiments show that compared to Docker on local storage, Wharf can speed up image retrievals by up to 12x, has more stable performance, and introduces only a minor overhead when accessing data on distributed storage. Chao Zheng 0002, Lukas Rupprecht, Vasily Tarasov, Douglas Thain, Mohamed Mohamed 0001, Dimitrios Skourtis, Amit Warke, Dean Hildebrand |
SoCC | 4 |
| 2018 | An Algebra for Robust Workflow TransformationsabstractScientific workflows are often designed with a particular compute site in mind. As a user changes sites the workflow needs to adjust. These changes include moving from a cluster to a cloud, updating an operating system, or investigating failures on a new cluster. As a workflow is moved, its tasks do not fundamentally change, but the steps to configure, execute and evaluate tasks differ. When handling these changes it may be necessary to use a script to analyze execution failure or run a container to use the correct operating system. To improve workflow portability and robustness, it is necessary to have a rigorous method that allows transformations on a workflow. These transformations do not change the tasks, only the way tasks are invoked. Using technologies such as containers, resource managers, and scripts to transform workflows allow for portability, but combining these technologies can lead to complications with execution and error handling. We define an algebra to reason about task transformations at the workflow level and express it in a declarative form using JSON. We implemented this algebra in the Makeflow workflow system and demonstrate how transformations can be used for resource monitoring, failure analysis, and software deployment across three sites. Nicholas L. Hazekamp, Douglas Thain |
eScience | 2 |
| 2018 | A First Look at the JX Workflow LanguageabstractScientific workflows are typically expressed as a graph of logical tasks, each one representing a single program along with its input and output files. This poster introduces JX (JSON eXtended), a declarative language that can express complex workloads as an assembly of sub-graphs that can be partitioned in flexible ways. We present a case study of using JX to represent complex workflows for the Lifemapper biodiversity project. We evaluate partitioning approaches across several computing environments, including ND-Condor, IU-Jetstream, and SDSC-Comet, and show that a coarse partitioning results in faster turnaround times, reduced data transfer, and lower master utilization across all three systems. Timothy Shaffer, Kyle M. D. Sweeney, Nathaniel Kremer-Herman, Douglas Thain |
eScience | 4 |
| 2018 | MAKER as a Service: Moving HPC Applications to Jetstream CloudabstractAs cloud resources become more available as an execution platform, the need to transition applications between HPC and the cloud becomes a necessity. However, because of the complex setup and system specific demands of these applications, transition is difficult and may not scale as desired. Jetstream is a NSF funded cloud service that is aiming to provide these services for users in a dynamical allocated nature. In this work we look at three key areas to focus on when transitioning between resources: providing a portable reproducible environment, scaling between local and remote resources, and using feedback to the user for informing configuration and runtime decisions. Building on the MAKER bioinformatic application, we have deployed WQ-MAKER on the Jetstream cloud platform, helping to annotate over 30 genomes and accelerating performance from days to hours and weeks to days. Nicholas L. Hazekamp, Upendra Kumar Devisetty, Nirav C. Merchant, Douglas Thain |
IC2E | 4 |
| 2018 | Automatic Dependency Management for Scientific Applications on ClustersabstractSoftware installation remains a challenge in scientific computing. End users require custom software stacks that are not provided through commodity channels. The resulting effort needed to install software delays research in the first place, and creates friction for moving applications to new resources and new users. Ideally, end-users should be able to manage their own software stacks without requiring administrator privileges. To that end, we describe vc3-builder, a tool for deploying software environments automatically on clusters. Its primary application comes in cloud and opportunistic computing, where deployment must be performed in batch as a side effect of job execution. vc3-builder uses workflow technologies as a means of exploiting cluster resources for building in a portable way. We demonstrate the use of vc3-builder on three applications with complex dependencies: MAKER, Octave, and CVMFS, building and running on three different cluster facilities in sequential, parallel, and distributed modes. Benjamín Tovar, Nicholas L. Hazekamp, Nathaniel Kremer-Herman, Douglas Thain |
IC2E | 4 |
| 2018 | A lightweight model for right-sizing master-worker applications
Nathaniel Kremer-Herman, Benjamín Tovar, Douglas Thain |
SC | 3 |
| 2018 | SHADHO: Massively Scalable Hardware-Aware Distributed Hyperparameter OptimizationabstractComputer vision is experiencing an AI renaissance, in which machine learning models are expediting important breakthroughs in academic research and commercial applications. Effectively training these models, however, is not trivial due in part to hyperparameters: user-configured values that control a model's ability to learn from data. Existing hyperparameter optimization methods are highly parallel but make no effort to balance the search across heterogeneous hardware or to prioritize searching high-impact spaces. In this paper, we introduce a framework for massively Scalable Hardware-Aware Distributed Hyperparameter Optimization (SHADHO). Our framework calculates the relative complexity of each search space and monitors performance on the learning task over all trials. These metrics are then used as heuristics to assign hyperparameters to distributed workers based on their hardware. We first demonstrate that our framework achieves double the throughput of a standard distributed hyperparameter optimization framework by optimizing SVM for MNIST using 150 distributed workers. We then conduct model search with SHADHO over the course of one week using 74 GPUs across two compute clusters to optimize U-Net for a cell segmentation task, discovering 515 models that achieve a lower validation loss than standard U-Net. Jeffery Kinnison, Nathaniel Kremer-Herman, Douglas Thain, Walter J. Scheirer |
WACV | 3 |
| 2018 | Combining Static and Dynamic Storage Management for Data Intensive Scientific WorkflowsabstractWorkflow management systems are widely used to express and execute highly parallel applications. For data-intensive workflows, storage can be the constraining resource: The number of tasks running at once must be artificially limited to not overflow the space available in the filesystem. It is all too easy for a user to dispatch a workflow which consumes all available storage and disrupts all system users. To address these issues, we present a three-tiered approach to workflow storage management: (1) A static analysis algorithm which analyzes the storage needs of a workflow before execution, giving a realistic prediction of success or failure. (2) An online storage management algorithm which accounts for the storage needed by future tasks to avoid deadlock at runtime. (3) A task containment system which limits storage consumption of individual tasks, enabling the strong guarantees of the static analysis and dynamic management algorithms. We demonstrate the application of these techniques on three complex workflows. Nicholas L. Hazekamp, Nathaniel Kremer-Herman, Benjamín Tovar, Haiyan Meng, Olivia Choudhury, Scott J. Emrich, Douglas Thain |
IEEE Trans. Parallel Distributed Syst. | 7 |
| 2018 | A Job Sizing Strategy for High-Throughput Scientific WorkflowsabstractThe user of a computing facility must make a critical decision when submitting jobs for execution: how many resources (such as cores, memory, and disk) should be requested for each job? If the request is too small, the job may fail due to resource exhaustion; if the request is too large, the job may succeed, but resources will be wasted. This decision is especially important when running hundreds of thousands of jobs in a high throughput workflow, which may exhibit complex, long tailed distributions of resource consumption. In this paper, we present a strategy for solving the job sizing problem: (1) applications are monitored and measured in user-space as they run; (2) the resource usage is collected into an online archive; and (3) jobs are automatically sized according to historical data in order to maximize throughput or minimize waste. We evaluate the solution analytically, and present case studies of applying the technique to high throughput physics and bioinformatics workflows consisting of hundreds of thousands of jobs, demonstrating an increase in throughput of 10-400 percent compared to naive approaches. Benjamín Tovar, Rafael Ferreira da Silva, Gideon Juve, Ewa Deelman, William E. Allcock, Douglas Thain, Miron Livny |
IEEE Trans. Parallel Distributed Syst. | 6 |
| 2017 | Deploying High Throughput Scientific Workflows on Container Schedulers with Makeflow and MesosabstractWorkflows are a widely used abstraction for describing large scientific applications and running them on distributed systems. However, most workflow systems have been silent on the question of what execution environment each task in the workflow is expected to run in. Consequently, a workflow may run successfully in the environment it was created, but fail on other platforms due to the differences in execution environment. Container-based schedulers have recently arisen as a potential solution to this problem, adopting containers to distribute computing resources and deliver well-defined execution environments to applications. In this paper, we consider how to connect workflow system to container schedulers with minimal performance loss and higher system efficiency. As an example of current technology, we use Makeflow and Mesos. We present five design challenges, and address them by using four configurations that connecting workflow system to container scheduler from different level of the infrastructure. In order to take full advantage of the resource sharing schema of Mesos, we enable the resource monitor of Makeflow to dynamically update the task resource requirement. We explore the performance of a large bioinformatics workflow, and observe that using Makeflow, Work Queue and the Resource monitor together not only increase the transfer throughput but also achieves highest resource usage rate. Chao Zheng 0002, Benjamín Tovar, Douglas Thain |
CCGrid | 3 |
| 2017 | Towards Scalable and Dynamic Social Sensing Using A Distributed Computing FrameworkabstractWith the rapid growth of online social media and ubiquitous Internet connectivity, social sensing has emerged as a new crowdsourcing application paradigm of collecting observations (often called claims) about the physical environment from humans or devices on their behalf. A fundamental problem in social sensing applications lies in effectively ascertaining the correctness of claims and the reliability of data sources without knowing either of them a priori, which is referred to as truth discovery. While significant progress has been made to solve the truth discovery problem, some important challenges have not been well addressed yet. First, existing truth discovery solutions did not fully solve the dynamic truth discovery problem where the ground truth of claims changes over time. Second, many current solutions are not scalable to large-scale social sensing events because of the centralized nature of their truth discovery algorithms. Third, the heterogeneity and unpredictability of the social sensing data traffic pose additional challenges to the resource allocation and system responsiveness. In this paper, we developed a Scalable Streaming Truth Discovery (SSTD) solution to address the above challenges. In this paper, we developed a Scalable Streaming Truth Discovery (SSTD) solution to address the above challenges. In particular, we first developed a dynamic truth discovery scheme based on Hidden Markov Models (HMM) to effectively infer the evolving truth of reported claims. We further developed a distributed framework to implement the dynamic truth discovery scheme using Work Queue in HTCondor system. We also integrated the SSTD scheme with an optimal workload allocation mechanism to dynamically allocate the resources (e.g., cores, memories) to the truth discovery tasks based on their computation requirements. We evaluated SSTD through real world social sensing applications using Twitter data feeds. The evaluation results on three real-world data traces (i.e., Boston Bombing, Paris Shooting and College Football) show that the SSTD scheme is scalable and outperforms the state-of-the-art truth discovery methods in terms of both effectiveness and efficiency. Daniel Yue Zhang, Chao Zheng 0002, Dong Wang 0002, Douglas Thain, Xin Mu, Gregory R. Madey, Chao Huang 0001 |
ICDCS | 4 |
| 2017 | Balancing push and pull in Confuga, an active storage cluster file system for scientific workflowsabstractSummary Most big‐data analysis systems require users to adopt restricted abstractions to achieve scaling and system stability. While highly effective at establishing data locality and eliminating interdependencies, this approach is not easily incorporated into scientific workflows that are often complex and irregular graphs of sequential programs with multiple dependencies. To address this, we have developed an active storage cluster file system named Confuga which harnesses the file information already available in the workflow to enable efficient and controlled distribution of dependencies across active storage nodes. Confuga is built upon the idea of leveraging a job's namespace to eliminate unknown transfers and to plan the replication of all job dependencies. Replication is carried out through two opposing transfer methodologies: centrally managed push transfers and distributed pulls. We evaluate the effectiveness of the two transfer mechanisms using workflows that stress the ability of the cluster to replicate dependencies. Ultimately, we show that a balance of the two approaches achieves optimal file distribution. This is shown in two bioinformatics workflows where a careful balance of the two mechanisms leads to 48% and 77% improvements over only push or pull. Copyright © 2016 John Wiley & Sons, Ltd. Patrick Donnelly, Douglas Thain |
Concurr. Comput. Pract. Exp. | 2 |
| 2017 | Report on the first workshop on negative and null results in eScienceabstractNew techniques and technologies, such as the use of large-scale computing, influence research approaches, methods, and scales and are rapidly changing the scientific landscape. Research projects in eScience ∗ thus start with many assumptions and many unknowns and are often complex. While the scientific process is sometimes viewed, at least in hindsight, as a linear progression from one good idea to the next, it is in fact fraught with false starts, wrong assumptions, and dead ends. The increasing reliance on computation adds to the scope of problems that occur. Researchers invest a significant amount of time and effort in their research. Funding agencies similarly make large investments to support such research, on the assumption that most of the research will be successful. When the research assumptions and hypotheses turn out to be false, causing results that are "negative" or "null", the natural bias is to judge that the research project "failed." The history of science, however, shows that negative results may be an opportunity to revolutionize a field of study. For example, Fleming noticed that his flu cultures were contaminated by mold, but that there was infection around that mold, leading to his discovery of Penicillin. Similarly, a project today may fail because of the misuse or failure of computational support. Such "failures" actually indicate that there is an opportunity for the cyberinfrastructure research community to improve computing resources and tools. The interaction of these modes of failure is multi-faceted. Negative results have been difficult to find in published papers in all scientific domains. We identify three reasons for this. First, negative results may not be identified as such but simply considered mistakes. Such cases may never be investigated further. Secondly, paper referees may demand a higher standard from such results, because they are more difficult to understand or challenge the conventional narrative. Third, researchers may self-select against publishing such results in light of the previous point. This paper contributes to the discussion about null or negative results in eScience. It also attempts to organize concepts about negative or null results in eScience in the form of a taxonomy. Falsifiability is the concept that a given statement can be refuted by a real-world measurement or observation 4. The empirical sciences are dominated by the construction of such statements and efforts to confirm or refute them. In eScience, such statements are rarely formally presented in a refutable manner. eScience projects typically merge goals from the physical science with computer science and engineering aspects. A failure in eScience may often be attributed to a computer engineering failure (software defects or unresolved performance shortcomings) or a collaboration misfit (the groups never came together). However, many important statements are never answered definitely, such as whether a given computational approach is effective for the physical science investigation. The formalization and confirmation/refutation of such statements have the potential to prevent efforts lost to engineering aspects. Post-mortem analysis of failed experiments provides "clues suggesting deeper lying forces," as Galison notes in How Experiments End, "Any historical reconstruction that ignores what seems in retrospect to be erroneous will be an inadequate account" 5. This means that the study of errors is not only relevant to students of history, as these "forces" can guide future investigations, suggest fundamental problems in experimental approaches, or even challenge prevailing theories. For example, relational database systems have been a well-accepted solution for information structuring, storage, and retrieval. This model is now being challenged by other database concepts, largely motivated by the need to cope with increasing data volumes. While the boundary between failure and success is not sharp in transitioning from relational to noSQL databases, the transition demonstrates a need to adapt and improve. This is often the normal path in research in computer science and cyberinfrastructure in particular, which could learn a lot from the various "failures" in eScience projects. Similarly, "short-term" examples of such "failures" are abound, such as the limit of being able to be a part of at most 16 Unix groups in NFS. It is likely that someone architecting an eScience collaboration system will face this limitation rather quickly. Are negative results as valuable as positive results in general? As Ayer points out, "What justifies scientific procedure ... is the success of the predictions to which it gives rise" 6. Following this line of thought, negative results are subordinate to the positive results that validate useful predictions. A negative result invalidates a previously held prediction, challenging or demolishing a theory or model. It is, however, incomplete. A negative result is an opportunity to pick up the pieces and fix the theory. Negative results are thus an important reminder of the limitations of science at any given point in time. "We forget about unpredictability when it is our turn to predict," Taleb says in The Black Swan 7, a book that attempts to analyze tumultuous events, including several scientific cases. Taleb makes the case that studying such cognitive upheavals is worthwhile in its own right. Professionals who act with the history of failed ideas and efforts in mind will be more resilient against similar changes in the future. Taleb describes the social and mental impact of experiencing (repeated) failure, indicating that without support, researchers can easily become demoralized and shy away from challenging, long-term problems. However, he notes that "Your finding nothing is very valuable... —hey, you know where not to look" 7. Venues such as the ERROR workshop are intended to encourage discussion of specific negative results. By co-locating with the eScience conference in Munich, the workshop attracted significant attention from this scientific community, with about 20 participants. The workshop accepted four papers out of six submitted after a peer review process. Each paper was reviewed by two to three members of the program committee, who evaluated the works based on originality, scientific rigor, significance, and presentation. The accepted papers were presented orally at the workshop, followed by a panel discussion on the topic "theory versus practice in eScience: gaps and gaping holes." The first presentation, by Gomes et al. 9, considered problems regarding interoperability between scientific workflows. Specifically, they discussed the problem of reusing workflows previously developed and implemented using one particular scientific workflow management framework with another one. To solve this problem, the authors developed an "intermediate" workflow language, with the idea that this intermediate language would preserve the workflow's semantic information across frameworks. However, they observed a loss of information about workflow semantics during the translation from a first workflow language to the "neutral" language and from the "neutral" language to the second workflow language. This happens because there are no ideal or standardized semantics for workflow languages, which is the key negative result in this research. A solution proposed to this problem is through the adoption of workflow patterns to describe richer workflow semantics. The second presentation, by Groen and Portgies Zwart 10, provided a high level overview of the authors' experience in constructing a distributed supercomputing system: CosmoGrid. The authors discussed how ambitious ideas can often be stymied by site-local resource allocation decisions. One of the negative results is the conclusion that harnessing multiple large machines is not feasible and therefore one should focus on harnessing a larger number of smaller machines. Additionally, the authors pointed out that a task as simple as getting software installed is significantly difficult at major computational sites, exposing the often overlooked reality of working with large scale computational infrastructures. The third presentation, by Cebrian et al. 11, presented an experience in designing two separate cache stores—for private and shared data— for multicore system architectures. The premise of the work was to improve efficiency by excluding private and shared read-only cache contents from coherency management. From the experiments and analysis of results obtained with this approach, the authors concluded that systems are less efficient with this kind of design, which is a negative result. This is because the overhead of classification mechanisms and increased concentration of access to shared data cause a bandwidth bottleneck to a particular portion of cache, resulting in higher latencies. The fourth presentation, by Jackson et al. 12, discussed an experience of performing an experiment in the context of a larger body of work. The experiment was on latency measurement between nodes within a single cluster, as well as across different clusters. Some of the main takeaways from the experiments as described by the authors were the technical and administrative obstacles faced when the experiment involves dependencies on several independently managed computational systems across administrative boundaries. Despite these obstacles, the authors were able to produce a significant body of latency data. One finding from this dataset was that an exhaustive study of latencies among systems was not necessarily, by itself, a good predictor of actual application performance. The topic of discussion for the panel session was "theory versus practice in eScience: gaps and gaping holes." This theme emerged from a recurring observation in the submitted papers, in which negative or null results are attributed to a mismatch between expectations, which are based on theory, and what is found in reality. The panelists were Daniel S. Katz, Simon Portegies Zwart, Kyle Chard, Juan M. Cebrián, and Gary Jackson. Each panelist spoke for 2–3 min, and then there was an open discussion between panelists and the audience. The rest of this subsection presents the highlights of the panelists talks and the discussion that followed. Jackson spoke about the importance of not losing research focus because of infrastructure complexities and problems—the proverbial "missing the forest for trees." Chard said that there are no good definitions of eScience, although we provide one taken from the eScience conference series website in Section 2. Furthermore, he raised questions as to how the scientific process, which had been mostly unchanged for hundreds of years and has long review cycles, has recently been changing with online data publication and open access journals. Katz spoke of the phenomenon of failures among research projects and endeavors by quoting from the opening of Tolstoy's Anna Karenina, "Happy families are all alike; every unhappy family is unhappy in its own way." He noted that successful research endeavors must have all their critical factors right to be successful and failing even one of them could jeopardize the complete project. This is popularly known as Anna Karenina Principle 13. Katz also emphasized the importance and value of scientific results in general and negative results in particular with the question/statement: "How do we decide if there is value in a result?" Portgies Zwart spoke about the importance of understanding the difference between core computer science and other sciences, as well as the scientists associated with each one. He argued that computer science is currently undergoing a crisis because it is hard to find interesting problems, because of competition with the industry. In particular, computer science is challenged by reproducibility. One solution, he suggested, is an establishment of a software museum to prevent loss of software. The open discussion that followed focused on diverse topics such as software preservation, publication and credit, training, and the definition of negative or null results. It began with participants expressing concerns about issues related to software in particular. Scenarios were discussed that introduce the "gaps between eScience theory and practice" connecting technologies, ideas, and people. One gap is that digital products in general and software in particular, including methods and knowledge (algorithms), can be lost over time, sometimes known as bit rot 14. Software hosting services such as GitHub can address this problem to a certain extent by preserving the files, but they still require much human effort to preserve the function delivered by the software as meaningful. A curation service for algorithms could be another solution, requiring additional effort. Can the software and algorithm hosting services be linked as concepts and implementation? Preservation and publication of negative results are a challenge. While there are no technical barriers, from a publishing culture point of view, there are few or no incentives for publishing negative results. In the presence of such incentives, people would develop the culture about explaining not only what they did but also why they did not do so in some other way. And, if the negative results are actually published, it is likely that the same approach will not be taken by other researchers and groups. Another identified gap is the lack of a comprehensive understanding of negative results because of the lack of a conceptual framework, for example, a taxonomy. The discussion also raised the gap introduced by a lack of a credit model for discovering, identifying, and reporting negative results. For example, can negative results and/or methods to obtain them be patented? For instance, who receives credit if a succession of graduate students working on a problem arrive at a negative result followed by positive result? Or what happens to the positive results obtained before further investigation leads to their negation and nullification? Another gap arises from the lack of proper eScience training of domain scientists—can domain scientists be trained to become eScientists? This also applies to principal investigators, many of whom were trained in an era when science was carried out differently than it is carried out today. Training imparting the knowledge of modern computational methods and capabilities could play a key role in filling such a gap. The difference between incomplete (such as obtained from samples of insufficient size) and negative results can also be unclear. This can often result in negative results that are subject to interpretation. For instance, in an MD simulation 14, it cannot be shown if sampling was sufficient. In the same vein, should incremental competitive results be considered negative? In general, there is no standard on how many simulated timesteps are needed to obtain the correct answer, although sometimes, one can validate against lab experiments. Another topic raised during the discussion compared research in academia versus science in the commercial sector. One prominent sentiment expressed in the discussion was that, in some areas, research done in the commercial domain is "ahead" of research in academia, particularly, where industry has larger-scale problems and data than academia. One possible reason for this could be that there are more negative and null results in academia compared to industry. However, academic research can be transferred to the commercial domain and vice versa. For instance, the patent system is in place to enable commercial contribution to public research. One question that arises here is should there be a distinction between commercial and academic research? How will this distinction manifest itself? It was noted that there is an asymmetry between positive and negative results: With negative results, it is more likely that some error was made. It may be harder to truly prove a negative result. For instance an "existence proof" is sufficient for a positive result while a "non-possible proof" is needed for negative results. The workshop led to concrete outcomes before, during, and as a follow-up of its realization. After the announcement of the workshop, the Mozilla Science Foundation hosted a guest post about the workshop by the organizers 15. The post discussed the importance of the theme of "negative" results and the goals for the workshop. One of the outcomes of the panel discussion was the call for a taxonomy of negative and null results in eScience. We respond to this call in this paper by proposing a taxonomy in Section 5. In Figure 1, we present a taxonomy of eScience results. The three kinds of results in eScience are positive, null, and negative. In the taxonomy, we focus on negative and null results. Negative results may be caused by one or more of the following reasons: technological, technical, human, and domain. For instance, an erroneous result obtained because of insufficient precision resulting from a limitation of a system library is an example of negative result caused by technical and technological limitations. Similarly, a simulation algorithm resulting from a flawed understanding of a natural phenomenon could be considered a negative result caused by domain and human factors. A mismatch between the problem/data size and the technology/methodology used is an example of a technological reason. Examples of technical causes include software bugs and cyberinfrastructure faults. Human causes include both incidental issues such as mistakes in measurements and systemic issues such as a false hypothesis or an insufficient sample size. Null results are obtained because of the lack of discriminating conditions to confirm or refute a hypothesis. Such situation may be caused by similar reasons as for negative results. An example of technological/technical reason is a statistical test has poor performance on the data because the implementation uses limited precision. Null results can also have a human cause, when insufficient samples are used in the experiments or when some bias in the data goes unnoticed. We are aware of two workshops with similar themes in related fields (Information and Communication Technologies). The first is NoISE (Workshop on Negative or Inconclusive Results in Semantic Web) 19. The second is NOPE (Workshop on Negative Outcomes, Post-mortems, and Experiences) 20. Both the workshops were organized for the first time in 2015. Similarly to the current special issue, there have been two special issues in the prominent journals focused on negative or null results: the Journal on Negative results in Empirical Software Engineering 21 and PLOS ONE Collection 22. The PLOS ONE Collection focuses on inconclusive results as a distinct type of negative results in addition to null results. Similarly to the workshops, both special issues were launched for the first time in 2015. In this section, we discuss some of the key implications that are drawn from the previous sections. These are the issues that are directly impacted by the occurrence of negative and null results in the research as conducted by the scientific community. We classify these implications into four categories: technological, technical, cultural, and domain specific. Publication, credit, and citation of the work that has yielded negative results are an important consideration from the research community point of view. Citations and credit are important measures of success for a research publication. Given the current trends of publishing positive results, it is a crucial decision for a researcher to invest efforts in publishing a negative result. Technical issues such as hardware faults and software bugs often go undetected until late in the research work. In these cases, the negative results are not necessarily of the same nature as the science domain unless the domain is computer science itself. It becomes difficult for a domain scientist to draw value from the publication and dissemination of such results. As a consequence, they often are ignored or fixed after the results were obtained and disseminated. For example, a bug in the third party library call up the toolchain of an application that limited the results of the actual science in scale or precision can be considered a negative result. In some cases, the problems are mismatched to the available computational infrastructure. Sometimes, the problems are too small for a given environment, leading to an inefficient use. In others cases, the problems are too large for the infrastructure, leading to generation of incomplete or no results. Publishing details about such cases can benefit the community by allowing it to better understand how to more optimally match problems and solutions. Finally, a cumulative effect from more than one cause is possible. One of the biggest concern about such issues is that they often go unnoticed by the larger community and hence the appropriate correction measures are not adapted. We are grateful to the program committee members, the panelists, and the authors who supported the realization of this first workshop. We also thank the reviewers of this special issue. The work by Katz was supported in part by the National Science Foundation while working at the Foundation. Any opinion, finding, and conclusions or recommendations expressed in this material are those of the author(s) and do not necessarily reflect the views of the National Science Foundation. Ketan Maheshwari, Daniel S. Katz, Sílvia Delgado Olabarriaga, Justin M. Wozniak, Douglas Thain |
Concurr. Comput. Pract. Exp. | 5 |
| 2017 | Designing Self-Tuning Split-Map-Merge Applications for High Cost-Efficiency in the CloudabstractCloud platforms are attractive for executing large concurrent applications that require access to a pool of resources for concurrently executing the partitions of their workloads. Historically, application designers have tuned concurrent applications for specific hardware and platforms. But such approaches are not viable in cloud platforms as applications can be deployed on a variety of platforms and the operating environments can vary in each deployment. In this work, we argue and demonstrate that concurrent applications in cloud platforms must be self-tuning. First, we show that applications must incorporate a model of the overheads of operation. Second, we show that applications must determine their resource requirements and tune their operation to the operating conditions using estimations from the model. We build two self-tuning applications, E-Sort and E-MAKER, and demonstrate their ability to achieve high cost-efficiency by determining the right scale of partitions and resources to use for operation and adapting their behavior according to the characteristics of the deployed environment. Dinesh Rajan, Douglas Thain |
IEEE Trans. Cloud Comput. | 2 |
| 2016 | PRUNE: A preserving run environment for reproducible scientific computingabstractComputing as a whole suffers from a crisis of reproducibility. Programs executed in one context are astonishingly hard to reproduce in another context, resulting in wasted effort by people and general distrust of results produced by computer. The root of the problem lies in the fact that every program has implicit dependencies on data and execution environment which are rarely understood by the end user. To address this problem, we present PRUNE, the Preserving Run Environment. In PRUNE, every task to be executed is wrapped in a functional interface and coupled with a strictly defined environment. The task is then executed by PRUNE rather than the user to ensure reproducibility. As a scientific workflow evolves in PRUNE, a growing but immutable tree of derived data is created. The provenance of every item in the system can be precisely described, facilitating sharing and modification between collaborating researchers, along with efficient management of limited storage space. We present the user interface and the initial prototype of PRUNE, and demonstrate its application in matching records and comparing surnames in U.S. Censuses. Peter Ivie, Douglas Thain |
eScience | 2 |
| 2016 | Conducting reproducible research with Umbrella: Tracking, creating, and preserving execution environmentsabstractPublishing scientific results without the detailed execution environments describing how the results were collected makes it difficult or even impossible for the reader to reproduce the work. However, the configurations of the execution environments are too complex to be described easily by authors. To solve this problem, we propose a framework facilitating the conduct of reproducible research by tracking, creating, and preserving the comprehensive execution environments with Umbrella. The framework includes a lightweight, persistent and deployable execution environment specification, an execution engine which creates the specified execution environments, and an archiver which archives an execution environment into persistent storage services like Amazon S3 and Open Science Framework (OSF). The execution engine utilizes sandbox techniques like virtual machines (VMs), Linux containers and user-space tracers, to create an execution environment, and allows common dependencies like base OS images to be shared by sandboxes for different applications. We evaluate our framework by utilizing it to reproduce three scientific applications from epidemiology, scene rendering, and high energy physics. We evaluate the time and space overhead of reproducing these applications, and the effectiveness of the chosen archive unit and mounting mechanism for allowing different applications to share dependencies. Our results show that these applications can be reproduced using different sandbox techniques successfully and efficiently, even through the overhead and performance slightly vary. Haiyan Meng, Douglas Thain, Alexander Vyushkov, Matthias Wolf 0003, Anna Woodard |
eScience | 2 |
| 2016 | DistIA: a cost-effective dynamic impact analysis for distributed programsabstractDynamic impact analysis is a fundamental technique for understanding the impact of specific program entities, or changes to them, on the rest of the program for concrete executions. However, existing techniques are either inapplicable or of very limited utility for distributed programs running in multiple concurrent processes. This paper presents DistIA, a dynamic analysis of distributed systems that predicts impacts propagated both within and across process boundaries by partially ordering distributed method-execution events, inferring causality from the ordered events, and exploiting message-passing semantics. We applied DistIA to large distributed systems of various architectures and sizes, for which it on average finishes the entire analysis within one minute and safely reduces impact-set sizes by over 43% relative to existing options with runtime overhead less than 8%. Moreover, two case studies initially demonstrate the precision of DistIA and its utility in distributed system understanding. While conservative thus subject to false positives, DistIA balances precision and efficiency to offer cost-effective options for evolving distributed programs. Haipeng Cai, Douglas Thain |
ASE | 2 |
| 2016 | DiaPro: Unifying Dynamic Impact Analyses for Improved and Variable Cost-EffectivenessabstractImpact analysis not only assists developers with change planning and management, but also facilitates a range of other client analyses, such as testing and debugging. In particular, for developers working in the context of specific program executions, dynamic impact analysis is usually more desirable than static approaches, as it produces more manageable and relevant results with respect to those concrete executions. However, existing techniques for this analysis mostly lie on two extremes: either fast, but too imprecise, or more precise, yet overly expensive. In practice, both more cost-effective techniques and variable cost-effectiveness trade-offs are in demand to fit a variety of usage scenarios and budgets of impact analysis. This article aims to fill the gap between these two extremes with an array of cost-effective analyses and, more broadly, to explore the cost and effectiveness dimensions in the design space of impact analysis. We present the development and evaluation of D ia P ro , a framework that unifies a series of impact analyses, including three new hybrid techniques that combine static and dynamic analyses. Harnessing both static dependencies and multiple forms of dynamic data including method-execution events, statement coverage, and dynamic points-to sets, D ia P ro prunes false-positive impacts with varying strength for variant effectiveness and overheads. The framework also facilitates an in-depth examination of the effects of various program information on the cost-effectiveness of impact analysis. We applied D ia P ro to ten Java applications in diverse scales and domains, evaluating it thoroughly on both arbitrary and repository-based queries from those applications. We show that the three new analyses are all significantly more effective than existing alternatives while remaining efficient, and the D ia P ro framework, as a whole, provides flexible cost-effectiveness choices for impact analysis with the best options for variable needs and budgets. Our study results also suggest that hybrid techniques tend to be much more cost-effective than purely dynamic approaches, in general, and that statement coverage has mostly stronger effects than dynamic points-to sets on the cost-effectiveness of dynamic impact analysis, while static dependencies have even stronger effects than both forms of dynamic data. Haipeng Cai, Raúl A. Santelices, Douglas Thain |
ACM Trans. Softw. Eng. Methodol. | 3 |
| 2015 | Confuga: Scalable Data Intensive Computing for POSIX WorkflowsabstractToday's big-data analysis systems achieve performance and scalability by requiring end users to embrace a novel programming model. This approach is highly effective whose the objective is to compute relatively simple functions on colossal amounts of data, but it is not a good match for a scientific computing environment which depends on complex applications written for the conventional POSIX environment. To address this gap, we introduce Conjugal, a scalable data-intensive computing system that is largely compatible with the POSIX environment. Conjugal brings together the workflow model of scientific computing with the storage architecture of other big data systems. Conjugal accepts large workflows of standard POSIX applications arranged into graphs, and then executes them in a cluster, exploiting both parallelism and data-locality. By making use of the workload structure, Conjugal is able to avoid the long-standing problems of metadata scalability and load instability found in many large scale computing and storage systems. We show that CompUSA's approach to load control offers improvements of up to 228% in cluster network utilization and 23% reductions in workflow execution time. Patrick Donnelly, Nicholas L. Hazekamp, Douglas Thain |
CCGRID | 3 |
| 2015 | Balancing Thread-Level and Task-Level Parallelism for Data-Intensive Workloads on Clusters and CloudsabstractThe runtime configuration of parallel and distributed applications remains a mysterious art. To tune an application on a particular system, the end-user must choose the number of machines, the number of cores per task, the data partitioning strategy, and so on, all of which result in a combinatorial explosion of choices. While one might try to exhaustively evaluate all choices in search of the optimal, the end user's goal is simply to run the application once with reasonable performance by avoiding terrible configurations. To address this problem, we present a hybrid technique based on regression models for tuning data intensive bioinformatics applications: the sequential computational kernel is characterized empirically and then incorporated into an ab initio model of the distributed system. We demonstrate this technique on the commonly-used applications BWA, Bowtie2, and BLASR and validate the accuracy of our proposed models on clouds and clusters. Olivia Choudhury, Dinesh Rajan, Nicholas L. Hazekamp, Sandra Gesing, Douglas Thain, Scott J. Emrich |
CLUSTER | 5 |
| 2015 | Practical Resource Monitoring for Robust High Throughput ComputingabstractRobust high throughput computing requires effective monitoring and enforcement of a variety of resources including CPU cores, memory, disk, and network traffic. Without effective monitoring and enforcement, it is easy to overload machines, causing failures and slowdowns, or underutilize machines, which results in wasted opportunities. This paper explores how to describe, measure, and enforce resources used by computational tasks. We focus on tasks running in distributed execution systems, in which a task requests the resources it needs, and the execution system ensures the availability of such resources. This presents two non-trivial problems: how to measure the resources consumed by a task, and how to monitor and report resource exhaustion in a robust and timely manner. For both of these tasks, operating systems have a variety of mechanisms with different degrees of availability, accuracy, overhead, and intrusiveness. We describe various forms of monitoring and the available mechanisms in contemporary operating systems. We then present two specific monitoring tools that choose different tradeoffs in overhead and accuracy, and evaluate them on a selection of benchmarks. Gideon Juve, Benjamín Tovar, Rafael Ferreira da Silva, Dariusz Król 0002, Douglas Thain, Ewa Deelman, William E. Allcock, Miron Livny |
CLUSTER | 5 |
| 2015 | Scaling Data Intensive Physics Applications to 10k Cores on Non-dedicated Clusters with LobsterabstractThe high energy physics (HEP) community relies upon a global network of computing and data centers to analyze data produced by multiple experiments at the Large Hadron Collider (LHC). However, this global network does not satisfy all research needs. Ambitious researchers often wish to harness computing resources that are not integrated into the global network, including private clusters, commercial clouds, and other production grids. To enable these use cases, we have constructed Lobster, a system for deploying data intensive high throughput applications on non-dedicated clusters. This requires solving multiple problems related to non-dedicated resources, including work decomposition, software delivery, concurrency management, data access, data merging, and performance troubleshooting. With these techniques, we demonstrate Lobster running effectively on 10k cores, producing throughput at a level comparable with some of the largest dedicated clusters in the LHC infrastructure. Anna Woodard, Matthias Wolf 0003, Charles Müller, Nil Valls, Benjamín Tovar, Patrick Donnelly, Peter Ivie, Kenyi Hurtado Anampa, Paul R. Brenner, Douglas Thain, Kevin Lannon, Michael D. Hildreth |
CLUSTER | 10 |
| 2015 | Scaling Up Bioinformatics Workflows with Dynamic Job Expansion: A Case Study Using Galaxy and MakeflowabstractLogical workflow management systems provide a user-friendly portal through which data can be processed using a sequence of standard tools. These logical workflows are a natural way to express the high level intent of the user, and to share the structure and the results with other users. However, logical workflows are not necessarily suited to expressing parallelism for very large runs. As the amount of data is scaled up, the run time of each node in the logical workflow may become extreme. We propose a technique of job expansion to solve this problem. When job expansion is applied to a logical workflow, each node in the workflow is itself expanded into a large performance workflow that may consist of hundreds to thousands of tasks that can be executed in parallel, thus enabling high concurrency and scalability. From the user's perspective, nothing has changed and the logical workflow remains in its original form. To demonstrate this technique, we have applied job expansion to a selection of bioinformatics applications running in the Galaxy workflow management system. Each job in the workflow is expanded into a highly parallel workflow executed using Makeflow, which is well suited to express high levels of parallelism. Work Queue is then utilized for execution because of its ability to quickly dispatch tasks and cache files for later reuse. After applying job expansion, we improve the execution time of BWA 18X and GATK 402X, with a total speedup of 61.5X on the workflow. We also take a look at the systems behavior since its launch to analyze its effectiveness. Nicholas L. Hazekamp, Joseph Sarro, Olivia Choudhury, Sandra Gesing, Scott J. Emrich, Douglas Thain |
e-Science | 6 |
| 2014 | Accelerating Comparative Genomics Workflows in a Distributed Environment with Optimized Data PartitioningabstractThe advent of new sequencing technology has generated massive amounts of biological data at unprecedented rates. High-throughput bioinformatics tools are required to keep pace with this. Here, we implement a workflow-based model for parallelizing the data intensive task of genome alignment and variant calling with BWA and GATK's Haplotype Caller. We explore different approaches of partitioning data and how each affect the run time. We observe granularity-based partitioning for BWA and alignment-based partitioning for Halo type Caller to be the optimal choices for the pipeline. We identify the various challenges encountered while developing such an application and provide an insight into addressing them. We report significant performance improvements, from 12 days to 4 hours, while running the BWA-GATK pipeline using 100 nodes for analyzing high-coverage oak tree data. Olivia Choudhury, Nicholas L. Hazekamp, Douglas Thain, Scott J. Emrich |
CCGRID | 3 |
| 2014 | Expanding Tasks of Logical Workflows Into Independent Workflows for Improved ScalabilityabstractWorkflow Management Systems, such as Galaxy and Taverna, provide a portal through which data can be processed using a sequence of different tools. This sequence allows for the creation of a logical workflow that describes the process. However, when the data workload becomes large enough the time spent in each logical step increases making it difficult to run the workflow fast and efficiently. The proposed solutions is to use task level expansion. Task expansion aims to take each step of the logical workflow and expand it into a new self-contained workflow. These workflows would allow for greater scalability and concurrency by creating more tasks. The resulting workflows will be used indistinguishably from the original tool, but perform more quickly and efficiently. The concept was applied to the BWA tool in Galaxy and we were able to see a 7.36 times speedup in runtime on our 32 GB dataset. Nicholas L. Hazekamp, Olivia Choudhury, Sandra Gesing, Scott J. Emrich, Douglas Thain |
CCGRID | 5 |
| 2014 | Opportunistic High Energy Physics Computing in User Space with ParrotabstractThe computing needs of high energy physics experiments like the Compact Muon Solenoid experiment at the Large Hadron Collider currently exceed the available dedicated computational resources, hence motivating a push to leverage opportunistic resources. However, access to opportunistic resources faces many obstacles, not the least of which is making available the complex software stack typically associated with such computations. This paper describes a framework constructed using existing software packages to distribute the needed software to opportunistic resources without the need for the job to have root-level privileges. Preliminary tests with this framework have demonstrated the feasibility of the approach and identified bottlenecks as well as reliability issues which must be resolved in order to make this approach viable for broad use. Dillon Skeehan, Paul R. Brenner, Benjamín Tovar, Douglas Thain, Nil Valls, Anna Woodard, Matthias Wolf 0003, T. Pearson, S. Lynch, Kevin Lannon |
CCGRID | 4 |
| 2014 | Adapting bioinformatics applications for heterogeneous systems: a case studyabstractSUMMARY The advent of new sequencing technologies has generated extremely large amounts of information. To successfully apply bioinformatics tools to such large datasets, they need to exhibit scalability and ideally elasticity in diverse computing environments. We describe the application of previously obtained lessons to a new workflow with and without shared file storage. Because the original workflows have an intractable sequential running times on large datasets, we propose lessons and results for refactoring bioinformatics tools for elastic scaling on personal clouds. Our case studies describe the various challenges faced when constructing such a workflow, from dealing with failure detection, to managing dependencies, to handling the quirks of the underlying operating systems. The practice of scaling bioinformatics tools is increasingly commonplace. As such, this hands‐on application of refactoring techniques can serve as a valuable guide. Significantly, our customized Makeflow framework enabled generalizable deployment on a wider variety of systems while substantially reducing wall clock runtimes using hundreds of cores. Copyright © 2012 John Wiley & Sons, Ltd. Irena Lanc, Peter Bui, Douglas Thain, Scott J. Emrich |
Concurr. Comput. Pract. Exp. | 3 |
| 2013 | Case Studies in Designing Elastic ApplicationsabstractClusters, clouds, and grids offer access to large scale computational resources at low cost. This is especially appealing to scientific applications that require a very large scale to compete in the research space. However, the resources available across these platforms differ significantly in their availability, hardware, environment, performance, cost of use, and more. This requires the use of elastic applications that can adapt to the resources available at run-time, transparently handling heterogeneity and failures. In this paper, we present case studies of several elastic applications built using the Work Queue programming framework. From this experience, we offer six general guidelines for the design and implementation of elastic applications that run on thousands of processors. Dinesh Rajan, Andrew Thrasher, Badi Abdul-Wahid, Jesús A. Izaguirre, Scott J. Emrich, Douglas Thain |
CCGRID | 6 |
| 2013 | Making work queue cluster-friendly for data intensive scientific applicationsabstractResearchers with large-scale data-intensive applications often wish to scale up applications to run on multiple clusters, employing a middleware layer for resource management across clusters. However, at the very largest scales, such middleware is often “unfriendly” to individual clusters, which are usually designed to support communication within the cluster, not outside of it. To address this problem we have modified the Work Queue master-worker application framework to support a hierarchical configuration that more closely matches the physical architecture of existing clusters. Using a synthetic application we explore the properties of the system and evaluate its performance under multiple configurations, with varying worker reliability, network capabilities, and data requirements. We show that by matching the software and hardware architectures more closely we can gain both a modest improvement in runtime and a dramatic reduction in network footprint at the master. We then run a scalable molecular dynamics application (AWE) to examine the impact of hierarchy on performance, cost and efficiency for real scientific applications and see a 96% reduction in network footprint, making it much more palatable to system operators and opening the possibility of increasing the application scale by another order of magnitude or more. Michael Albrecht, Dinesh Rajan, Douglas Thain |
CLUSTER | 3 |
| 2012 | Fine-Grained Access Control in the Chirp Distributed File SystemabstractAlthough the distributed file system is a widely used technology in local area networks, it has seen less use on the wide area networks that connect clusters, clouds, and grids. One reason for this is access control: existing file system technologies require either the client machine to be fully trusted, or the client process to hold a high value user credential, neither of which is practical in large scale systems. To address this problem, we have designed a system for fine-grained access control which dramatically reduces the amount of trust required of a batch job accessing a distributed file system. We have implemented this system in the context of the Chirp user-level distributed file system used in clusters, clouds, and grids, but the concepts can be applied to almost any other storage system. The system is evaluated to show that performance and scalability are similar to other authentication methods. The paper concludes with a discussion of integrating the authentication system into workflow systems. Patrick Donnelly, Douglas Thain |
CCGRID | 2 |
| 2012 | Resource Management for Elastic Cloud WorkflowsabstractCloud computing systems have joined campus and private grids as powerful and highly scalable environments for scientific computing. Furthermore, distributed applications are typically expressed in a form that allows them to run on an arbitrary number of nodes while tolerating failures and changes in available resources. This flexibility introduces problems relating to how many nodes an application can use, and how they should be allocated. In this paper, we explore these problems by presenting a general purpose architecture for scalable cloud applications, and describe inherent resource management problems. We address these challenges by developing methods for runtime measurement of the number of nodes an application can use, for appropriately placing masters and workers, and for matching workers to masters. Finally, we propose a resource management mechanism that allows automatic resource allocation and flexible resource distribution. These techniques are presented in the context of our specific cloud architecture, but the lessons apply to any system where competing elastic applications must be right-sized to the available resources. Douglas Thain |
CCGRID | 2 |
| 2012 | Folding proteins at 500 ns/hour with Work QueueabstractMolecular modeling is a field that traditionally has large computational costs. Until recently, most simulation techniques relied on long trajectories, which inherently have poor scalability. A new class of methods is proposed that requires only a large number of short calculations, and for which minimal communication between computer nodes is required. We considered one of the more accurate variants called Accelerated Weighted Ensemble Dynamics (AWE) and for which distributed computing can be made efficient. We implemented AWE using the Work Queue framework for task management and applied it to an all atom protein model (Fip35 WW domain). We can run with excellent scalability by simultaneously utilizing heterogeneous resources from multiple computing platforms such as clouds (Amazon EC2, Microsoft Azure), dedicated clusters, grids, on multiple architectures (CPU/GPU, 32/64bit), and in a dynamic environment in which processes are regularly added or removed from the pool. This has allowed us to achieve an aggregate sampling rate of over 500 ns/hour. As a comparison, a single process typically achieves 0.1 ns/hour. Badi Abdul-Wahid, Dinesh Rajan, Haoyun Feng, Eric Darve, Douglas Thain, Jesús A. Izaguirre |
eScience | 6 |
| 2012 | A system for management of Computational Fluid Dynamics simulations for civil engineeringabstractWe introduce a web-based system for management of Computational Fluid Dynamics(CFD) simulations. This system provides an interface for users, on a web-browser, to have an intuitive, user-friendly means of dispatching and controlling long-running simulations. CFD presents a challenge to its users due to the complexity of its internal mathematics, the high computational demands of its simulations and the complexity of inputs to its simulations and related tasks. We designed this system to be as extensible as possible in order to be suitable for many different civil engineering applications. The front-end of this system is a webserver, which provides the user interface. The back-end is responsible for starting and stopping jobs as requested. There are also numerous components specifically for facilitating CFD computation. We discuss our experience with presenting this system to real users and the future ambitions for this project. Peter Sempolinski, Douglas Thain, Daniel Wei, Ahsan Kareem |
eScience | 2 |
| 2012 | Scripting distributed scientific workflows using WeaverabstractSUMMARY Weaver is a high‐level distributed computing framework that enables researchers to construct scalable scientific data‐processing workflows. Instead of developing a new workflow language, we introduce a domain‐specific language built on top of Python called Weaver, which takes advantage of users' familiarity with the programming language, minimizes barriers to adoption, and allows for integration with a rich ecosystem of existing software. In this paper, we provide an overview of Weaver's programming model, which allows users to organize and specify scientific workflows by using a collection of datasets, functions, and abstractions. We also explain how these workflow specifications are compiled into a directed acyclic graph that is used by the Makeflow workflow manager to dispatch work to a variety of distributed execution platforms. To demonstrate the power and benefits of using the framework in constructing scientific research applications, the paper examines four distinct real‐world applications scripted using Weaver and analyzes the performance, scalability, and impact of the distributed generated scientific workflows. Copyright © 2011 John Wiley & Sons, Ltd. Peter Bui, Andrew Thrasher, Rory Carmichael, Irena Lanc, Patrick Donnelly, Douglas Thain |
Concurr. Comput. Pract. Exp. | 7 |
| 2012 | ROARS: a robust object archival system for data intensive scientific computing
Hoang Bui, Peter Bui, Patrick J. Flynn, Douglas Thain |
Distributed Parallel Databases | 4 |
| 2012 | A Framework for Scalable Genome Assembly on Clusters, Clouds, and GridsabstractBioinformatics researchers need efficient means to process large collections of genomic sequence data. One application of interest, genome assembly, has great potential for parallelization; however, most previous attempts at parallelization require uncommon high-end hardware. This paper introduces the Scalable Assembler at Notre Dame (SAND) framework that can achieve significant speedup using large numbers of commodity machines harnessed from clusters, clouds, and grids. SAND interfaces with the Celera open-source assembly toolkit, replacing two independent sequential modules with scalable parallel alternatives: the candidate selector exploits distributed memory capacity, and the sequence aligner exploits distributed computing capacity. For large problems, these modules provide robust task and data management while also achieving speedup with high efficiency. We show results for several data sets ranging from 738 thousand to over 320 million alignments using resources ranging from a small cluster to more than a thousand nodes spanning three institutions. Christopher Moretti, Andrew Thrasher, Michael Olson, Scott J. Emrich, Douglas Thain |
IEEE Trans. Parallel Distributed Syst. | 6 |
| 2011 | Converting a High Performance Application to an Elastic Cloud ApplicationabstractOver the past decade, high performance applications have embraced parallel programming and computing models. While parallel computing offers advantages such as good utilization of dedicated hardware resources, it also has several drawbacks such as poor fault-tolerance, scalability, and ability to harness available resources during run-time. The advent of cloud computing presents a viable and promising alternative to parallel computing because of its advantages in offering a distributed computing model. In this work, we establish directives that serve as guidelines for the design and implementation or identification of a suitable cloud computing framework to build or convert a high performance application to run in the cloud. We show that following these directives leads to an elastic implementation that has better scalability, run-time resource adaptability, fault tolerance, and portability across cloud computing platforms, while requiring minimal effort and intervention from the user. We illustrate this by converting an MPI implementation of replica exchange, a parallel tempering molecular dynamics application, to an elastic cloud application using the Work Queue framework that adheres to these directive. We observe better scalability and resource adaptability of this elastic application on multiple platforms, including a homogeneous cluster environment (SGE) and heterogeneous cloud computing environments such as Microsoft Azure and Amazon EC2. Dinesh Rajan, Anthony Canino, Jesús A. Izaguirre, Douglas Thain |
CloudCom | 4 |
| 2011 | Expert-Citizen Engineering: "Crowdsourcing" Skilled CitizensabstractCitizen Engineering (CE) is a concept that engages a cohort of physically dispersed citizens connected by the Internet to collaboratively solve real-world problems through massive cooperation. With advances in information technology, we can build transformative cyber-infrastructures to effectively leverage the ''wisdom of crowds''. Regarding the citizen engineers who function as main contributors, there is a wide spectrum of human resources that crowd sourcing system designers can harness - from amateurs/hobbyists, lacking practical experience, to experts/licensed engineers, with years of professional training. As such, we are encouraged to investigate proper approaches to design CEs that can sufficiently engage and support expert citizens who have unique needs that may be different from those of amateur citizen engineers. In this study, we focused on a system designed for engaging high-end users - expert citizens. Our experiment is based on a web site -- ''Expert Citizen Engineering Experiment'' developed to fulfill a sophisticated civil engineering task. The conclusions generated from this experiment provide guidance for future CE project designs, where skilled users are the main contributors. Based on our observations and post-experiment interviews, we believe that expert citizen engineers have higher expectations on computation platform capacity and system stability compared to average citizen engineers. Meanwhile, it should be acknowledged that in the domain of civil engineering, high reliability and trustworthiness are particularly emphasized. Zhi Zhai, Peter Sempolinski, Douglas Thain, Gregory R. Madey, Daniel Wei, Ahsan Kareem |
DASC | 3 |
| 2011 | Biocompute 2.0: an improved collaborative workspace for data intensive bio-scienceabstractSUMMARY The explosion of data in the biological community requires scalable and flexible portals for bioinformatics. To help address this need, we proposed characteristics needed for rigorous, reproducible, and collaborative resources for data‐intensive science. Implementing a system with these characteristics exposed challenges in user interface, data distribution, and workflow description/execution. We describe ongoing responses to these and other challenges. Our Data‐Action‐Queue design pattern addresses user interface and system organization concepts. A dynamic data distribution mechanism lays the foundation for the management of persistent datasets. Makeflow facilitates the simple description and execution of complex multi‐part jobs and forms the kernel of a module system powering diverse bioinformatics applications. Our improved web portal, Biocompute 2.0, has been in production use since the summer of 2010. Through it and its predecessor, we have provided over 56 years of CPU time through its five modules—BLAST, SSAHA, SHRIMP, BWA, and SNPEXP—to research groups at three universities. In this paper, we describe the goals and interface to the system, its architecture and performance, and the insights gained in its development. Copyright © 2011 John Wiley & Sons, Ltd. Rory Carmichael, Patrick Braga-Henebry, Douglas Thain, Scott J. Emrich |
Concurr. Comput. Pract. Exp. | 3 |
| 2010 | Attaching Cloud Storage to a Campus Grid Using Parrot, Chirp, and HadoopabstractThe Hadoop file system is a large scale distributed file system used to manage and quickly process extremely large data sets. We want to utilize Hadoop to assist with data-intensive workloads in a distributed campus grid environment. Unfortunately, the Hadoop file system is not designed to work in such an environment easily or securely. We present a solution that bridges the Chirp distributed file system to Hadoop for simple access to large data sets. Chirp layers on top of Hadoop many grid computing desirables including simple deployment without special privileges, easy access via Parrot, and strong and flexible security Access Control Lists (ACL). We discuss the challenges involved in using Hadoop on a campus grid and evaluate the performance of the combined systems. Patrick Donnelly, Peter Bui, Douglas Thain |
CloudCom | 3 |
| 2010 | A Comparison and Critique of Eucalyptus, OpenNebula and NimbusabstractEucalyptus, Open Nebula and Nimbus are three major open-source cloud-computing software platforms. The overall function of these systems is to manage the provisioning of virtual machines for a cloud providing infrastructure-as-a-service. These various open-source projects provide an important alternative for those who do not wish to use a commercially provided cloud. We provide a comparison and analysis of each of these systems. We begin with a short summary comparing the current raw feature set of these projects. After that, we deepen our analysis by describing how these cloud management frameworks relate to the many other software components required to create a functioning cloud computing system. We also analyse the overall structure of each of these projects and address how the differing features and implementations reflect the different goals of each of these projects. Lastly, we discuss some of the common challenges that emerge in setting up any of these frameworks and suggest avenues of further research and development. These include the problem of fair scheduling in absence of money, eviction or preemption, the difficulties of network configuration, and the frequent lack of clean abstractions. Peter Sempolinski, Douglas Thain |
CloudCom | 2 |
| 2010 | Grid, Cluster and Cloud Computing
Kate Keahey, Domenico Laforenza, Alexander Reinefeld, Pierluigi Ritrovato, Douglas Thain, Nancy Wilkins-Diehr |
Euro-Par (1) | 5 |
| 2010 | ROARS: a scalable repository for data intensive scientific computingabstractAs scientific research becomes more data intensive, there is an increasing need for scalable, reliable, and high performance storage systems. Such data repositories must provide both data archival services and rich metadata, and cleanly integrate with large scale computing resources. ROARS is a hybrid approach to distributed storage that provides both large, robust, scalable storage and efficient rich metadata queries for scientific applications. In this paper, we demonstrate that ROARS is capable of importing and exporting large quantities of data, migrating data to new storage nodes, providing robust fault tolerance, and generating materialized views based on metadata queries. Our experimental results demonstrate that ROARS' aggregate throughput scales with the number of concurrent clients while providing fault-tolerant data access. ROARS is currently being used to store 5.1TB of data in our local biometrics repository. Hoang Bui, Peter Bui, Patrick J. Flynn, Douglas Thain |
HPDC | 4 |
| 2010 | Towards long term data quality in a large scale biometrics experimentabstractQuality of data plays a very important role in any scientific research. In this paper we present some of the challenges that we face in managing and maintaining data quality for a terabyte scale biometrics repository. We have developed a step by step model to capture, ingest, validate, and prepare data for biometrics research. During these processes, there are many hidden errors which can be introduced into the data. Those errors can affect the overall quality of data, and thus can skew the results of biometrics research. We discuss necessary steps we have taken to reduce and eliminate the errors. Steps such as data replication, automated data validation, and logging metadata changes are both necessary and crucial to improve the quality and reliability of our data. Hoang Bui, Diane Wright, Clarence Helm, Rachel Witty, Patrick J. Flynn, Douglas Thain |
HPDC | 6 |
| 2010 | Weaver: integrating distributed computing abstractions into scientific workflows using PythonabstractWeaver is a high-level framework that enables researchers to integrate distributed computing abstractions into their scientific workflows. Rather than develop a new workflow language, we built Weaver on top of the Python programming language. As such, Weaver takes advantage of users' familiarity with Python, minimizes barriers to adoption, and allows for integration with existing software. In this paper, we introduce Weaver's programming model, which consists of datasets, functions, and abstractions that users combine to organize and specify large-scale scientific workflows. We also explain how these specifications are compiled into a directed acyclic graph used by a workflow manager that dispatches the work to a variety of distributed computing engines. To examine how Weaver is used in scientific research, we present three example applications that demonstrate Weaver's ability to integrate into existing workflows and incorporate optimized distributed computing abstraction tools. Peter Bui, Douglas Thain |
HPDC | 3 |
| 2010 | Biocompute: towards a collaborative workspace for data intensive bio-scienceabstractThe explosion of data in the biological community demands the development of more scalable and flexible portals for bioinformatic computation. To address this need, we put forth characteristics needed for rigorous, reproducible, and collaborative resources for data intensive science. Implementing a system with these characteristics exposed challenges in user interface, data distribution, and workflow description/execution. We describe several responses to these challenges. The Data-Action-Queue metaphor addresses user interface and system organization concepts. A dynamic data distribution mechanism lays the foundation for the management of persistent datasets. The Makeflow workflow facilitates the simple description and execution of complex multipart jobs. The resulting web portal, Biocompute, has been in production use at the University of Notre Dame's Bioinformatics Core Facility since the summer of 2009. It has provided over seven years of CPU time through its three sequence search modules --- BLAST, SSAHA, and SHRIMP --- to ten biological and bioinformatic research groups spanning three universities. In this paper we describe the goals and interface to the system, its architecture and performance, and the insights gained in its development. Rory Carmichael, Patrick Braga-Henebry, Douglas Thain, Scott J. Emrich |
HPDC | 3 |
| 2010 | Visualizing massively multithreaded applications with ThreadScopeabstractAbstract As highly parallel multicore machines become commonplace, programs must exhibit more concurrency to exploit the available hardware. Many multithreaded programming models already encourage programmers to create hundreds or thousands of short‐lived threads that interact in complex ways. Programmers need to be able to analyze, tune, and troubleshoot these large‐scale multithreaded programs. To address this problem, we present ThreadScope: a tool for tracing, visualizing, and analyzing massively multithreaded programs. ThreadScope extracts the machine‐independent program structure from execution trace data from a variety of tracing tools and displays it as a graph of dependent execution blocks and memory objects, enabling identification of synchronization and structural problems, even if they did not occur in the traced run. It also uses graph‐based analysis to identify potential problems. We demonstrate the use of ThreadScope to view program structure, memory access patterns, and synchronization problems in three programming environments and seven applications. Copyright © 2009 John Wiley & Sons, Ltd. Kyle B. Wheeler, Douglas Thain |
Concurr. Comput. Pract. Exp. | 2 |
| 2010 | All-Pairs: An Abstraction for Data-Intensive Computing on Campus GridsabstractToday, campus grids provide users with easy access to thousands of CPUs. However, it is not always easy for nonexpert users to harness these systems effectively. A large workload composed in what seems to be the obvious way by a naive user may accidentally abuse shared resources and achieve very poor performance. To address this problem, we argue that campus grids should provide end users with high-level abstractions that allow for the easy expression and efficient execution of data-intensive workloads. We present one example of an abstraction—All-Pairs—that fits the needs of several applications in biometrics, bioinformatics, and data mining. We demonstrate that an optimized All-Pairs abstraction is both easier to use than the underlying system, achieve performance orders of magnitude better than the obvious but naive approach, and is both faster and more efficient than a tuned conventional approach. This abstraction has been in production use for one year on a 500 CPU campus grid at the University of Notre Dame and has been used to carry out a groundbreaking analysis of biometric data. Christopher Moretti, Hoang Bui, Karen Hollingsworth, Brandon Rich, Patrick J. Flynn, Douglas Thain |
IEEE Trans. Parallel Distributed Syst. | 6 |
| 2009 | The quest for scalable support of data-intensive workloads in distributed systemsabstractData-intensive applications involving the analysis of large datasets often require large amounts of compute and storage resources, for which locality can be crucial to high throughput and performance. We propose a data approach that acquires compute and storage resources dynamically, replicates in response to demand, and schedules computations close to data. As demand increases, more resources are acquired, thus allowing faster response to subsequent requests that refer to the same data; when demand drops, resources are released. This approach can provide the benefits of dedicated hardware without the associated high costs, depending on workload and resource characteristics. To explore the feasibility of diffusion, we offer both a theoretical and an empirical analysis. We define an abstract model for diffusion, introduce new scheduling policies with heuristics to optimize real-world performance, and develop a competitive online cache eviction policy. We also offer many empirical experiments to explore the benefits of dynamically expanding and contracting resources based on load, to improve system responsiveness while keeping wasted resources small. We show performance improvements of one to two orders of magnitude across three diverse workloads when compared to the performance of parallel file systems with throughputs approaching 80 Gb/s on a modest cluster of 200 processors. We also compare diffusion with a best model for active storage, contrasting the difference between a pull-model found in diffusion and a push-model found in active storage. Ioan Raicu, Ian T. Foster, Yong Zhao 0009, Philip Little, Christopher Moretti, Amitabh Chaudhary, Douglas Thain |
HPDC | 7 |
| 2009 | Harnessing parallelism in multicore clusters with the all-pairs and wavefront abstractionsabstractBoth distributed systems and multicore computers are difficult programming environments. Although the expert programmer may be able to tune distributed and multicore computers to achieve high performance, the non-expert may struggle to achieve a program that even functions correctly. Christopher Moretti, Scott J. Emrich, Kenneth Judd, Douglas Thain |
HPDC | 5 |
| 2009 | Chirp: a practical global filesystem for cluster and Grid computing
Douglas Thain, Christopher Moretti, Jeffrey Hemmes |
J. Grid Comput. | 1 |
| 2009 | Reflections on the virtues of modularity: a case study in linux security modulesabstractAbstract Developing a modular system that properly supports a range of security models is challenging. The work presented here details our experiences with the modularLinuxsecurity framework called Linux Security Modules, or LSMs. Throughout our experiences we discovered that the developers of the LSM framework made certain tradeoffs for speed and simplicity during implementation, and consequently leaving the framework incomplete. Our experiences show at which points the theory of the LSM differs from reality, and details how these differences play out when developing and using a custom LSM. Copyright © 2009 John Wiley & Sons, Ltd. Andrew Blaich, Douglas Thain, Aaron Striegel |
Softw. Pract. Exp. | 2 |
| 2008 | BXGrid: A Data Repository and Workflow Abstraction for Biometrics ResearchabstractResearch in the field of biometrics depends on the effective management of large amounts of data and computation. Current research projects in biometrics acquire many terabytes of images and video of subjects in many different modes and situations, annotated with detailed metadata. To study the effectiveness of new algorithms for identifying people, researchers must exhaustively compare large numbers of measurements with a variety of custom functions. The quality of the end results is often dependent upon the sheer amount of data marshalled to support it.To address these challenges, we are constructing BXGrid, an end-to-end computing system for conducting biometrics research. BXGrid assists with the entire research process from data acquisition all the way to generating results for publication. Because the entire chain of research is kept consistently within one system, multiple users may easily share tools and results, building off of each other's work. BXGrid also helps to ensure scientific integrity by automating a variety of consistency checks, external data audits, and reproduction of existing results. Hoang Bui, Deborah Thomas, Michael Kelly, Christopher Lyon, Douglas Thain, Patrick J. Flynn |
eScience | 5 |
| 2008 | Using Small Abstractions to Program Large Distributed SystemsabstractDistributed systems such as clusters, clouds and grids remain a difficult platform for executing large data intensive workloads. Even sophisticated users struggle to shape complex workloads into the simple 'assembly language' of file transfer and job submission. To address this, problem, we advocate the use of abstractions, which are simple frameworks for expressing large structured problems. In this talk, we will discuss three examples of abstractions developed at the University of Notre Dame for scientific applications. In each case, we have been able to scale up workloads one to wto orders of magnitude larger than we previously feasible. Through each example, we will address some persistent obstacles in the field of distributed computing. Douglas Thain, Christopher Moretti, Hoang Bui, Nitesh V. Chawla, Patrick J. Flynn |
eScience | 1 |
| 2008 | DataLab: transactional data-parallel computing on an active storage cloudabstractActive storage clouds are an attractive platform for executing large data intensive workloads found in many fields of science. However, active storage presents new system management challenges. A large system of fault-prone machines with local persistent state can easily degenerate into a mess of unreferenced data and runaway computations. Our solution to this problem is DataLab, a software framework for running data parallel workloads on active storage clusters. DataLab provides a simple language for expressing workloads, works with legacy application codes, and achieves robustness through the use of distributed transactions. Our prototype implementation scales to 250 nodes on a large biometric image processing workload. Brandon Rich, Douglas Thain |
HPDC | 2 |
| 2008 | Scaling up Classifiers to Cloud ComputersabstractAs the size of available datasets has grown from Megabytes to Gigabytes and now into Terabytes, machine learning algorithms and computing infrastructures have continuously evolved in an effort to keep pace. But at large scales, mining for useful patterns still presents challenges in terms of data management as well as computation. These issues can be addressed by dividing both data and computation to build ensembles of classifiers in a distributed fashion, but trade-offs in cost, performance, and accuracy must be considered when designing or selecting an appropriate architecture. In this paper, we present an abstraction for scalable data mining that allows us to explore these trade-offs. Data and computation are distributed to a computing cloud with minimal effort from the user, and multiple models for data management are available depending on the workload and system configuration. We demonstrate the performance and scalability characteristics of our ensembles using a wide variety of datasets and algorithms on a Condor-based pool with Chirp to handle the storage. Christopher Moretti, Karsten Steinhaeuser, Douglas Thain, Nitesh V. Chawla |
ICDM | 3 |
| 2008 | Data mining on the grid for the gridabstractBoth users and administrators of computing grids are presented with enormous challenges in debugging and troubleshooting. Diagnosing a problem with one application on one machine is hard enough, but diagnosing problems in workloads of millions of jobs running on thousands of machines is a problem of a new order of magnitude. Suppose that a user submits one million jobs to a grid, only to discover some time later that half of them have failed, Users of large scale systems need tools that describe the overall situation, indicating what problems are commonplace versus occasional, and which are deterministic versus random. Machine learning techniques can be used to debug these kinds of problems in large scale systems. We present a comprehensive framework from data to knowledge discovery as an important step towards achieving this vision. Nitesh V. Chawla, Douglas Thain, Ryan Lichtenwalter, David A. Cieslak |
IPDPS | 2 |
| 2008 | All-pairs: An abstraction for data-intensive cloud computingabstractAlthough modern parallel and distributed computing systems provide easy access to large amounts of computing power, it is not always easy for non-expert users to harness these large systems effectively. A large workload composed in what seems to be the obvious way by a naive user may accidentally abuse shared resources and achieve very poor performance. To address this problem, we propose that production systems should provide end users with high-level abstractions that allow for the easy expression and efficient execution of data intensive workloads. We present one example of an abstraction - all-pairs - that fits the needs of several data-intensive scientific applications. We demonstrate that an optimized all-pairs abstraction is both easier to use than the underlying system, and achieves performance orders of magnitude better than the obvious but naive approach, and twice as fast as a hand-optimized conventional approach. Christopher Moretti, Jared Bulosan, Douglas Thain, Patrick J. Flynn |
IPDPS | 3 |
| 2008 | Qthreads: An API for programming with millions of lightweight threadsabstractLarge scale hardware-supported multithreading, an attractive means of increasing computational power, benefits significantly from low per-thread costs. Hardware support for lightweight threads is a developing area of research. Each architecture with such support provides a unique interface, hindering development for them and comparisons between them. A portable abstraction that provides basic lightweight thread control and synchronization primitives is needed. Such an abstraction would assist in exploring both the architectural needs of large scale threading and the semantic power of existing languages. Managing thread resources is a problem that must be addressed if massive parallelism is to be popularized. The qthread abstraction enables development of large-scale multithreading applications on commodity architectures. This paper introduces the qthread API and its Unix implementation, discusses resource management, and presents performance results from the HPCCG benchmark. Kyle B. Wheeler, Richard C. Murphy, Douglas Thain |
IPDPS | 3 |
| 2008 | ENAVis: Enterprise Network Activities Visualization
Qi Liao 0002, Andrew Blaich, Aaron Striegel, Douglas Thain |
LISA | 4 |
| 2008 | Making the best of a bad situation: Prioritized storage management in GEMS
Justin M. Wozniak, Paul R. Brenner, Douglas Thain, Aaron Striegel, Jesús A. Izaguirre |
Future Gener. Comput. Syst. | 3 |
| 2008 | Biomolecular committor probability calculation enabled by processing in network storage
Paul R. Brenner, Justin M. Wozniak, Douglas Thain, Aaron Striegel, Jeffrey W. Peng, Jesús A. Izaguirre |
Parallel Comput. | 3 |
| 2007 | Biomolecular Path Sampling Enabled by Processing in Network StorageabstractComputationally complex and data intensive atomic scale biomolecular simulation is enabled via processing in network storage (PINS): a novel distributed system framework to overcome bandwidth, compute, storage, and security challenges inherent to the wide area computation and storage grid. High throughput data generation requirements for our scientific target are overcome through novel aggregate bandwidth capabilities. Biomolecular simulation methods are correlated with the client tools, hybrid database/file server (GEMS), computation engine (Condor), virtual file system adapter (Parrot), and local file servers (Chirp). PINS performance is reported for the path sampling of a solvated protein domain requiring over 1000 simulations with total output data generation on the order of 1TB. Paul R. Brenner, Justin M. Wozniak, Douglas Thain, Aaron Striegel, Jeffrey W. Peng, Jesús A. Izaguirre |
IPDPS | 3 |
| 2007 | Challenges in Executing Data Intensive Biometric Workloads on a Desktop GridabstractDesktop grids have traditionally focused on executing computation intensive workloads. Can they also be used to execute data-intensive workloads? To answer this question, we present a case study of a data intensive biometric application which is infeasible to process on a single machine. We evaluate the capacity of a desktop grid to store and deliver the data need to execute the workload, and compare several general techniques for data deployment. Selecting the most scalable technique, we execute and evaluate five large production workloads on a 350-CPU desktop grid. We observe that this technique is sensitive to many parameters, and propose that an ideal system should be responsible for choosing the proper decomposition of a workload. Christopher Moretti, Timothy C. Faltemier, Douglas Thain, Patrick J. Flynn |
IPDPS | 3 |
| 2006 | Positioning Dynamic Storage Caches for Transient DataabstractSimulations, experiments and observatories are generating a deluge of scientific data. Even more staggering is the ever growing application demand to process and assimilate these datasets. Application users perform a range of data operations, collaborate and share data in many novel ways. The current storage landscape is struggling to keep up with these trends in scientific data processing. Application users pay the price due to over-crowded shared filesystems, or expensive storage area networks, or not enough local storage, or high-latency archival or wide-area transfers. In order to sustain and maximize I/O bandwidth relative to increasing CPU speeds, applications must take advantage of large amounts of intermediate commodity storage, However, intermediate storage presents new challenges above and beyond the traditional distributed file system paradigm: persistent scheduling, storage/CPU coallo-cation, namespace management, lifetime management, and novel application interfaces. In this paper, we describe applications that require intermediate storage management, suggest several open research problems, and illustrate two systems - Freeloader and Tactical Storage - that attack different aspects of these problems Sudharshan S. Vazhkudai, Douglas Thain, Xiaosong Ma, Vincent W. Freeh |
CLUSTER | 2 |
| 2006 | Troubleshooting Distributed Systems via Data MiningabstractThrough massive parallelism, distributed systems enable the multiplication of productivity. Unfortunately, increasing the scale of available machines to users will also multiply debugging when failure occurs. Data mining allows the extraction of patterns within large amounts of data and therefore forms the foundation for a useful method of debugging, particularly within such distributed systems. This paper outlines a successful application of data mining in troubleshooting distributed systems, proposes a framework for further study, and speculates on other future work. 1 David A. Cieslak, Douglas Thain, Nitesh V. Chawla |
HPDC | 2 |
| 2006 | Transparent access to Grid resources for user softwareabstractAbstract Grid computing promises access to large amounts of computing power, but so far adoption of Grid computing has been limited to highly specialized experts for three reasons. First, users are used to batch systems, and interfaces to Grid software are often complex and different to those in batch systems. Second, users are used to having transparent file access, which Grid software does not conveniently provide. Third, efforts to achieve wide‐spread coordination of computers while solving the first two problems is hampered when clusters are on private networks. Here we bring together a variety of software that allows users to almost transparently use Grid resources as if they were local resources while providing transparent access to files, even when private networks intervene. As a motivating example, the BaBar Monte Carlo production system is deployed on a truly distributed environment, the European DataGrid, without any modification to the application itself. Copyright © 2005 John Wiley & Sons, Ltd. Sander Klous, Jaime Frey, Se-Chang Son, Douglas Thain, Alain J. Roy, Miron Livny, Jo F. J. van den Brand |
Concurr. Comput. Pract. Exp. | 4 |
| 2006 | How to measure a large open-source distributed systemabstractHow can we measure the impact of an open-source software package over time? When a system has no price, no purchase contracts and no buyers or sellers it can be difficult to judge its impact on the world. To explore this issue, we have instrumented the Condor distributed batch system in a variety of ways and observed its growth to over 50 000 CPUs at over 1000 sites over five years. Instrumentation methods include automatic updates by e-mail and user datagram protocol (UDP), annotated download records and a voluntary user survey. Each of these metrics has various strengths and weaknesses that we are able to compare and contrast. We also explore the ethical and legal issues surrounding automatic data collection. Surprisingly, we discover that objections to automatic data collection are higher among people that are not using the Condor software. We conclude with some practical advice for further research into the measurement of software systems. Copyright © 2006 John Wiley & Sons, Ltd. Douglas Thain, Todd Tannenbaum, Miron Livny |
Concurr. Comput. Pract. Exp. | 1 |
| 2005 | Identity boxing: secure user-level containment for the gridabstractToday, a public key infrastructure allows grid users to be identified with strong cryptographic credentials and and a descriptive, globally-unique name such as /O=UnivNowhere/CN=Fred. This powerful security infrastructure allows users to perform a single login and then access a variety of remote resources on the grid without further authentication steps. However, once connected to a specific system, a user's grid credentials must somehow be mapped to a local namespace. This creates a significant burden upon the administrator of each site to manage a continuously-changing user list. Large systems have worked around this by employing the old insecure standby of shared user accounts. A single user may be known by a different account name at every single site that he or she accesses, in addition to a variety of identity names given by certificate authorities. In order to access a resource, the user may need to have a local account generated. In order to share resources, each user must know the local identities of users that he/she wishes to share with. To solve these problems, we introduce the technique of identity boxing. An identity box is a well-defined execution space in which all processes and resources are associated with an external identity that need not have any relationship to the set of local accounts. That is, within an identity box, a program runs with an explicit grid identity string rather than with a simple integer UID. As a program executes, all access controls are performed using the high level name rather than the low-level account information. A single Unix account may be used to securely manage several identity boxes simultaneously, thus eliminating the need to services to run as root merely to change identities. Douglas Thain |
HPDC | 1 |
| 2005 | Generosity and gluttony in GEMS: grid enabled molecular simulationsabstractBiomolecular simulations produce more output data than can be managed effectively by traditional computing systems. Researchers need distributed systems that allow the pooling of resources, the sharing of simulation data, and the reliable publication of both tentative and final results. To address this need, we have designed GEMS, a system that enables biomolecular researchers to store, search, and share large scale simulation data. The primary design problem is striking a balance between generosity and gluttony. On one hand, storage providers wish to be generous and share resources with their collaborators. On the other hand, an unchecked data producer can be gluttonous and easily replicate data unnecessarily until it fills all available space. To balance generosity and gluttony, GEMS allows both storage providers and data producers to state and enforce policies on the consumption of storage and the replication of data. By taking advantage of known properties of simulation data, the system is able to distinguish between high value final results that must be preserved and low value intermediate results that can be deleted and regenerated if necessary. We have built a prototype of GEMS on a cluster of workstations and demonstrate its ability to store new data, to replicate within policy limits, and to recover from failures. Justin M. Wozniak, Paul R. Brenner, Douglas Thain, Aaron Striegel, Jesús A. Izaguirre |
HPDC | 3 |
| 2005 | Identity Boxing: A New Technique for Consistent Global IdentityabstractToday, users of the grid may easily authenticate themselves to computing resources around the world using a public key security infrastructure. However, users are forced to employ a patchwork of local identities, each assigned by a different local authority. This forces each grid system to provide a mapping from global to local identities, creating a significant administrative burden and inhibiting many possibilities of data sharing. To remedy this, we introduce the technique of identity boxing. This technique allows a high-level identity to be attached directly to each process and resource that a user employs, rendering the local account name irrelevant. This allows a grid user to be known by the same name consistently at all sites, thus reducing administrative burdens and enabling new forms of sharing. We have implemented identity boxing at the user level within a secure system-call interposition agent and applied it to a distributed storage and execution system. The performance overhead of this implementation is only 0.7 to 6.5 percent for a selection of scientific applications, but as high as 35 percent for a metadata-intensive software build. We conclude with some reflections on how the operating system might be modified to better support grid computing. Douglas Thain |
SC | 1 |
| 2005 | Separating Abstractions from Resources in a Tactical Storage SystemabstractSharing data and storage space in a distributed system remains a difficult task for ordinary users, who are constrained to the fixed abstractions and resources provided by administrators. To remedy this situation, we introduce the concept of a tactical storage system (TSS) that separates storage abstractions from storage resources, leaving users free to create, reconfigure, and destroy abstractions as their needs change. In this paper, we describe how a TSS can provide a variety of filesystem and database abstractions for unmodified applications without requiring special privileges or kernel changes. A TSS provides performance competitive with NFS for single clients and also scales well for multiple servers and multiple clients. A prototype TSS of 120 disks and 6 TB of storage has been deployed at the University of Notre Dame and used for applications in high energy physics and bioinformatics. Douglas Thain, Sander Klous, Justin M. Wozniak, Paul R. Brenner, Aaron Striegel, Jesús A. Izaguirre |
SC | 1 |
| 2005 | Distributed computing in practice: the Condor experienceabstractAbstract Since 1984, the Condor project has enabled ordinary users to do extraordinary computing. Today, the project continues to explore the social and technical problems of cooperative computing on scales ranging from the desktop to the world‐wide computational Grid. In this paper, we provide the history and philosophy of the Condor project and describe how it has interacted with other projects and evolved along with the field of distributed computing. We outline the core components of the Condor system and describe how the technology of computing must correspond to social structures. Throughout, we reflect on the lessons of experience and chart the course travelled by research ideas as they grow into production systems. Copyright © 2005 John Wiley & Sons, Ltd. Douglas Thain, Todd Tannenbaum, Miron Livny |
Concurr. Pract. Exp. | 1 |
| 2004 | Explicit Control in the Batch-Aware Distributed File System
John Bent, Douglas Thain, Andrea C. Arpaci-Dusseau, Remzi H. Arpaci-Dusseau, Miron Livny |
NSDI | 2 |
| 2003 | XtremWeb & Condor sharing resources between Internet connected Condor poolsabstractGrid computing presents two major challenges for deploying large scale applications across wide area networks gathering volunteers PC and clusters/parallel computers as computational resources: security and fault tolerance. This paper presents a lightweight Grid solution for the deployment of multi-parameters applications on a set of clusters protected by firewalls. The system uses a hierarchical design based on Condor for managing each cluster locally and XtremWeb for enabling resource sharing among the clusters. We discuss the security and fault tolerance mechanisms used for this design and demonstrate the usefulness of the approach measuring the performances of a multi-parameters bio-chemistry application deployed on two sites: University of Wisconsin/Madison and Paris South University. This experiment shows that we can efficiently and safely harness the computational power of about 200 PC distributed on two geographic sites. Oleg Lodygensky, Gilles Fedak, Franck Cappello, Vincent Néri, Miron Livny, Douglas Thain |
CCGRID | 6 |
| 2003 | Pipeline and Batch Sharing in Grid WorkloadsabstractWe present a study of six batch-pipeline scientific workloads that are candidates for execution on computational grids. Whereas other studies focus on the behavior of single applications, this study characterizes workloads composed of pipelines of sequential processes that use file storage for communication and also share measurements of the memory, CPU, and I/O requirements of individual components as well as analyses of I/O sharing within complete batches. We conclude with a discussion of the ramifications of these workloads for end-to-end scalability and overall system design. Douglas Thain, John Bent, Andrea C. Arpaci-Dusseau, Remzi H. Arpaci-Dusseau, Miron Livny |
HPDC | 1 |
| 2003 | The Ethernet Approach to Grid ComputingabstractDespite many competitors, Ethernet became the dominant protocol for local area networking due to its simplicity, robustness, and efficiency in wide variety of conditions and technology. Reflecting on the current frailty of much software, grid and otherwise, we propose that the Ethernet approach to resource sharing is an effective and reliable technique for combining coarse-grained software when failures are common and poorly detailed. This approach involves placing several simple but important responsibilities on client software to acquire shared resources conservatively, to back off during periods of failure, and to inform competing clients when resources are in contention. We present a simple scripting language that simplifies and encourages the Ethernet approach, and demonstrate its use in several grid computing scenarios, including job submission, disk allocation, and data replication. We conclude with a discussion of the limitations of this approach, and describe how it is uniquely suited to high-level programming. Douglas Thain, Miron Livny |
HPDC | 1 |
| 2002 | Error Scope on a Computational Grid: Theory and PracticeabstractError propagation is a central problem in grid computing. We re-learned this while adding a Java feature to the Condor computational grid. Our initial experience with the system was negative, due to the large number of new ways in which the system could fail. To reason about this problem, we developed a theory of error propagation. Central to our theory is the concept of an error's scope, defined as the portion of a system that it invalidates. With this theory in hand, we recognized that the expanded system did not properly consider the scope of errors it discovered. We modified the system according to our theory, and succeeded in making it a more robust platform for distributed computing. Douglas Thain, Miron Livny |
HPDC | 1 |
| 2002 | Caveat Emptor: Making Grid Services Dependable from the Client SideabstractGrid computing relies on fragile partnerships. Clients with hundreds or even thousands of pending service requests must seek out and form temporary alliances with remote servers eager to satisfy them. Yet, despite the high quality and reliability of these servers and their software, unexpected events and behavior are common. Communication networks, power systems, operating systems, middleware and operator intervention all conspire to attack even the most carefully arranged client-server interaction. To survive in such an imperfect world, customers of grid resources must be equipped with resilient client software that tolerates failures while aggressively representing their interests. Following our tradition of developing technology that harnesses the power of opportunistic resources, the Condor Project is actively engaged in developing the basic mechanisms for building dependable and effective grid computing clients. Guided by our experience and the practical needs of production users in disciplines as diverse as astronomy and sociology, the Project aims to equip users with powerful software that complements the reliability of the servers that they exploit. Our most visible product is the Condor-G job manager. Other research ventures, including the full Condor distributed system, offer valuable lessons in dependable client-side management. Dependability has been explored in a number of branches of computing, ranging from database systems to programming languages. The hard-earned lessons from these fields are also essential to grid computing. Fundamental concepts such as timeouts, logging, checkpoints, transactions, leases, and atomic operations must be employed and expressed in basic protocols and interfaces for CPU and I/O access. Without these techniques, clients and servers lose track of the other’s state, leading to missed opportunities, wasted resources, incorrect results, and unnecessary failures. This principle is espoused in systems such as Condor-G and protocols such as the most recent version of GRAM. In a Grid environment we must never view failure as a disaster. Rather, failures occur at every level and every interface, and must be expected and structured. No single failure must bring a computation to a halt, nor can any type of failure be retried indefinitely. Jobs may be retracted even from systems deemed reliable when better performance may be found elsewhere. In addition, we must always be careful to determine whether the source of a failure lies in the system or in the job itself. Examples of this principle are found in the DAGMan meta-scheduler and the fault-tolerant shell. Miron Livny, Douglas Thain |
PRDC | 2 |
| 2001 | The Kangaroo Approach to Data Movement on the GridabstractAccess to remote data is one of the principal challenges of Grid computing. While performing I/O, Grid applications must be prepared for server crashes, performance variations and exhausted resources. To achieve high throughput in such a hostile environment, applications need a resilient service that moves data while hiding errors and latencies. We illustrate this idea with Kangaroo, a simple data movement system that makes opportunistic use of disks and networks to keep applications running. We demonstrate that Kangaroo can achieve better end-to-end performance than traditional data movement techniques, even though its individual components do not achieve high performance. Douglas Thain, Jim Basney, Se-Chang Son, Miron Livny |
HPDC | 1 |
| 2001 | Gathering at the well: creating communities for grid I/OabstractGrid applications have demanding I/O needs. Schedulers must bring jobs and data in close proximity in order to satisfy throughput, scalability, and policy requirements. Most systems accomplish this by making either jobs or data mobile. We propose a system that allows jobs and data to meet by binding execution and storage sites together into I/O communities which then participate in the wide-area system. The relationships between participants in a community may be expressed by the ClassAd framework. Extensions to the framework allow community members to express indirect relations. We demonstrate our implementation of I/O communities by improving the performance of a key high-energy physics simulation on an international distributed system. Douglas Thain, John Bent, Andrea C. Arpaci-Dusseau, Remzi H. Arpaci-Dusseau, Miron Livny |
SC | 1 |
| 2000 | Bypass: A Tool for Building Split Execution SystemsabstractSplit execution is a common model for providing a friendly environment on a foreign machine. In this model, a remotely executing process sends some or all of its system calls back to a home environment for execution. Unfortunately, hand-coding split execution systems for experimentation and research is difficult and error-prone. We have built a tool, called Bypass, for quickly producing portable and correct split execution systems for unmodified legacy applications. We demonstrate Bypass by using it to transparently connect a POSIX application to a simple data staging system based on the Globus toolkit. Douglas Thain, Miron Livny |
HPDC | 1 |