EDBT 2026 Demo / reviewers in the wild / expert
Jon B. Weissman
dblp:52/4225
· DBLP profile ↗
66ranked-venue papers
15as first author
4since 2021 · last 2025
0000-0001-8799-317XORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 54 · 15 first-author · 3 since 2021Software engineering, systems software and programming languages · 4Computer networks · 1Applied, interdisciplinary, general and emerging computing · 1
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2025 | ASTRA: Association, Spatial proximity and Temporal Relevance based Adaptive prefetching for Edge ARabstractMobile Augmented Reality (MAR) applications face performance challenges due to their high computational demands and need for low-latency responses. Traditional approaches like on-device storage or reactive data fetching from the cloud often result in limited augmented reality (AR) experiences. Edge caching, which caches AR objects closer to the user, provides a promising solution. However, existing edge caching approaches do not consider AR-specific features such as AR object sizes, user interactions, user’s field of view and physical location in a coherent manner. This paper investigates how to further optimize edge caching by employing AR-aware prefetching techniques. We present ASTRA, a prefetching framework tailored for mobile augmented reality edge caches. It integrates object associations derived from user interaction patterns with spatial awareness based on the user’s physical location and field of view. This approach employs an association factor per object that considers the recency of object co-access; and a lazy fetching strategy that prioritizes prefetching only when the user is in close proximity to the virtual objects. Furthermore, ASTRA incorporates an adaptive tuning algorithm for minimum support in association rule generation to minimize the computation overhead, making it a distinct and effective solution for enhancing user experience in AR applications by ensuring timely virtual object availability.Through extensive evaluation using both synthetic and real-world workloads, we demonstrate that ASTRA significantly improves cache hit rates compared to current prefetching algorithms, achieving gains in hit rate of upto 35% and end-to-end latency by upto 14%. Further, we demonstrate that the adaptive tuning algorithm that automatically tunes minimum support further improves the hit rate of ASTRA by 10%. Our findings demonstrate the potential of ASTRA to substantially enhance the user experience in MAR applications by ensuring the timely availability of virtual objects. Nikhil Sreekumar, Abhishek Chandra, Jon B. Weissman |
IC2E | 3 |
| 2024 | Jingle: IoT-Informed Autoscaling for Efficient Resource Management in Edge ComputingabstractEdge computing is increasingly applied to various systems for its proximity to end-users and data sources. To facilitate the deployment of diverse edge-native applications, container technology has emerged as a favored solution due to its simplicity in development and resource management. However, deploying edge applications at scale can quickly overwhelm edge resources, potentially leading to violations of service-level objectives (SLOs). Scheduling edge containerized applications to meet SLOs while efficiently managing resources is a significant challenge. In this paper, we introduce Jingle, an autoscaler for edge clusters designed to efficiently scale edge-native applications. Jingle utilizes application performance metrics and domain-specific insights collected from IoT devices to construct a hybrid model. This hybrid model combines a predictive-reactive module with a lightweight learning model. We demonstrate Jingle’s effectiveness through a real-world deployment in a classroom setting, managing two edge-native applications across edge configurations. Our experimental results show that Jingle can fulfill SLO requirements while requiring up to 50% fewer containers than a state-of-the-art cloud scheduler, which highlights its resource management efficiency and SLO compliance. Abhishek Chandra, Jon B. Weissman |
CCGrid | 3 |
| 2022 | Towards Elasticity in Heterogeneous Edge-dense EnvironmentsabstractEdge computing has enabled a large set of emerging edge applications by exploiting data proximity and offloading computation-intensive workloads to nearby edge servers. However, supporting edge application users at scale poses challenges due to limited point-of-presence edge sites and constrained elasticity. In this paper, we introduce a densely-distributed edge resource model that leverages capacity-constrained volunteer edge nodes to support elastic computation offloading. Our model also enables the use of geo-distributed edge nodes to further support elasticity. Collectively, these features raise the issue of edge selection. We present a distributed edge selection approach that relies on client-centric views of available edge nodes to optimize average end-to-end latency, with considerations of system heterogeneity, resource contention and node churn. Elasticity is achieved by fine-grained performance probing, dynamic load balancing, and proactive multi-edge node connections per client. Evaluations are conducted in both real-world volunteer environments and emulated platforms to show how a common edge application, namely AR-based cognitive assistance, can benefit from our approach and deliver low-latency responses to distributed users at scale. Zhiying Liang, Nikhil Sreekumar, Sumanth Kaushik 0001, Abhishek Chandra, Jon B. Weissman |
ICDCS | 6 |
| 2022 | Network Cost-Aware Geo-Distributed Data Analytics SystemabstractMany geo-distributed data analytics (GDA) systems have focused on the network performance-bottleneck: inter-data center network bandwidth to improve performance. Unfortunately, these systems may encounter acost-bottleneck(${\$}$) because they have not considered data transfer cost (${\$}$), one of the most expensive and heterogeneous resources in a multi-cloud environment. In this article, we presentKimchi, a network cost-aware GDA system to meet the cost-performance tradeoff by exploiting data transfer cost heterogeneity to avoid the cost-bottleneck. Kimchi determines cost-aware task placement decisions for scheduling tasks given inputs including data transfer cost, network bandwidth, input data size and locations, and desired cost-performance tradeoff preference. In addition, Kimchi is also mindful of data transfer cost in the presence of dynamics. Kimchi has been applied to two common GDA MapReduce models: synchronous barrier and asynchronous push-based shuffle. A Kimchi prototype has been implemented on Spark, and experiments show that it reduces cost by 5%$\scriptstyle \sim$24% without impacting performance and reduces query execution time by 45%$\scriptstyle \sim$70% without impacting cost compared to other baseline approaches centralized, vanilla Spark, and bandwidth-aware (e.g., Iridium). More importantly, Kimchi allows applications to explore a much richer cost-performance tradeoff space in a multi-cloud environment. Kwangsung Oh, Minmin Zhang, Abhishek Chandra, Jon B. Weissman |
IEEE Trans. Parallel Distributed Syst. | 4 |
| 2020 | A Network Cost-aware Geo-distributed Data Analytics SystemabstractMany geo-distributed data analytics (GDA) systems have focused on the network performance-bottleneck: interdata center network bandwidth to improve performance. Unfortunately, these systems may encounter a cost-bottleneck ($) because they have not considered data transfer cost ($), one of the most expensive and heterogeneous resources in a multi-cloud environment. In this paper, we present Kimchi, a network cost-aware GDA system to meet the cost-performance tradeoff by exploiting data transfer cost heterogeneity to avoid the cost-bottleneck. Kimchi determines cost-aware task placement decisions for scheduling tasks given inputs including data transfer cost, network bandwidth, input data size and locations, and desired cost-performance tradeoff preference. In addition, Kim- chi is also mindful of data transfer cost in the presence of dynamics. A Kimchi prototype has been implemented on Spark and experiments show that it reduces cost by 14% ~ 24% without impacting performance and reduces query execution time by 45% ~ 70% without impacting cost compared to other baseline approaches centralized, vanilla Spark, and bandwidth-aware (e.g. Iridium). More importantly, Kimchi allows applications to explore a much richer cost-performance tradeoff space in a multi-cloud environment. Kwangsung Oh, Abhishek Chandra, Jon B. Weissman |
CCGRID | 3 |
| 2020 | Position Paper: Towards a Robust Edge-Native Storage SystemabstractEdge environments are generating an increasingly large amount of data due to the proliferation of edge devices. Accommodating this large influx of data at edge servers is a challenging issue. While some data can be processed as it is generated, others must be stored for later access. This paper proposes the features that a new edge-native storage system must possess including support for user mobility and node fluctuation. To motivate this, we first describe several emerging edge applications and their data needs. We then describe the challenges in meeting these needs. We then evaluate an out-of-the-box cloud storage system, Cassandra, to assess it's suitability as an edge storage system due to many edge-friendly features. We determined that while a cloud-based storage system can be ported to the edge meeting some of the challenges, other challenges require new solutions. Based on the challenges and the results of Cassandra case study, we propose a set of design principles for a new edge-native storage system. Nikhil Sreekumar, Abhishek Chandra, Jon B. Weissman |
SEC | 3 |
| 2020 | WASP: Wide-area Adaptive Stream ProcessingabstractAdaptability is critical for stream processing systems to ensure stable, low-latency, and high-throughput processing of long-running queries. Such adaptability is particularly challenging for wide-area stream processing due to the highly dynamic nature of the wide-area environment, which includes unpredictable workload patterns, variable network bandwidth, occurrence of stragglers, and failures. Unfortunately, existing adaptation techniques typically achieve these performance goals by compromising the quality/accuracy of the results, and they are often application-dependent. In this work, we rethink the adaptability property of wide-area stream processing systems and propose a resource-aware adaptation framework, called WASP. WASP adapts queries through a combination of multiple techniques: task re-assignment, operator scaling, and query re-planning, and applies them in a WAN-aware manner. It is able to automatically determine which adaptation action to take depending on the type of queries, dynamics, and optimization goals. We have implemented a WASP prototype on Apache Flink. Experimental evaluation with the YSB benchmark and a real Twitter trace shows that WASP can handle various dynamics without compromising the quality of the results. Albert Jonathan, Abhishek Chandra, Jon B. Weissman |
Middleware | 3 |
| 2020 | Wiera: Policy-Driven Multi-Tiered Geo-Distributed Cloud Storage SystemabstractMulti-tiered geo-distributed cloud storage systems must tame complexity at many levels: uniform APIs for storage access, supporting flexible storage policies that meet a wide array of application metrics, determining an optimal data placement, handling uncertain network dynamics and access dynamism, and operating across many levels of heterogeneity both within and across data-centers (DCs). In this paper, we present an integrated solution called Wiera. Wiera enables the specification of data management policies both within a local DC and across DCs. Such policies enable the user to optimize for cost, performance, reliability, durability, and consistency, and to express their tradeoffs. In addition, Wiera determines an optimal data placement for the user to meet their desired tradeoffs easily in such an environment. A key aspect of Wiera is first-class support for dynamism due to network, workload, and access patterns changes. As far as we know, Wiera is the first geo-distributed cloud storage system which handles dynamism actively at run-time. Wiera allowsunmodified applicationsto reap the benefits of flexible data/storage policies by externalizing the policy specification. We show how Wiera enables a rich specification of dynamic policies using a concise notation and describe the design and implementation of the system. We have implemented a Wiera prototype on multiple cloud environments, AWS and Azure, that illustrates potential benefits from managing dynamics and in using multiple cloud storage tiers both within and across DCs. Kwangsung Oh, Nan Qin, Abhishek Chandra, Jon B. Weissman |
IEEE Trans. Parallel Distributed Syst. | 4 |
| 2018 | Multi-Query Optimization in Wide-Area Streaming AnalyticsabstractWide-area data analytics has gained much attention in recent years due to the increasing need for analyzing data that are geographically distributed. Many of such queries often require real-time analysis on data streams that are continuously being generated across multiple locations. Yet, analyzing these geo-distributed data streams in a timely manner is very challenging due to the highly heterogeneous and limited bandwidth availability of the wide-area network (WAN). This paper examines the opportunity of applying multi-query optimization in the context of wide-area streaming analytics, with the goal of utilizing WAN bandwidth efficiently while achieving high throughput and low latency execution. Our approach is based on the insight that many streaming analytics queries often exhibit common executions, whether in consuming a common set of input data or performing the same data processing. In this work, we study different types of sharing opportunities and propose a practical online algorithm that allows streaming analytics queries to share their common executions incrementally. We further address the importance of WAN awareness in applying multi-query optimization. Without WAN awareness, sharing executions in a wide-area environment may lead to performance degradation. We have implemented our WAN-aware multi-query optimization in a prototype implementation based on Apache Flink. Experimental evaluation using Twitter traces on a real wide-area system deployment across geo-distributed EC2 data centers shows that our technique is able to achieve 21% higher throughput while saving WAN bandwidth consumption by 33% compared to a WAN-aware, sharing-agnostic system. Albert Jonathan, Abhishek Chandra, Jon B. Weissman |
SoCC | 3 |
| 2018 | Dynamically negotiating capacity between on-demand and batch clusters
Kate Keahey, Pierre Riteau, Jon B. Weissman |
SC | 4 |
| 2017 | TripS: automated multi-tiered data placement in a geo-distributed cloud environmentabstractExploiting the cloud storage hierarchy both within and across data-centers of different cloud providers empowers Internet applications to choose data centers (DCs) and storage services based on storage needs. However, using multiple storage services across multiple data centers brings a complex data placement problem that depends on a large number of factors including, e.g., desired goals, storage and network characteristics, and pricing policies. In addition, dynamics e.g., changing user locations and access patterns, make it impossible to determine the best data placement statically. In this paper, we present TripS, a lightweight system that considers both data center locations and storage tiers to determine the data placement for geo-distributed storage systems. Such systems make use of TripS by providing inputs including SLA, consistency model, fault tolerance, latency information, and cost information. With given inputs, TripS models and solves the data placement problem using mixed integer linear programming (MILP) to determine data placement. In addition, to adapt quickly to dynamics, we introduce the notion of Target Locale List (TLL), a pro-active approach to avoid expensive re-evaluation of the optimal placement. The TripS prototype is running on Wiera, a policy driven geo-distributed storage system, to show how a storage system can easily utilize TripS for data placement. We evaluate TripS/Wiera on multiple data centers of AWS and Azure. The results show that TripS/Wiera can reduce cost 14.96% ∼ 98.1% based on workloads in comparison with other works' approaches and can handle both short- and long-term dynamics to avoid SLA violations. Kwangsung Oh, Abhishek Chandra, Jon B. Weissman |
SYSTOR | 3 |
| 2017 | Nebula: Distributed Edge Cloud for Data Intensive ComputingabstractCentralized cloud infrastructures have become the popular platforms for data-intensive computing today. However, they suffer from inefficient data mobility due to the centralization of cloud resources, and hence, are highly unsuited for geo-distributed data-intensive applications where the data may be spread at multiple geographical locations. In this paper, we present Nebula: a dispersed edge cloud infrastructure that explores the use of voluntary resources for both computation and data storage. We describe the lightweight Nebula architecture that enables distributed data-intensive computing through a number of optimization techniques including location-aware data and computation placement, replication, and recovery. We evaluate Nebula performance on an emulated volunteer platform that spans over 50 PlanetLab nodes distributed across Europe, and show how a common data-intensive computing framework, MapReduce, can be easily deployed and run on Nebula. We show Nebula MapReduce is robust to a wide array of failures and substantially outperforms other wide-area versions based on emulated existing systems. Albert Jonathan, Mathew Ryden, Kwangsung Oh, Abhishek Chandra, Jon B. Weissman |
IEEE Trans. Parallel Distributed Syst. | 5 |
| 2016 | Wiera: Towards Flexible Multi-Tiered Geo-Distributed Cloud Storage InstancesabstractGeo-distributed cloud storage systems must tame complexity at many levels: uniform APIs for storage access, supporting flexible storage policies that meet a wide array of application metrics, handling uncertain network dynamics and access dynamism, and operating across many levels of heterogeneity both within and across data-centers. In this paper, we present an integrated solution called Wiera. Wiera extends our earlier cloud storage system, Tiera, that is targeted to multi-tiered policy-based single cloud storage, to the wide-area and multiple data-centers (even across different providers). Wiera enables the specification of global data management policies built on top of local Tiera policies. Such policies enable the user to optimize for cost, performance, reliability, durability, and consistency, both within and across data-centers, and to express their tradeoffs. A key aspect of Wiera is first-class support for dynamism due to network, workload, and access patterns changes. Wiera policies can adapt to changes in user workload, poorly performing data tiers, failures, and changes in user metrics (e.g., cost). Wiera allows unmodified applications to reap the benefits of flexible data/storage policies by externalizing the policy specification. As far as we know, Wiera is the first geo-distributed cloud storage system which handles dynamism actively at run-time. We show how Wiera enables a rich specification of dynamic policies using a concise notation and describe the design and implementation of the system. We have implemented a Wiera prototype on multiple cloud environments, AWS and Azure, that illustrates potential benefits from managing dynamics and in using multiple cloud storage tiers both within and across data-centers. Kwangsung Oh, Abhishek Chandra, Jon B. Weissman |
HPDC | 3 |
| 2016 | Awan: Locality-Aware Resource Manager for Geo-Distributed Data-Intensive ApplicationsabstractToday, many organizations need to operate on data that is distributed around the globe. This is inevitable due to the nature of data that is generated in different locations such as video feeds from distributed cameras, log files from distributed servers, and many others. Although centralized cloud platforms have been widely used for data-intensive applications, such systems are not suitable for processing geo-distributed data due to high data transfer overheads. An alternative approach is to use an Edge Cloud which reduces the network cost of transferring data by distributing its computations globally. While the Edge Cloud is attractive for geo-distributed data-intensive applications, extending existing cluster computing frameworks to a wide-area environment must account for locality. We propose Awan : a new locality-aware resource manager for geo-distributed data-intensive applications. Awan allows resource sharing between multiple computing frameworks while enabling high locality scheduling within each framework. Our experiments with the Nebula Edge Cloud on PlanetLab show that Awan achieves up to a 28% increase in locality scheduling which reduces the average job turnaround time by approximately 18% compared to existing cluster management mechanisms. Albert Jonathan, Abhishek Chandra, Jon B. Weissman |
IC2E | 3 |
| 2016 | Integrating Abstractions to Enhance the Execution of Distributed ApplicationsabstractOne of the factors that limits the scale, performance, and sophistication of distributed applications is the difficulty of concurrently executing them on multiple distributed computing resources. In part, this is due to a poor understanding of the general properties and performance of the coupling between applications and dynamic resources. This paper addresses this issue by integrating abstractions representing distributed applications, resources, and execution processes into a pilot-based middleware. The middleware provides a platform that can specify distributed applications, execute them on multiple resource and for different configurations, and is instrumented to support investigative analysis. We analyzed the execution of distributed applications using experiments that measure the benefits of using multiple resources, the late-binding of scheduling decisions, and the use of backfill scheduling. Matteo Turilli, Zhao Zhang 0007, André Merzky, Michael Wilde, Jon B. Weissman, Daniel S. Katz, Shantenu Jha |
IPDPS | 6 |
| 2016 | End-to-End Optimization for Geo-Distributed MapReduceabstractMapReduce has proven remarkably effective for a wide variety of data-intensive applications, but it was designed to run on large single-site homogeneous clusters. Researchers have begun to explore the extent to which the original MapReduce assumptions can be relaxed, including skewed workloads, iterative applications, and heterogeneous computing environments. This paper continues this exploration by applying MapReduce across geo-distributed data over geo-distributed computation resources. Using Hadoop, we show that network and node heterogeneity and the lack of data locality lead to poor performance, because the interaction of MapReduce phases becomes pronounced in the presence of heterogeneous network behavior. To address these problems, we take a two-pronged approach: We first develop a model-driven optimization that serves as an oracle, providing high-level insights. We then apply these insights to design cross-phase optimization techniques that we implement and demonstrate in a real-world MapReduce implementation. Experimental results in both Amazon EC2 and PlanetLab show the potential of these techniques as performance is improved by 7-18 percent depending on the execution environment and application. Benjamin Heintz, Abhishek Chandra, Ramesh K. Sitaraman, Jon B. Weissman |
IEEE Trans. Cloud Comput. | 4 |
| 2015 | Cloud-Based, User-Centric Mobile Application OptimizationabstractThe abundance of compute and storage resources available in the cloud makes it well-suited to addressing the limitations of mobile devices. We explore the use of cloud infrastructure to optimize content-centric mobile applications, which can have high communication and storage requirements, based on the analysis of user activity. We present two specific optimizations, precaching and prefetching, as well as the design and implementation of a middleware framework that allows mobile application developers to easily utilize these techniques. Our framework is fully generalizable to any content-centric mobile application, a large and growing class of Internet applications. A news aggregation application is used as a case study to evaluate our implementation. We make use of a cosine similarity scheme to identify users with similar interests, which in turn is used to determine what content to prefetch. Various cache algorithms, implemented for our framework, are also considered. A workload trace and simulation are used to measure the performance of the application and framework. We observe a dramatic improvement in application performance due to use of our framework with a reasonable amount of overhead. Our system also significantly outperforms a baseline implementation that performs the same optimizations without taking user activity into account. John Kolb, Prashant Chaudhary, Alexander Schillinger, Abhishek Chandra, Jon B. Weissman |
IC2E | 5 |
| 2015 | Elastic job bundling: an adaptive resource request strategy for large-scale parallel applicationsabstractIn today's batch queue HPC cluster systems, the user submits a job requesting a fixed number of processors. The system will not start the job until all of the requested resources become available simultaneously. When cluster workload is high, large sized jobs will experience long waiting time due to this policy. In this paper, we propose a new approach that dynamically decomposes a large job into smaller ones to reduce waiting time, and lets the application expand across multiple subjobs while continuously achieving progress. This approach has three benefits: (i) application turnaround time is reduced, (ii) system fragmentation is diminished, and (iii) fairness is promoted. Our approach does not depend on job queue time prediction but exploits available backfill opportunities. Simulation results have shown that our approach can reduce application mean turnaround time by up to 48%. Jon B. Weissman |
SC | 2 |
| 2014 | Nebula: Distributed Edge Cloud for Data Intensive ComputingabstractCentralized cloud infrastructures have become the de-facto platform for data-intensive computing today. However, they suffer from inefficient data mobility due to the centralization of cloud resources, and hence, are highly unsuited for dispersed-data-intensive applications, where the data may be spread at multiple geographical locations. In this paper, we present Nebula: a dispersed cloud infrastructure that uses voluntary edge resources for both computation and data storage. We describe the lightweight Nebula architecture that enables distributed data-intensive computing through a number of optimizations including location-aware data and computation placement, replication, and recovery. We evaluate Nebula's performance on an emulated volunteer platform that spans over 50 PlanetLab nodes distributed across Europe, and show how a common data-intensive computing framework, MapReduce, can be easily deployed and run on Nebula. We show Nebula MapReduce is robust to a wide array of failures and substantially outperforms other wide-area versions based on a BOINC like model. Mathew Ryden, Kwangsung Oh, Abhishek Chandra, Jon B. Weissman |
IC2E | 4 |
| 2014 | Tiera: towards flexible multi-tiered cloud storage instancesabstractCloud providers offer an array of storage services that represent different points along the performance, cost, and durability spectrum. If an application desires the composite benefits of multiple storage tiers, then it must manage the complexity of different interfaces to these storage services and their diverse policies. We believe that it is possible to provide the benefits of customized tiered cloud storage to applications without compromising simplicity using a lightweight middleware. In this paper, we introduce Tiera, a middleware that enables the provision of multi-tiered cloud storage instances that are easy to specify, flexible, and enable a rich array of storage policies and desired metrics to be realized. Tiera's novelty lies in the first-class support for encapsulated tiered cloud storage, ease of programmability of data management policies, and support for runtime replacement and addition of policies and tiers. Tiera enables an application to realize a desired metric (e.g., low latency or low cost) by selecting different storage services that constitute a Tiera instance, and easily specifying a policy, using event and response pairs, to manage the life cycle of data stored in the instance. We illustrate the benefits of Tiera through a prototype implemented on the Amazon cloud. By deploying unmodified MySQL database engine and a TPC-W Web bookstore application on Tiera, we are able to improve their respective throughputs by 47% -- 125% and 46% -- 69%, over standard deployments. We further show the flexibility of Tiera in achieving different desired application metrics with minimal overhead. Ajaykrishna Raghavan, Abhishek Chandra, Jon B. Weissman |
Middleware | 3 |
| 2014 | A Security-enabled Grid System for MINDS Distributed Data Mining
Jinoh Kim, Jon B. Weissman |
J. Grid Comput. | 3 |
| 2013 | Cross-Phase Optimization in MapReduceabstractMap Reduce has been designed to accommodate large-scale data-intensive workloads running on large single-site homogeneous clusters. Researchers have begun to explore the extent to which the original Map Reduce assumptions can be relaxed including skewed workloads, iterative applications, and heterogeneous computing environments. Our work continues this exploration by applying Map Reduce across widely distributed data over distributed computation resources. This problem arises when datasets are generated at multiple sites as is common in many scientific domains and increasingly e-commerce applications. It also occurs when multi-site resources such as geographically separated data centers are applied to the same Map Reduce job. Using Hadoop, we show that the absence of network and node homogeneity and locality of data lead to poor performance. The problem is that interaction of Map Reduce phases becomes pronounced in the presence of heterogeneous network behavior. In this paper, we propose new cross-phase optimization techniques that enable independent Map Reduce phases to influence one another. We propose techniques that optimize the push and map phases to enable push-map overlap and to allow map behavior to feed back into push dynamics. Similarly, we propose techniques that optimize the map and reduce phases to enable shuffle cost to feed back and affect map scheduling decisions. We evaluate the benefits of our techniques in both Amazon EC2 and Planet Lab. The experimental results show the potential of these techniques as performance is improved from 7%-18% depending on the execution environment and application. Benjamin Heintz, Abhishek Chandra, Jon B. Weissman |
IC2E | 4 |
| 2013 | Distributed computing practice for large-scale science and engineering applicationsabstractSUMMARY It is generally accepted that the ability to develop large‐scale distributed applications has lagged seriously behind other developments in cyberinfrastructure. In this paper, we provide insight into how such applications have been developed and an understanding of why developing applications for distributed infrastructure is hard. Our approach is unique in the sense that it is centered around half a dozen existing scientific applications; we posit that these scientific applications are representative of the characteristics, requirements, as well as the challenges of the bulk of current distributed applications on production cyberinfrastructure (such as the US TeraGrid). We provide a novel and comprehensive analysis of such distributed scientific applications. Specifically, we survey existing models and methods for large‐scale distributed applications and identify commonalities, recurring structures, patterns and abstractions. We find that there are many ad hoc solutions employed to develop and execute distributed applications, which result in a lack of generality and the inability of distributed applications to be extensible and independent of infrastructure details. In our analysis, we introduce the notion of application vectors: a novel way of understanding the structure of distributed applications. Important contributions of this paper include identifying patterns that are derived from a wide range of real distributed applications, as well as an integrated approach to analyzing applications, programming systems and patterns, resulting in the ability to provide a critical assessment of the current practice of developing, deploying and executing distributed applications. Gaps and omissions in the state of the art are identified, and directions for future research are outlined. Copyright © 2012 John Wiley & Sons, Ltd. Shantenu Jha, Murray Cole, Daniel S. Katz, Manish Parashar, Omer F. Rana, Jon B. Weissman |
Concurr. Comput. Pract. Exp. | 6 |
| 2012 | Sharing-Aware Cloud-Based Mobile OutsourcingabstractMobile devices, such as smart phones and tablets, are becoming the universal interface to online services and applications. However, such devices have limited computational power and battery life, which limits their ability to execute resource-intensive applications. Computation outsourcing to external resources has been proposed as a technique to alleviate this problem. Most existing work on mobile outsourcing has focused on either single application optimization or outsourcing to fixed, local resources, with the assumption that wide-area latency is prohibitively high. However, the opportunity of improving the outsourcing performance by utilizing the relation among multiple applications and optimizing the server provisioning is neglected. In this paper, we present the design and implementation of an Android/Amazon EC2-based mobile application outsourcing framework, leveraging the cloud for scalability, elasticity, and multi-user code/data sharing. Using this framework, we empirically demonstrate that the cloud is not only feasible but desirable as an offloading platform for latency-tolerant applications. We have proposed to use data mining techniques to detect data sharing across multiple applications, and developed novel scheduling algorithms that exploit such data sharing for better outsourcing performance. Additionally, our platform is designed to dynamically scale to support a large number of mobile users concurrently. Experiments show that our proposed techniques and algorithms substantially improve application performance, while achieving high efficiency in terms of computation resource and network usage. Chonglei Mei, Abhishek Chandra, Jon B. Weissman |
IEEE CLOUD | 5 |
| 2012 | Reducing Data Transfer in Service-Oriented Architectures: The Circulate ApproachabstractAs the number of services and the size of data involved in workflows increases, centralized orchestration techniques are reaching the limits of scalability. When relying on web services without third-party data transfer, a standard orchestration model needs to pass all data through a centralized engine, which results in unnecessary data transfer and the engine to become a bottleneck to the execution of a workflow. As a solution, this paper presents and evaluates Circulate, an alternative service-oriented architecture which facilitates an orchestration model of central control in combination with a choreography model of optimized distributed data transport. Extensive performance analysis through the PlanetLab framework is conducted on a web service-based implementation over a range of Internet-scale configurations which mirror scientific workflow environments. Performance analysis concludes that our architecture's optimized model of data transport speeds up the execution time of workflows, consistently outperforms standard orchestration and scales with data and node size. Furthermore, Circulate is a less-intrusive solution as individual services do not have to be reconfigured in order to take part in a workflow. Adam Barker, Jon B. Weissman, Jano I. van Hemert |
IEEE Trans. Serv. Comput. | 2 |
| 2011 | ViDeDup: An Application-Aware Framework for Video De-duplication
Atul Katiyar, Jon B. Weissman |
HotStorage | 2 |
| 2011 | A Scheduling and Certification Algorithm for Defeating Collusion in Desktop GridsabstractBy exploiting idle time on volunteer machines, desktop grids provide a way to execute large sets of tasks with negligible maintenance and low cost. Although desktop grids are attractive for their scalability and low cost, relying on external resources may compromise the correctness of application execution due to the well-known unreliability of nodes. In this paper, we consider a very challenging threat model: correlated errors caused either by organized groups of cheaters that may collude to produce incorrect results, or by buggy or so-called "unofficial" clients. By using a previously described on-line algorithm for detecting collusion and characterizing the participant behaviors, we propose a scheduling and result certification algorithm that tackles collusion. Using several real-life traces, we show that our approach minimizes both replication overhead and the number of incorrectly certified results. Louis-Claude Canon, Emmanuel Jeannot, Jon B. Weissman |
ICDCS | 3 |
| 2011 | Passive Network Performance Estimation for Large-Scale, Data-Intensive ComputingabstractDistributed computing applications are increasingly utilizing distributed data sources. However, the unpredictable cost of data access in large-scale computing infrastructures can lead to severe performance bottlenecks. Providing predictability in data access is, thus, essential to accommodate the large set of newly emerging large-scale, data-intensive computing applications. In this regard, accurate estimation of network performance is crucial to meeting the performance goals of such applications. Passive estimation based on past measurements is attractive for its relatively small overhead compared to relying on explicit probing. In this paper, we take a passive approach for network performance estimation. Our approach is different from existing passive techniques that rely either on past direct measurements of pairs of nodes or on topological similarities. Instead, we exploit secondhand measurements collected by other nodes without any topological restrictions. In this paper, we present Overlay Passive Estimation of Network performance (OPEN), a scalable framework providing end-to-end network performance estimation based on secondhand measurements, and discuss how OPEN achieves cost-effective estimation in a large-scale infrastructure. Our extensive experimental results show that OPEN estimation can be applicable for replica and resource selections commonly used in distributed computing. Jinoh Kim, Abhishek Chandra, Jon B. Weissman |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2010 | A dynamic approach for characterizing collusion in desktop gridsabstractBy exploiting idle time on volunteer machines, desktop grids provide a way to execute large sets of tasks with negligible maintenance and low cost. Although desktop grids are attractive for cost-conscious projects, relying on external resources may compromise the correctness of application execution due to the well-known unreliability of nodes. In this paper, we consider the most challenging threat model: organized groups of cheaters that may collude to produce incorrect results. We propose two on-line algorithms for detecting collusion and characterizing the participant behaviors. Using several real-life traces, we show that our approach is accurate and efficient in identifying collusion and in estimating group behavior. Louis-Claude Canon, Emmanuel Jeannot, Jon B. Weissman |
IPDPS | 3 |
| 2010 | Scheduling Multisource Divisible Loads on Arbitrary NetworksabstractScheduling multisource divisible loads is a challenging task as different sources should cooperate and share their computing power with others to balance their loads and minimize total computational time. In this study, we attempt to address a generalized divisible load scheduling problem for handling loads from multiple sources on arbitrary networks. This problem is all the more challenging as 1) the topology is arbitrary, 2) in such networks, it is difficult to decide from which source and which route a processing node should receive loads, and 3) processing nodes must be allocated to different sources when they become available. We study two distinct cases of interest, static case and dynamic case, and propose two novel strategies, referred to as static scheduling strategy (SSS) and dynamic scheduling strategy (DSS), respectively. Both strategies work in an iterative fashion. In each iteration, they will use a novel graph partitioning (GP) scheme to partition the network such that each source in the network gains a portion of network resources and then these sources cooperate to process their loads. We analyze the performance of DSS using queuing theory and derive upper bounds on a load's average waiting time and a source's average queue length. We use simulation to verify the usefulness and effectiveness of SSS and DSS. Our findings reveal an interesting ¿load insensitive¿ propertyof SSS and also verify the theoretical upper bound of average queue length at each source in the dynamic case. Jingxi Jia, Bharadwaj Veeravalli, Jon B. Weissman |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2009 | Introduction
Jon B. Weissman, Lex Wolters, David Abramson 0001, Marty Humphrey |
Euro-Par | 1 |
| 2009 | Adaptive middleware supporting scalable performance for high-end network services
Byoung-Dai Lee, Jon B. Weissman, Young-Kwang Nam |
J. Netw. Comput. Appl. | 2 |
| 2009 | Using Data Accessibility for Resource Selection in Large-Scale Distributed SystemsabstractLarge-scale distributed systems provide an attractive scalable infrastructure for network applications. However, the loosely coupled nature of this environment can make data access unpredictable, and in the limit, unavailable. We introduce the notion of accessibility to capture both availability and performance. An increasing number of data-intensive applications require not only considerations of node computation power but also accessibility for adequate job allocations. For instance, selecting a node with intolerably slow connections can offset any benefit to running on a fast node. In this paper, we present accessibility-aware resource selection techniques by which it is possible to choose nodes that will have efficient data access to remote data sources. We show that the local data access observations collected from a node's neighbors are sufficient to characterize accessibility for that node. By conducting trace-based, synthetic experiments on PlanetLab, we show that the resource selection heuristics guided by this principle significantly outperform conventional techniques such as latency-based or random allocations. The suggested techniques are also shown to be stable even under churn despite the loss of prior observations. Jinoh Kim, Abhishek Chandra, Jon B. Weissman |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2008 | Orchestrating Data-Centric WorkflowsabstractWhen orchestrating data-centric workflows as are commonly found in the sciences, centralised servers can become a bottleneck to the performance of a workflow; output from service invocations are normally transferred via a centralised orchestration engine, when they should be passed directly to where they are needed at the next service in the workflow. To address this performance bottleneck, this paper presents a lightweight hybrid workflow architecture and concrete API, based on a centralised control flow, distributed data flow model. Our architecture maintains the robustness and simplicity of centralised orchestration, but facilitates choreography by allowing services to exchange data directly with one another, reducing data that needs to be transferred through a centralised server. Furthermore our architecture is standards compliment, flexible and is a non-disruptive solution; service definitions do not have to be altered prior to enactment. Adam Barker, Jon B. Weissman, Jano I. van Hemert |
CCGRID | 2 |
| 2008 | Eliminating the middleman: peer-to-peer dataflowabstractEfficiently executing large-scale, data-intensive workflows such as Montage must take into account the volume and pattern of communication. When orchestrating data-centric workflows, centralised servers common to standard workflow systems can become a bottleneck to performance. However, standards-based workflow systems that rely on centralisation, e.g., Web service based frameworks, have many other benefits such as a wide user base and sustained support. Adam Barker, Jon B. Weissman, Jano I. van Hemert |
HPDC | 2 |
| 2008 | Accessibility-Based Resource Selection in Loosely-Coupled Distributed SystemsabstractLarge-scale distributed systems provide an attractive scalable infrastructure for network applications. However,the loosely-coupled nature of this environment can make data access unpredictable, and in the limit, unavailable. We introduce the notion of accessibility to capture both availability and performance. An increasing number of data intensive applications require not only considerations of node computation power but also accessibility for adequate job allocations. For instance, selecting a node with intolerably slow connections can offset any benefit to running on a fast node. In this paper, we present accessibility-aware resource selection techniques by which it is possible to choose nodes that will have efficient data access to remote data sources. We show that the local data access observations collected from a node's neighbors are sufficient to characterize accessibility for that node. We then present resource selection heuristics guided by this principle, and show that they significantly out perform standard techniques. The suggested techniques are also shown to be stable even under churn despite the loss of prior observations. Jinoh Kim, Abhishek Chandra, Jon B. Weissman |
ICDCS | 3 |
| 2008 | Exploring the throughput-fairness tradeoff of deadline scheduling in heterogeneous computing environmentsabstractThe scalability and computing power of large-scale computational platforms has made them attractive for hosting compute-intensive time-critical applications. Many of these applications are composed of computational tasks that require specific deadlines to be met for successful completion. In this paper, we show that combining redundant scheduling with deadline-based scheduling in these systems leads to a fundamental tradeoff between throughput and fairness. We propose a new scheduling algorithm called Limited Resource Earliest Deadline (LRED) that couples redundant scheduling with deadline-driven scheduling in a flexible way by using a simple tunable parameter to exploit this tradeoff. Our evaluation of LRED shows that LRED provides a powerful mechanism to achieve desired throughput or fairness under high loads and low timeliness environments. Vasumathi Sundaram, Abhishek Chandra, Jon B. Weissman |
SIGMETRICS | 3 |
| 2007 | Exploiting Heterogeneity for Collective Data Downloading in Volunteer-based NetworksabstractScientific computing is being increasingly deployed over volunteer-based distributed computing environments consisting of idle resources on donated user machines. A fundamental challenge in these environments is the dissemination of data to the computation nodes, with the successful completion of jobs being driven by the efficiency of collective data download across compute nodes, and not only the individual download times. This paper considers the use of a data network consisting of data distributed across a set of data servers, and focuses on the server selection problem: how do individual nodes select a server for downloading data to minimize the communication makespan - the maximal download time for a data file. Through experiments conducted on a pastry network running on PlanetLab, we demonstrate that nodes in a volunteer-based network are heterogeneous in terms of several metrics, such as bandwidth, load, and capacity, which impact their download behavior. We propose new server selection heuristics that incorporate these metrics, and demonstrate that these heuristics outperform traditional proximity-based server selection, reducing average makespans by at least 30%. We further show that incorporating information about download concurrency avoids overloading servers, and improves performance by about 17-43% over heuristics considering only proximity and bandwidth. Jinoh Kim, Abhishek Chandra, Jon B. Weissman |
CCGRID | 3 |
| 2007 | Ridge: combining reliability and performance in open grid platformsabstractLarge-scale donation-based distributed infrastructures need to cope with the inherent unreliability of participant nodes. A widely-used work scheduling technique in such environments is to redundantly schedule the out sourced computations to a number of nodes. We present the design and implementation of RIDGE, a reliability aware system which uses a node's prior performance and behavior to make more effective scheduling decisions. We have implemented RIDGE on top of the BOINC distributed computing infrastructure and have evaluated its performance on a live test bed consisting of 120 PlanetLab nodes. Our experimental results show that RIDGE is able to match or surpass the throughput of the best vanilla BOINC configuration under different reliability environments, by automatically adapting to the characteristics of the underlying environment. In addition, RIDGE is able to provide much lower work unit makes pans compared to BOINC, which indicates its desirability in service-oriented environments with time constraints. Krishnaveni Budati, Jason D. Sonnek, Abhishek Chandra, Jon B. Weissman |
HPDC | 4 |
| 2007 | NGS: Service Adaptation in Open Grid PlatformsabstractLarge-scale donation-based distributed infrastructures need to cope with the inherent unreliability of participant nodes. A widely-used work scheduling technique in such environments is to redundantly schedule the outsourced computations to a number of nodes. We present the design and implementation of RIDGE, a reliability-aware system which uses a node's prior performance and behavior to make more effective scheduling decisions. We have implemented RIDGE on top of the BOINC distributed computing infrastructure and have evaluated its performance on a live PlanetLab testbed. Our experimental results show that RIDGE is able to match or surpass the throughput of the best BOINC configuration by automatically adapting to the characteristics of the underlying environment. In addition, RIDGE is able to provide much lower workunit makespans compared to BOINC. RIDGE is also able to produce significantly lower communication makespans for downloading clients. Collectively, the results suggest that RIDGE has great promise for service-oriented environments with time constraints. Krishnaveni Budati, Jinoh Kim, Abhishek Chandra, Jon B. Weissman |
IPDPS | 4 |
| 2007 | A Robust Spanning Tree Topology for Data Collection and Dissemination in Distributed EnvironmentsabstractLarge-scale distributed applications are subject to frequent disruptions due to resource contention and failure. Such disruptions are inherently unpredictable and, therefore, robustness is a desirable property for the distributed operating environment. In this work, we describe and evaluate a robust topology for applications that operate on a spanning tree overlay network. Unlike previous work that is adaptive or reactive in nature, we take a proactive approach to robustness. The topology itself is able to simultaneously withstand disturbances and exhibit good performance. We present both centralized and distributed algorithms to construct the topology, and then demonstrate its effectiveness through analysis and simulation of two classes of distributed applications: Data collection in sensor networks and data dissemination in divisible load scheduling. The results show that our robust spanning trees achieve a desirable trade-off for two opposing metrics where traditional forms of spanning trees do not. In particular, the trees generated by our algorithms exhibit both resilience to data loss and low power consumption for sensor networks. When used as the overlay network for divisible load scheduling, they display both robustness to link congestion and low values for the makespan of the schedule Darin England, Bharadwaj Veeravalli, Jon B. Weissman |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2007 | Adaptive Reputation-Based Scheduling on Unreliable Distributed InfrastructuresabstractThis paper addresses the inherent unreliability and instability of worker nodes in large-scale donation-based distributed infrastructures such as peer-to-peer and grid systems. We present adaptive scheduling techniques that can mitigate this uncertainty and significantly outperform current approaches. In this work, we consider nodes that execute tasks via donated computational resources and may behave erratically or maliciously. We present a model in which reliability is not a binary property, but a statistical one based on a node's prior performance and behavior. We use this model to construct several reputation-based scheduling algorithms that employ estimated reliability ratings of worker nodes for efficient task allocation. Our scheduling algorithms are designed to adapt to changing system conditions, as well as nonstationary node reliability. Through simulation, we demonstrate that our algorithms can significantly improve throughput while maintaining a very high success rate of task completion. Our results suggest that reputation-based scheduling can handle a wide variety of worker populations, including nonstationary behavior, with overhead that scales well with system size. We also show that our adaptation mechanism allows the application designer fine-grain control over the desired performance metrics. Jason D. Sonnek, Abhishek Chandra, Jon B. Weissman |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2006 | Reputation-Based Scheduling on Unreliable Distributed InfrastructuresabstractThis paper presents a design and analysis of scheduling techniques to cope with the inherent unreliability and instability of worker nodes in large-scale donation-based distributed infrastructures such as P2P and Grid systems. In particular, we focus on nodes that execute tasks via donated computational resources and may behave erratically or maliciously. We present a model in which reliability is not a binary property but a statistical one based on a node’s prior performance and behavior. We use this model to construct several reputation-based scheduling algorithms that employ estimated reliability ratings of worker nodes for efficient task allocation. Through simulation of a BOINC-like distributed computing infrastructure, we demonstrate that our algorithms can significantly improve throughput, while maintaining a very high success rate of task completion. Jason D. Sonnek, Mukesh Nathan, Abhishek Chandra, Jon B. Weissman |
ICDCS | 4 |
| 2005 | Supporting the dynamic grid service lifecycleabstractThis paper presents an architecture and implementation for a dynamic OGSA-based grid service architecture that extends GT3 to support dynamic service hosting - where to host and re-host a service within the grid in response to service demand and resource fluctuation. Our model goes beyond current OGSI implementations in which the service is presumed to be pre-installed at all sites (and only service instantiation is dynamic). In dynamic virtual organizations (VOs), we believe dynamic service hosting provides an important flexibility. Our model also defines several new adaptive grid service classes that support adaptation at multiple levels. Dynamic service deployment allows new services to be added or replaced without taking down a site for reconfiguration and allows a VO to respond effectively to dynamic resource availability and demand. The preliminary results suggest that the cost of dynamic installation, deployment, and invocation, is tolerable. Jon B. Weissman, Darin England |
CCGRID | 1 |
| 2005 | A new metric for robustness with application to job schedulingabstractScheduling strategies for parallel and distributed computing have mostly been oriented toward performance, while striving to achieve some notion of fairness. With the increase in size, complexity, and heterogeneity of today's computing environments, we argue that, in addition to performance metrics, scheduling algorithms should be designed for robustness. That is, they should have the ability to maintain performance under a wide variety of operating conditions. Although robustness is easy to define, there are no widely used metrics for this property. To this end, we present a methodology for characterizing and measuring the robustness of a system to a specific disturbance. The methodology is easily applied to many types of computing systems and it does not require sophisticated mathematical models. To illustrate its use, we show three applications of our technique to job scheduling; one supporting a previous result with respect to backfilling, one examining overload control in a streaming video server, and one comparing two different scheduling strategies for a distributed network service. The last example also demonstrates how consideration of robustness leads to better system design as we were able to devise a new and effective scheduling heuristic. Darin England, Jon B. Weissman, Jayashree Sadagopan |
HPDC | 2 |
| 2004 | A Genetic Algorithm Based Approach for Scheduling Decomposable Data Grid ApplicationsabstractData grid technology promises geographically distributed scientists to access and share physically distributed resources such as compute resource, networks, storage, and most importantly data collections for large-scale data intensive problems. Because of the massive size and distributed nature of these datasets, scheduling data grid applications must consider communication and computation simultaneously to achieve high performance. In many data grid applications, data can be decomposed into multiple independent sub datasets and distributed for parallel execution and analysis. We exploit this property and propose a novel genetic algorithm based approach that automatically decomposes data onto communication and computation resources. The proposed GA-based scheduler takes advantage of the parallelism of decomposable data grid applications to achieve the desired performance level. We evaluate the proposed approach comparing with other algorithms. Simulation results show that the proposed GA-based approach can be a competitive choice for scheduling large data grid applications in terms of both scheduling overhead and the relative solution quality as compared to other algorithms. Jon B. Weissman |
ICPP | 2 |
| 2004 | Costs and Benefits of Load Sharing in the Computational Grid
Darin England, Jon B. Weissman |
JSSPP | 2 |
| 2003 | Adaptive Resource Selection for Grid-Enabled Network ServicesabstractDue to the popularity of high-speed networks and advances in packaging and interface technologies, there has been significant efforts for providing high performance applications as network services that can be accessed remotely across the network, thus promoting sharing of both software and hardware. For high-demand network services, in particular, it will often be the case that the network services are installed at multiple sites so that each participating site can handle parts of client requests. We label such services as grid-enabled network services. In this paper, we present two adaptive site selection heuristics that do not depend on accurate predictions of completion times of service requests: weight queue length based heuristic and multi-level queue based selection. Byoung-Dai Lee, Jon B. Weissman |
NCA | 2 |
| 2003 | Integrated scheduling: the best of both worlds
Jon B. Weissman, Lakshman Rao Abburi, Darin England |
J. Parallel Distributed Comput. | 1 |
| 2003 | Guest editor introduction: special issue on Computational Grids
Jon B. Weissman, Richard Wolski |
J. Parallel Distributed Comput. | 1 |
| 2002 | Community Services: A Toolkit for Rapid Deployment of Network ServicesabstractAdvances in packaging and interface technologies have made it possible for software components to be shared across the network through encapsulation and offered as network services. They allow end-users to focus on their applications and obtain remote services when needed simply by invoking them across the network. Many groups have built significant infrastructures for providing domain-specific high performance services. However, transforming high performance applications into network services is labor intensive and time consuming because there is little existing infrastructure to utilize. In this paper, we propose a software toolkit and runtime infrastructure for rapid deployment of network services. Byoung-Dai Lee, Jon B. Weissman |
CLUSTER | 2 |
| 2002 | The Virtual Service Grid: an architecture for delivering high-end network servicesabstractAbstract This paper presents the design of a new system architecture, Virtual Service Grid (VSG), for delivering high‐performance network services. The VSG is based on the concept of the virtual service which provides location, replication, and fault transparency to clients accessing remotely deployed high‐end services. One of the novel features of the virtual service is the ability to self‐scale in response to client demand. The VSG exploits network and service information to make adaptive dynamic replica selection, creation, and deletion decisions. We describe the VSG architecture, middleware, and replica management policies. We have deployed the VSG on a wide‐area Internet testbed to evaluate its performance. The results indicate that the VSG can deliver efficient performance for a wide range of client workloads, both in terms of reduced response time and in the utilization of system resources. Copyright © 2002 John Wiley & Sons, Ltd. Jon B. Weissman, Byoung-Dai Lee |
Concurr. Comput. Pract. Exp. | 1 |
| 2002 | Predicting the Cost and Benefit of Adapting Data Parallel Applications in Clusters
Jon B. Weissman |
J. Parallel Distributed Comput. | 1 |
| 2001 | Applying Grid Technologies to BioinformaticsabstractThe science of bioinformatics provides researchers with the tools necessary to unravel the mysteries of life and evolution, discover cures for disease, and control the evolution of living organisms. To assist researchers in managing the growing data processing and management demands associated with bioinformatics, we have created a production system that draws upon Grid based technologies to control several aspects of the process. We briefly discuss system architecture, results, and future directions of the project. Michael Karo, Christopher Dwan, John L. Freeman, Ernest F. Retzel, Jon B. Weissman, Miron Livny |
HPDC | 5 |
| 2001 | Dynamic Replica Management in the Service GridabstractAs the Internet is evolving away from providing simple connectivity towards providing more sophisticated services, it is difficult to provide efficient delivery of high-demand services to end users, due to the dynamic sharing of the network and connected servers. To address this problem, we propose the service grid architecture that incorporates dynamic replication and deletion of services. Byoung-Dai Lee, Jon B. Weissman |
HPDC | 2 |
| 2001 | Optimizing Remote File Access for Parallel and Distributed Network Applications
Jon B. Weissman, Mahesh K. Marina, Michael Gingras |
J. Parallel Distributed Comput. | 1 |
| 1999 | Fault Tolerant Computing on the Grid: What are My Options?abstractHigh-performance distributed computing across wide-area networks has become an active topic of research. Achieving large-scale distributed computing in a seamless manner introduces a number of difficult problems. This paper examines one of the most critical problems, fault tolerance. We have examined fault tolerance options for a common class of high-performance parallel applications, single-program-multiple-data (SPMD). Performance models for two fault tolerance methods, checkpoint-recovery (CR) and wide-area replication (WR), have been developed. These models enable quantitative comparisons of the two methods as applied to SPMD applications. Jon B. Weissman |
HPDC | 1 |
| 1999 | Prophet: automated scheduling of SPMD programs in workstation networksabstractObtaining efficient execution of parallel programs in workstation networks is a difficult problem for the user. Unlike dedicated parallel computer resources, network resources are shared, heterogeneous, vary in availability, and offer communication performance that is still an order of magnitude slower than parallel computer interconnection networks. Prophet, a system that automatically schedules data parallel SPMD programs in workstation networks for the user, has been developed. Prophet uses application and resource information to select the appropriate type and number of workstations, divide the application into component tasks and data across these workstations, and assign tasks to workstations. This system has been integrated into the Mentat parallel processing system developed at the University of Virginia. A suite of scientific Mentat applications has been scheduled using Prophet on a heterogeneous workstation network. The results are promising and demonstrate that scheduling SPMD applications can be automated with good performance. Copyright © 1999 John Wiley & Sons, Ltd. Jon B. Weissman |
Concurr. Pract. Exp. | 1 |
| 1998 | Metascheduling: A Scheduling Model for Metacomputing SystemsabstractMetacomputing is the seamless application of geographically separated distributed computing resources to user applications. We consider the scheduling of meta applications; applications consisting of multiple components that may communicate and interact over the course of the application. Components may be schedulable computations, remote servers or databases, remote instruments, humans in the loop, etc. We divide applications into three categories-concurrent, parallel, and pipeline. Concurrent is the classic meta application in which a set of components each running in a single site are executing concurrently and exchanging data. Parallel is a special case of concurrent in which a component is replicated and distributed across multiple sites. Pipeline applications consist of components connected in a chain like fashion. Jon B. Weissman |
HPDC | 1 |
| 1998 | Gallop: The Benefits of Wide-Area Computing for Parallel Processing
Jon B. Weissman |
J. Parallel Distributed Comput. | 1 |
| 1997 | Run-time Support for Scheduling Parallel Applications in Heterogeneous NOWsabstractThis paper describes the current state of Prophet-a system that provides run-time scheduling support for parallel applications in heterogeneous workstation networks. Prior work on Prophet demonstrated that scheduling SPMD applications could be effectively automated with excellent performance. Enhancements have been made to Prophet to broaden its use to other application types including parallel pipelines, and to make more effective use of dynamic system state information to further improve performance. The results indicate that both SPMD and parallel pipeline applications can be scheduled to produce reduced completion time by exploiting the application structure and run-time information. Jon B. Weissman |
HPDC | 1 |
| 1996 | A Federated Model for Scheduling in Wide-Area SystemsabstractA model for scheduling in wide area systems is described. The model is federated and utilizes a collection of local site schedulers that control the use of their resources. The wide area scheduler consults the local site schedulers to obtain candidate machine schedules. A set of issues and challenges inherent to wide area scheduling are also described and the proposed model is shown to address many of these problems. A distributed algorithm for wide area scheduling is presented and relies upon information made available about the resource needs of user jobs. The wide area scheduler will be implemented in Legion, a wide area computing system developed at the University of Virginia. Jon B. Weissman, Andrew S. Grimshaw |
HPDC | 1 |
| 1996 | Portable Run-Time Support for Dynamic Object-Oriented Parallel ProcessingabstractMentat is an object-oriented parallel processing system designed to simplify the task of writing portable parallel programs for parallel machines and workstation networks. The Mentat compiler and run-time system work together to automatically manage the communication and synchronization between objects. The run-time system marshals member function arguments, schedules objects on processors, and dynamically constructs and executes large-grain data dependence graphs. In this article we present the Mentat run-time system. We focus on three aspects—the software architecture, including the interface to the compiler and the structure and interaction of the principle components of the run-time system; the run-time overhead on a component-by-component basis for two platforms, a Sun SparcStation 2 and an Intel Paragon; and an analysis of the minimum granularity required for application programs to overcome the run-time overhead. Andrew S. Grimshaw, Jon B. Weissman, W. Timothy Strayer |
ACM Trans. Comput. Syst. | 2 |
| 1995 | A framework for partitioning parallel computations in heterogeneous environmentsabstractAbstract In the paper we present a framework for partitioning data parallel computations across a heterogeneous metasystem at runtime. The framework is guided by program and resource information which is made available to the system. Three difficult problems are handled by the framework: processor selection, task placement and heterogeneous data domain decomposition. Solving each of these problems contributes to reduced elapsed time. In particular, processor selection determines the best grain size at which to run the computation, task placement reduces communication cost, and data domain decomposition achieves processor load balance. We present results which indicate that excellent performance is achievable using the framework. The paper extends our earlier work on partitioning data parallel computations across a single‐level network of heterogeneous workstations. Jon B. Weissman, Andrew S. Grimshaw |
Concurr. Pract. Exp. | 1 |
| 1994 | Network Partitioning of Data Parallel ComputationsabstractPartitioning data parallel computations across a network of heterogeneous workstations is a difficult problem for the user. We have developed a runtime partitioning method for choosing the number and type of processors to apply to a data parallel computation, and a decomposition of the data domain in order to achieve reduced completion time. The partitioning method utilizes information about the problem in the form of callback functions and uses a set of topology-specific communication functions to estimate communication costs. We show that the method makes effective partitioning decisions and has runtime overhead that is easily tolerated. In particular we show that for two implementations of a canonical stencil application, minimum elapsed times are obtained for a range of problem sizes on a network of heterogeneous workstations.> Jon B. Weissman, Andrew S. Grimshaw |
HPDC | 1 |
| 1994 | Metasystems: An Approach Combining Parallel Processing and Heterogeneous Distributed Computing Systems
Andrew S. Grimshaw, Jon B. Weissman, Emily A. West, Edmond C. Loyot Jr. |
J. Parallel Distributed Comput. | 2 |