EDBT 2026 Demo / reviewers in the wild / expert
Henri Casanova
dblp:25/6606
· DBLP profile ↗
111ranked-venue papers
27as first author
13since 2021 · last 2025
0000-0001-6310-0365ORCID · corroborated
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 98 · 24 first-author · 8 since 2021Software engineering, systems software and programming languages · 7 · 1 first-author · 3 since 2021Applied, interdisciplinary, general and emerging computing · 5 · 1 first-author · 3 since 2021Computer networks · 1Theory of computation · 1
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2025 | Lowering entry barriers to developing custom simulators of distributed applications and platforms with SimGrid
Henri Casanova, Arnaud Giersch, Arnaud Legrand, Martin Quinson, Frédéric Suter |
Parallel Comput. | 1 |
| 2024 | Automated Calibration of a Simulator of MPI Application ExecutionsabstractThe traditional approach for assessing the performance of scientific applications on HPC platforms consists in executing these applications on these platforms. But conducting these real-world experiments comes with several difficulties. Besides being often time-, labor-, and resource-intensive, experiments are limited to application and platform configurations at hand, thus precluding the exploration of "what if?" scenarios. A way to resolve these difficulties is to resort to simulation. The main concern, then, is that of simulation accuracy. For a simulation to be accurate, the parameters that define the behavior of the simulation models can be calibrated with respect to ground-truth executions. Simulation calibration, in the current state of the art, relies, at best, on labor-intensive manual procedures. We propose an automated simulation calibration approach, and apply this approach to the specific context of the simulation of MPI applications on leadership class HPC platforms. This poster will motivate the development of this approach and detail our methodology and results. Yick Ching Wong, Frédéric Suter, Kshitij Mehta, Henri Casanova, Jesse McDonald |
e-Science | 4 |
| 2024 | Simulation of Large-Scale HPC Storage Systems: Challenges and MethodologiesabstractAs the scale of production HPC platforms increases, so does the computing and I/O performance gap, exacerbating the storage bottleneck. High-performance storage systems have been developed to alleviate this bottleneck, but many questions remain concerning their architecture, implementation, and configuration. Answering these questions via experimental campaigns proves arduous. First, some answers are required before deploying the system. Second, once a system hits production the experimental scope is limited by the system's specific configuration and by constraints of production use. In this work we identify challenges posed by the design and validation of a storage simulator. We then propose solutions implemented in Fives, a simulator of HPC workloads on platforms that comprise a parallel file system. We show how our simulator can be instantiated and calibrated for the accurate simulation of a production Lustre deployment. Julien Monniot, Francois Tessier, Henri Casanova, Gabriel Antoniu |
HiPC | 3 |
| 2024 | An exploration of online-simulation-driven portfolio scheduling in Workflow Management Systems
Jesse McDonald, John Dobbs, Yick Ching Wong, Rafael Ferreira da Silva, Henri Casanova |
Future Gener. Comput. Syst. | 5 |
| 2023 | WfCommons: Data Collection and Runtime Experiments using Multiple Workflow SystemsabstractScientific workflows have become ubiquitous across scientific fields, and their execution methods and systems continue to be the subject of research and development. Most experimental evaluations of these workflows rely on workflow instances, which can be either real-world or synthetic, to ensure relevance to current application domains or explore hypothetical/future scenarios. The WfCommons project addresses this need by providing data and tools to support such evaluations. In this paper, we present an overview of WfCommons and describe two recent developments. Firstly, we introduce a workflow execution "tracer" for Nextflow, which significantly enhances the set of real-world instances available in WfCommons. Secondly, we describe a workflow instance "translator" that enables the execution of any real-world or synthetic WfCommons workflow instance using Dask. Our contributions aim to provide researchers and practitioners with more comprehensive resources for evaluating scientific workflows. Henri Casanova, Kyle Berney, Serge Chastel, Rafael Ferreira da Silva |
COMPSAC | 1 |
| 2023 | Automated generation of scientific workflow generators with WfChef
Tainã Coleman, Henri Casanova, Rafael Ferreira da Silva |
Future Gener. Comput. Syst. | 2 |
| 2022 | On the Feasibility of Simulation-Driven Portfolio Scheduling for Cyberinfrastructure Runtime Systems
Henri Casanova, Yick Ching Wong, Loïc Pottier, Rafael Ferreira da Silva |
JSSPP | 1 |
| 2022 | WfCommons: A framework for enabling scientific workflow research and development
Tainã Coleman, Henri Casanova, Loïc Pottier, Manav Kaushik, Ewa Deelman, Rafael Ferreira da Silva |
Future Gener. Comput. Syst. | 2 |
| 2022 | Beyond Binary Search: Parallel In-Place Construction of Implicit Search Tree LayoutsabstractWe present parallel algorithms to efficiently permute a sorted array into the level-order binary search tree (BST), level-order B-tree (B-tree), and van Emde Boas (vEB) layoutsin-place. We analytically determine the complexity of our algorithms and empirically measure their performance. When considering the total time to permute the data in-place and to perform a series of search queries, the vEB layout provides the best performance on the CPU. Given an input of$N$N=537 million 64-bit integers, the benefits of query performance (compared to binary search) outweigh the cost of in-place permutation when performing as few as 0.37% of$N$Nqueries. On the GPU, results depend on the particular architecture, with the B-tree and vEB layouts performing the best. The number of queries necessary to reach the break-even point with binary search ranges from 1.3% to 8.9% of$N$N=1,074 million 32-bit integers. Kyle Berney, Henri Casanova, Benjamin Karsin, Nodari Sitchinava |
IEEE Trans. Computers | 2 |
| 2021 | Modeling the Linux page cache for accurate simulation of data-intensive applicationsabstractThe emergence of Big Data in recent years has resulted in a growing need for efficient data processing solutions. While infrastructures with sufficient compute power are available, the I/O bottleneck remains. The Linux page cache is an efficient approach to reduce I/O overheads, but few experimental studies of its interactions with Big Data applications exist, partly due to limitations of real-world experiments. Simulation is a popular approach to address these issues, however, existing simulation frameworks do not simulate page caching fully, or even at all. As a result, simulation-based performance studies of data-intensive applications can lead to misleading results and inaccurate conclusions.In this paper, we propose an I/O simulation model that captures the key features of the Linux page cache. We have implemented this model as part of the WRENCH workflow simulation framework, which itself builds on the popular Sim-Grid distributed systems simulation framework. Our model and its implementation enable the simulation of both single-threaded and multithreaded applications, and of both writeback and writethrough caches for local or network-based filesystems. We evaluate the accuracy of our model in different conditions, including sequential and concurrent applications, as well as local and remote I/Os. We find that our page cache model reduces the simulation error by up to an order of magnitude when compared to state-of-the-art, cacheless simulations. Our model is publicly available in the WRENCH framework, making it usable in a wide range of simulation studies. Hoang-Dung Do, Valérie Hayot-Sasson, Rafael Ferreira da Silva, Christopher Steele, Henri Casanova, Tristan Glatard |
CLUSTER | 5 |
| 2021 | WfChef: Automated Generation of Accurate Scientific Workflow GeneratorsabstractScientific workflow applications have become mainstream and their automated and efficient execution on large-scale compute platforms is the object of extensive research and development. For these efforts to be successful, a solid experimental methodology is needed to evaluate workflow algorithms and systems. A foundation for this methodology is the availability of realistic workflow instances. Dozens of workflow instances for a few scientific applications are available in public repositories. While these are invaluable, they are limited: workflow instances are not available for all application scales of interest. To address this limitation, previous work has developed generators of synthetic, but representative, workflow instances of arbitrary scales. These generators are popular, but implementing them is a manual, labor-intensive process that requires expert application knowledge. As a result, these generators only target a handful of applications, even though hundreds of applications use workflows in production.In this work, we present WfChef, a framework that fully automates the process of constructing a synthetic workflow generator for any scientific application. Based on an input set of workflow instances, WfChef automatically produces a synthetic workflow generator. We define and evaluate several metrics for quantifying the realism of the generated workflows. Using these metrics, we compare the realism of the workflows generated by WfChef generators to that of the workflows generated by the previously available, hand-crafted generators. We find that the WfChef generators not only require zero development effort (because it is automatically produced), but also generate workflows that are more realistic than those generated by hand-crafted generators. Tainã Coleman, Henri Casanova, Rafael Ferreira da Silva |
e-Science | 2 |
| 2021 | GLUME: A Strategy for Reducing Workflow Execution Times on Batch-Scheduled Platforms
Evan Hataishi, Pierre-François Dutot, Rafael Ferreira da Silva, Henri Casanova |
JSSPP | 4 |
| 2021 | Teaching parallel and distributed computing concepts in simulation with WRENCH
Henri Casanova, Ryan Tanaka, William Koch, Rafael Ferreira da Silva |
J. Parallel Distributed Comput. | 1 |
| 2020 | Modeling the Performance of Scientific Workflow Executions on HPC Platforms with Burst BuffersabstractScientific domains ranging from bioinformatics to astronomy and earth science rely on traditional high-performance computing (HPC) codes, often encapsulated in scientific workflows. In contrast to traditional HPC codes that employ a few programming and runtime approaches that are highly optimized for HPC platforms, scientific workflows are not necessarily optimized for these platforms. As an effort to reduce the gap between compute and I/O performance, HPC platforms have adopted intermediate storage layers known as burst buffers. A burst buffer (BB) is a fast storage layer positioned between the global parallel file system and the compute nodes. Two designs currently exist: (i) shared, where the BBs are located on dedicated nodes; and (ii) on-node, in which each compute node embeds a private BB. In this paper, using accurate simulations and realworld experiments, we study how to best use these new storage layers when executing scientific workflows. These applications are not necessarily optimized to run on HPC systems, and thus can exhibit I/O patterns that differ from that of HPC codes. Thus, we first characterize the I/O behaviors of a real-world workflow under different configuration scenarios on two leadership-class HPC systems (Cori at NERSC and Summit at ORNL). Then, we use these characterizations to calibrate a simulator for workflow executions on HPC systems featuring shared and private BBs. Last, we evaluate our approach against a large I/O-intensive workflow, and we provide insights on the performance levels and the potential limitations of these two BBs architectures. Loïc Pottier, Rafael Ferreira da Silva, Henri Casanova, Ewa Deelman |
CLUSTER | 3 |
| 2020 | Developing accurate and scalable simulators of production workflow management systems with WRENCH
Henri Casanova, Rafael Ferreira da Silva, Ryan Tanaka, Suraj Pandey, Gautam Jethwani, William Koch, Spencer Albrecht, James Oeth, Frédéric Suter |
Future Gener. Comput. Syst. | 1 |
| 2019 | Sparse 3-D NoCs with Inductive CouplingabstractWireless interconnects based on inductive coupling technology are compelling propositions for designing 3-D integrated chips. This work addresses the heat dissipation problem on such systems. Although effective cooling technologies have been proposed for systems designed based on Through Silicon Via (TSV), their application to systems that use inductive coupling is problematic because of increased wireless-communication distance. For this reason, we propose two methods for designing sparse 3-D chips layouts and Networks on Chip (NoCs) based on inductive coupling. The first method computes an optimized 3-D chip layout and then generates a randomized network topology for this layout. The second method uses a standard stack chip layout with a standard network topology as a starting point, and then deterministically transforms it into either a "staircase" or a "checkerboard" layout. We quantitatively compare the designs produced by these two methods in terms of network and application performance. Our main finding is that the first method produces designs that ultimately lead to higher parallel application performance, as demonstrated for nine OpenMP applications in the NAS Parallel Benchmarks. Michihiro Koibuchi, Lambert T. Leong, Tomohiro Totoki, Naoya Niwa, Hiroki Matsutani, Hideharu Amano, Henri Casanova |
DAC | 7 |
| 2019 | Bridging Concepts and Practice in eScience via Simulation-Driven EngineeringabstractThe CyberInfrastructure (CI) has been the object of intensive research and development in the last decade, resulting in a rich set of abstractions and interoperable software implementations that are used in production today for supporting ongoing and breakthrough scientific discoveries. A key challenge is the development of tools and application execution frameworks that are robust in current and emerging CI configurations, and that can anticipate the needs of upcoming CI applications. This paper presents WRENCH, a framework that enables simulation-driven engineering for evaluating and developing CI application execution frameworks. WRENCH provides a set of high-level simulation abstractions that serve as building blocks for developing custom simulators. These abstractions rely on the scalable and accurate simulation models that are provided by the SimGrid simulation framework. Consequently, WRENCH makes it possible to build, with minimum software development effort, simulators that that can accurately and scalably simulate a wide spectrum of large and complex CI scenarios. These simulators can then be used to evaluate and/or compare alternate platform, system, and algorithm designs, so as to drive the development of CI solutions for current and emerging applications. Rafael Ferreira da Silva, Henri Casanova, Ryan Tanaka, Frédéric Suter |
eScience | 2 |
| 2018 | Analysis-driven Engineering of Comparison-based Sorting Algorithms on GPUsabstractWe study the relationship between memory accesses, bank conflicts, thread multiplicity (also known as over-subscription) and instruction-level parallelism in comparison-based sorting algorithms for Graphics Processing Units (GPUs). We experimentally validate a proposed formula that relates these parameters with asymptotic analysis of the number of memory accesses by an algorithm. Using this formula we analyze and compare several GPU sorting algorithms, identifying key performance bottlenecks in each one of them. Based on this analysis we propose a GPU-efficient multiway merge-sort algorithm, GPU-MMS, which minimizes or eliminates these bottlenecks and balances various limiting factors for specific hardware. Benjamin Karsin, Volker Weichert, Henri Casanova, John Iacono, Nodari Sitchinava |
ICS | 3 |
| 2018 | Beyond Binary Search: Parallel In-Place Construction of Implicit Search Tree LayoutsabstractWe present parallel algorithms to efficiently permute a sorted array into the level-order binary search tree (BST), level-order B-tree (B-tree), and van Emde Boas (vEB) layouts in-place. We analytically determine the complexity of our algorithms and empirically measure their performance. Results indicate that on both CPU and GPU architectures B-tree layouts provide the best query performance. However, when considering the total time to permute the data and to perform a series of search queries, our vEB permutation provides the best performance on the CPU. We show that, given an input of N=500M 64-bit integers, the benefits of query performance (compared to binary search) outweigh the cost of in-place permutation using our algorithms when performing at least 5M queries (1% of N) and 27M queries (6% of N), on our CPU and GPU platforms, respectively. Kyle Berney, Henri Casanova, Alyssa Higuchi, Benjamin Karsin, Nodari Sitchinava |
IPDPS | 2 |
| 2018 | Computing the expected makespan of task graphs in the presence of silent errors
Henri Casanova, Julien Herrmann, Yves Robert |
Parallel Comput. | 1 |
| 2018 | Checkpointing Workflows for Fail-Stop ErrorsabstractWe consider the problem of orchestrating the execution of workflow applications structured as Directed Acyclic Graphs (DAGs) on parallel computing platforms that are subject to fail-stop failures. The objective is to minimize expected overall execution time, or makespan. A solution to this problem consists of a schedule of the workflow tasks on the available processors and of a decision of which application data to checkpoint to stable storage, so as to mitigate the impact of processor failures. To address this challenge, we consider a restricted class of graphs, Minimal Series-Parallel Graphs (M-SPGS), which is relevant to many real-world workflow applications. For this class of graphs, we propose a recursive list-scheduling algorithm that exploits the M-SPG structure to assign sub-graphs to individual processors, and uses dynamic programming to decide how to checkpoint these sub-graphs. We assess the performance of our algorithm for production workflow configurations, comparing it to an approach in which all application data is checkpointed and an approach in which no application data is checkpointed. Results demonstrate that our algorithm outperforms both the former approach, because of lower checkpointing overhead, and the latter approach, because of better resilience to failures. Li Han 0001, Louis-Claude Canon, Henri Casanova, Yves Robert, Frédéric Vivien |
IEEE Trans. Computers | 3 |
| 2017 | An Efficient Algorithm for the 1D Total Visibility-Index ProblemabstractLet T be a terrain, and let P be a set of points (locations) on its surface. An important problem in Geographic Information Science (GIS) is computing the visibility index of a point p on P, that is, the number of points in P that are visible from p. The total visibility-index problem asks for computing the visibility index of every point in P. Most applications of this problem involve 2-dimensional terrains represented by a grid of n × n square cells, where each cell is associated with an elevation value, and P consists of the center-points of these cells. Current approaches for computing the total visibility-index on such a terrain take at least quadratic time with respect to the number of the terrain cells. While finding a subquadratic solution to this 2D total visibility-index problem is an open problem, surprisingly, no subquadratic solution has been proposed for the one-dimensional (1D) version of the problem; in the 1D problem, the terrain is an x-monotone polyline, and P is the set of the polyline vertices. We present an O(n log2 n) algorithm that solves the 1D total visibility-index problem in the RAM model. Our algorithm is based on a geometric dualization technique, which reduces the problem into a set of instances of the red-blue line segment intersection counting problem. We also present a parallel version of this algorithm, which requires O(log2 n) time and O(n log2 n) work in the CREW PRAM model. We implement a naive O(n2) approach and three variations of our algorithm: one employing an existing red-blue line segment intersection algorithm and two new approaches that perform the intersection counting by leveraging features specific to our problem. We present experimental results for both serial and parallel implementations on large synthetic and real-world datasets, using two distinct hardware platforms. Results show that all variants of our algorithm outperform the naive approach by several orders of magnitude on large datasets. Furthermore, we show that our new intersection counting implementations achieve more than 8 times speedup over the existing red-blue line segment intersection algorithm. Our parallel implementation is able to process a terrain of 224 vertices in under 1 minute using 16 cores, achieving more than 7 times speedup over serial execution. Peyman Afshani, Mark de Berg, Henri Casanova, Benjamin Karsin, Colin Lambrechts, Nodari Sitchinava, Constantinos Tsirogiannis |
ALENEX | 3 |
| 2017 | Checkpointing Workflows for Fail-Stop ErrorsabstractWe consider the problem of orchestrating the execution of workflow applications structured as Directed Acyclic Graphs (DAGs) on parallel computing platforms that are subject to fail-stop failures. The objective is to minimize expected overall execution time, or makespan. A solution to this problem consists of a schedule of the workflow tasks on the available processors and of a decision of which application data to checkpoint to stable storage, so as to mitigate the impact of processor failures. For general DAGs this problem is hopelessly intractable. In fact, given a solution, computing its expected makespan is still a difficult problem. To address this challenge, we consider a restricted class of graphs, Minimal Series-Parallel Graphs (M-SPGS). It turns out that many real-world workflow applications are naturally structured as M-SPGS. For this class of graphs, we propose a recursive list-scheduling algorithm that exploits the M-SPG structure to assign sub-graphs to individual processors, and uses dynamic programming to decide which tasks in these sub-gaphs should be checkpointed. Furthermore, it is possible to efficiently compute the expected makespan for the solution produced by this algorithm, using a first-order approximation of task weights and existing evaluation algorithms for 2-state probabilistic DAGs. We assess the performance of our algorithm for production workflow configurations, comparing it to (i) an approach in which all application data is checkpointed, which corresponds to the standard way in which most production workflows are executed today; and (ii) an approach in which no application data is checkpointed. Our results demonstrate that our algorithm strikes a good compromise between these two approaches, leading to lower checkpointing overhead than the former and to better resilience to failure than the latter. Li Han 0001, Louis-Claude Canon, Henri Casanova, Yves Robert, Frédéric Vivien |
CLUSTER | 3 |
| 2017 | A Case for Uni-directional Network Topologies in Large-Scale ClustersabstractDesigning low-latency network topologies of switches is a key objective for next-generation large-scale clusters. Low latency is preconditioned on low hop counts, but existing network topologies have hop counts much larger than theoretical lower bounds. To alleviate this problem, we propose building network topologies based on uni-directional graphs that are known to have hop counts close to theoretical lower bounds. A practical difficulty with uni-directional topologies is switch-by-switch flow control, which we resolve by using hot-potato routing. Cycle-accurate network simulation experiments for various traffic patterns on uni-directional topologies show that hot-potato routing achieves performance comparable to that of conventional deadlock-free routing. Similar experiments are used to compare several uni-directional topologies to bi-directional topologies, showing that the former achieve significantly lower latency and higher throughput. We quantify end-to-end application performance for parallel application benchmarks via discrete-even simulation, showing that uni-directional topologies can lead to large application performance improvements over their bi-directional counterparts. Finally, we discuss practical issues for uni-directional topologies such as cabling complexity and cost, power consumption, and soft-error tolerance. Our results make a compelling case for considering uni-directional topologies for upcoming large-scale clusters. Michihiro Koibuchi, Tomohiro Totoki, Hiroki Matsutani, Hideharu Amano, Fabien Chaix, Ikki Fujiwara, Henri Casanova |
CLUSTER | 7 |
| 2017 | High-Bandwidth Low-Latency Approximate Interconnection NetworksabstractComputational applications are subject to various kinds of numerical errors, ranging from deterministic roundoff errors to soft errors caused by non-deterministic bit flips, which do not lead to application failure but corrupt application results. Non-deterministic bit flips are typically mitigated in hardware using various error correcting codes (ECC). But in practice, due to performance and cost concerns, these techniques do not guarantee error-free execution. On large-scale computing platforms, soft errors occur with non-negligible probability in RAM and on the CPU, and it has become clear that applications must tolerate them. For some applications, this tolerance is intrinsic as result quality can remain acceptable even in the presence of soft errors (e.g., data analysis applications, multimedia applications). Tolerance can also be built into the application, resolving data corruptions in software during application execution. By contrast, today's optical networks hold on to a rigid error-free standard, which imposes limits on network performance scalability. In this work we propose high-bandwidth, low-latency approximate networks with the following three features: (1) Optical links that exploit multi-level quadrature amplitude modulation (QAM) for achieving high bandwidth; (2) Avoidance of forward error correction (FEC), which makes optical link error-prone but affords lower latency; and (3) The use of symbol mapping coding between bit sequence and QAM to ensure data integrity that is sufficient for practical soft-error-tolerant applications. Discrete-event simulation results for application benchmarks show that approx networks achieve speedups up to 2.94 when compared to conventional networks. Daichi Fujiki, Kiyo Ishii, Ikki Fujiwara, Hiroki Matsutani, Hideharu Amano, Henri Casanova, Michihiro Koibuchi |
HPCA | 6 |
| 2016 | Distance Threshold Similarity Searches: Efficient Trajectory Indexing on the GPUabstractApplications in many domains perform searches over datasets that contain moving object trajectories. A common class of searches are similarity searches that attempt to identify trajectories with similar characteristics. In this work, we focus on the distance threshold similarity search that finds all trajectories within a given distance of a query trajectory over a time interval. This search involves large numbers of Euclidean moving distance calculations, thus making it a good candidate for execution on manycore platforms such as GPUs. However, low search response time is preconditioned on efficient indexing of trajectory data. We propose three indexing schemes designed for the GPU, with spatial, temporal and spatiotemporal selectivity. These schemes differ significantly from traditional tree-based indexing schemes that have been previously proposed for CPU executions. We evaluate implementations of our proposed indexing schemes using two synthetic and one real-world astrophysics dataset, showing under which conditions each scheme achieves high performance. Our broad finding is that a GPU implementation, provided an appropriate indexing scheme is used, can outperform a multithreaded CPU implementation that uses a state-of-the-art index tree. In particular, the performance improvement is large for regimes that are relevant for classes of real-world applications, thereby demonstrating that the GPU is an attractive platform for searching and processing moving object trajectories. Michael G. Gowanlock, Henri Casanova |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2015 | Efficient Batched Predecessor Search in Shared Memory on GPUsabstractMany-core Graphics Processing Units (GPUs) are being used for general-purpose computing. However, due to architectural features, for many problems it is challenging to design parallel algorithms that exploit the full compute power of GPUs. Among these features is the memory design. Although the issue of coalesced global memory access has been documented and studied extensively, another important architectural feature is the organization of shared memory into banks. The study of how bank conflicts impact algorithm performance has only recently begun to receive attention. In this work we study the predecessor search algorithm and the effects of bank conflicts on its execution time. Via complexity analysis we show that bank conflicts cause significant loss in parallelism for a naive algorithm. We then propose two improved algorithms: one that eliminates bank conflicts altogether but that uses a work inefficient linear search, and one that is work-optimal but that experiences a limited number of bank conflicts. We develop GPU implementations of these algorithms and present experimental results obtained on real-world hardware. These results validate our theoretical analysis of the naive algorithm and allow us to assess the performance of our algorithms in practice. Although both our improved algorithms outperform the naive algorithm, our main experimental finding is that our conflict-limited algorithm provides a larger performance gain. Benjamin Karsin, Henri Casanova, Nodari Sitchinava |
HiPC | 2 |
| 2015 | Augmenting low-latency HPC network with free-space optical linksabstractVarious network topologies can be used for deploying High Performance Computing (HPC) clusters. The network topology, which connects switches In cabinets on a machine room floor, is typically defined once and for all at system deployment time. For a diverse application workload, there are downsides to having a single wired topology. In this work, we propose using free-space optics (FSO) in large-scale systems so that a diverse application workload can be better supported. A high-density layout of FSO terminals on top of the cabinets is determined that allows line-of-sight communication between arbitrary cabinet pairs. We first show that our proposal reduces both end-to-end network latency and total cable length when compared to a wired topology. We then demonstrate that the use of FSO links improves the embedding/partitioning capabilities of a wired topology. More specifically, we show that a recently proposed random low-latency topology can be augmented with a reasonable number of FSO links to support multiple k-ary n-cube and fat tree embedded topologies. Finally, we investigate power-aware on/off link regulation techniques and show how adding/reconfiguring FSO links leads to both performance and power efficiency improvements. Ikki Fujiwara, Michihiro Koibuchi, Tomoya Ozaki, Hiroki Matsutani, Henri Casanova |
HPCA | 5 |
| 2015 | Indexing of Spatiotemporal Trajectories for Efficient Distance Threshold Similarity Searches on the GPUabstractApplications in many domains search moving object trajectory databases. The distance threshold search finds all trajectories within a given distance of a query trajectory. We develop three GPU distance threshold search implementations that use indexing techniques significantly different from those used in CPU implementations. We determine experimentally under which conditions each approach performs well using one real-world astrophysics dataset and two synthetic datasets. Overall, we find that the GPU is an attractive technology for a broad range of relevant trajectory database scenarios. Michael G. Gowanlock, Henri Casanova |
IPDPS | 2 |
| 2015 | Simulation of MPI applications with time-independent tracesabstractSummary Analyzing and understanding the performance behavior of parallel applications on parallel computing platforms is a long‐standing concern in the High Performance Computing community. When the targeted platforms are not available, simulation is a reasonable approach to obtain objective performance indicators and explore various hypothetical scenarios. In the context of applications implemented with the Message Passing Interface, two simulation methods have been proposed, on‐line simulation and off‐line simulation, both with their own drawbacks and advantages. In this work, we present an off‐line simulation framework, that is, one that simulates the execution of an application based on event traces obtained from an actual execution. The main novelty of this work, when compared to previously proposed off‐line simulators, is that traces that drive the simulation can be acquired on large, distributed, heterogeneous, and non‐dedicated platforms. As a result, the scalability of trace acquisition is increased, which is achieved by enforcing that traces contain no time‐related information. Moreover, our framework is based on a state‐of‐the‐art scalable, fast, and validated simulation kernel. We introduce the notion of performing off‐line simulation from time‐independent traces, propose and evaluate several trace acquisition strategies, describe our simulation framework, and assess its quality in terms of trace acquisition scalability, simulation accuracy, and simulation time. Copyright © 2014 John Wiley & Sons, Ltd. Henri Casanova, Frédéric Desprez, George S. Markomanolis, Frédéric Suter |
Concurr. Comput. Pract. Exp. | 1 |
| 2015 | On the impact of process replication on executions of large-scale parallel applications with coordinated checkpointing
Henri Casanova, Yves Robert, Frédéric Vivien, Dounia Zaidouni |
Future Gener. Comput. Syst. | 1 |
| 2015 | Selecting linear algebra kernel composition using response time predictionabstractSummary Numerical linear algebra libraries provide many kernels that can be composed to perform complex computations. For a given computation, there is typically a large number of functionally equivalent kernel compositions. Some of these compositions achieve better response times than others for particular data and when executed on a particular computer architecture. Previous research provides methods to enumerate (a subset of) these kernel compositions. In this work, we study the problem of determining the composition that yields the lowest response time. Our approach is based on a response time prediction for each candidate combination. While this prediction could in principle be obtained using analytical and/or empirical performance models, developing accurate such models is known to be challenging. Instead, we define a feature space that captures salient properties of kernel combinations and predict response time using supervised machine learning. We experiment with a standard set of machine learning algorithms and identify an effective algorithm for our kernel composition selection problem. Using this algorithm, our approach widely outperforms the strategy that would consist in always using the simplest kernel composition and is often close to the fastest kernel compositions among those evaluated. We quantify the potential benefit of our approach if it were to be implemented as part of an interactive computational tool. We find that although the potential benefit is substantial, a limiting factor is the kernel composition enumeration overhead. Copyright © 2014 John Wiley & Sons, Ltd. Aurélie Hurault, Kyungim Baek, Henri Casanova |
Softw. Pract. Exp. | 3 |
| 2015 | Swap-And-Randomize: A Method for Building Low-Latency HPC InterconnectsabstractRandom network topologies have been proposed to create low-diameter, low-latency interconnection networks in large-scale computing systems. However, these topologies are difficult to deploy in practice, especially when re-designing existing systems, because they lead to increased total cable length and cable packaging complexity. In this work we propose a new method for creating random topologies without increasing cable length: randomly swap link endpoints in a non-random topology that is already deployed across several cabinets in a machine room. We quantitatively evaluate topologies created in this manner using both graph analysis and cycle-accurate network simulation, including comparisons with non-random topologies and previously-proposed random topologies. Ikki Fujiwara, Michihiro Koibuchi, Hiroki Matsutani, Henri Casanova |
IEEE Trans. Parallel Distributed Syst. | 4 |
| 2014 | Distance threshold similarity searches on spatiotemporal trajectories using GPGPUabstractThe processing of moving object trajectories arises in many application domains. We focus on a trajectory similarity search, the distance threshold search, which finds all trajectories within a given distance of a query trajectory over a time interval. A multithreaded CPU implementation that makes use of an in-memory R-tree index can achieve high parallel efficiency. We propose a GPGPU implementation that avoids index-trees altogether and instead features a GPU-friendly indexing scheme. We show that our GPU implementation compares well to the CPU implementation. One interesting question is that of creating efficient query batches (so as to reduce both memory pressure and computation cost on the GPU). We design algorithms for creating such batches, and we find that using fixed-size batches is sufficient in practice. We develop an empirical response time model that can be used to pick a good batch size. Michael G. Gowanlock, Henri Casanova |
HiPC | 2 |
| 2014 | Cost-Optimal Execution of Boolean Query Trees with Shared StreamsabstractThe processing of queries expressed as trees of boolean operators applied to predicates on sensor data streams has several applications in mobile computing. Sensor data must be retrieved from the sensors, which incurs a cost, e.g., an energy expense that depletes the battery of a mobile query processing device. The objective is to determine the order in which predicates should be evaluated so as to shortcut part of the query evaluation and minimize the expected cost. This problem has been studied assuming that each data stream occurs at a single predicate. In this work we remove this assumption since it does not necessarily hold in practice. Our main results are an optimal algorithm for single-level trees and a proof of NP-completeness for DNF trees. For DNF trees, however, we show that there is an optimal predicate evaluation order that corresponds to a depth-first traversal. This result provides inspiration for a class of heuristics. We show that one of these heuristics largely outperforms other sensible heuristics, including a heuristic proposed in previous work. Henri Casanova, Lipyeow Lim, Yves Robert, Frédéric Vivien, Dounia Zaidouni |
IPDPS | 1 |
| 2014 | Skywalk: A Topology for HPC Networks with Low-Delay SwitchesabstractWith low-delay switches on the horizon, end-to-end latency in large-scale High Performance Computing (HPC) interconnects will be dominated by cable delays. In this context we define a new network topology, Skywalk, for deploying low-latency interconnects in upcoming HPC systems. Skywalk uses randomness to achieve low latency, but does so in a way that accounts for the physical layout of the topology so as to lead to further cable length and thus latency reductions. Via graph analysis and discrete-event simulation we show that Skywalk compares favorably (in terms of latency, cable length, and throughput) to traditional low-degree torus and moderate-degree hypercube topologies, to high-degree fully-connected Dragonfly topologies, to the HyperX topology, and to recently proposed fully random topologies. Ikki Fujiwara, Michihiro Koibuchi, Hiroki Matsutani, Henri Casanova |
IPDPS | 4 |
| 2014 | Versatile, scalable, and accurate simulation of distributed applications and platforms
Henri Casanova, Arnaud Giersch, Arnaud Legrand, Martin Quinson, Frédéric Suter |
J. Parallel Distributed Comput. | 1 |
| 2013 | Layout-conscious random topologies for HPC off-chip interconnectsabstractAs the scales of parallel applications and platforms increase the negative impact of communication latencies on performance becomes large. Random network topologies can be used to achieve low hop counts between nodes and thus low latency. However, random topologies lead to increased aggregate cable length and cable packaging complexity on a machine room floor. In this work we propose two new methods for generating random topologies and their physical layout on a floorplan: randomize links after optimizing the physical layout, or optimize the layout after randomizing links. The first method randomly swaps link endpoints for a given non-random topology for which a good physical layout is known. The resulting topology has the same cable length and cable packaging as the original topology, but achieves lower communication latency. The second method creates a random topology with random links picked so that they will not lead to a long physical cable length, and then solves a constrained optimization problem to compute a physical layout that minimizes aggregate cable length. We quantitatively compare these two methods using both graph analysis and cycle-accurate network simulation, including comparisons with previously proposed random topologies and non-random topologies. Michihiro Koibuchi, Ikki Fujiwara, Hiroki Matsutani, Henri Casanova |
HPCA | 4 |
| 2013 | Mapping Tightly-Coupled Applications on Volatile ResourcesabstractPlatforms that comprise volatile processors, such as desktop grids, have been traditionally used for executing independent-task applications. In this work we study the scheduling of tightly-coupled iterative master-worker applications onto volatile processors. The main challenge is that workers must be simultaneously available for the application to make progress. We consider two additional complications: one should take into account that workers can become temporarily reclaimed and, for data-intensive applications, one should account for the limited bandwidth between the master and the workers. In this context, our first contribution is a theoretical study of the scheduling problem in its off-line version, i.e., when processor availability is known in advance. Even in this case the problem is NP-hard. Our second contribution is an analytical approximation of the expectation of the time needed by a set of workers to complete a set of tasks and of the probability of success of this computation. This approximation relies on a Markovian assumption for the temporal availability of processors. Our third contribution is a set of heuristics, some of which use the above approximation to favor reliable processors in a sensible manner. We evaluate these heuristics in simulation. We identify some heuristics that significantly outperform their competitors and derive heuristic design guidelines. Henri Casanova, Fanny Dufossé, Yves Robert, Frédéric Vivien |
PDP | 1 |
| 2012 | Virtual Machine Resource Allocation for Service Hosting on Heterogeneous Distributed PlatformsabstractWe propose algorithms for allocating multiple resources to competing services running in virtual machines on heterogeneous distributed platforms. We develop a theoretical problem formulation and compare these algorithms via simulation experiments based in part on workload data supplied by Google. Our main finding is that vector packing approaches proposed in the homogeneous case can be extended to provide high-quality solutions in the heterogeneous case, and combined to provide a single efficient algorithm. We also consider the case when there may be bounded errors in estimates of performance-related resource needs. We provide a heuristic for compensating for such errors that performs well in simulation, as well as a proof of the worst-case competitive ratio for the single-resource, single-node case when there is no bound on the error. Mark Stillwell, Frédéric Vivien, Henri Casanova |
IPDPS | 3 |
| 2012 | A case for random shortcut topologies for HPC interconnectsabstractAs the scales of parallel applications and platforms increase the negative impact of communication latencies on performance becomes large. Fortunately, modern High Performance Computing (HPC) systems can exploit low-latency topologies of high-radix switches. In this context, we propose the use of random shortcut topologies, which are generated by augmenting classical topologies with random links. Using graph analysis we find that these topologies, when compared to non-random topologies of the same degree, lead to drastically reduced diameter and average shortest path length. The best results are obtained when adding random links to a ring topology, meaning that good random shortcut topologies can easily be generated for arbitrary numbers of switches. Using flit-level discrete event simulation we find that random shortcut topologies achieve throughput comparable to and latency lower than that of existing non-random topologies such as hypercubes and tori. Finally, we discuss and quantify practical challenges for random shortcut topologies, including routing scalability and larger physical cable lengths. Michihiro Koibuchi, Hiroki Matsutani, Hideharu Amano, D. Frank Hsu, Henri Casanova |
ISCA | 5 |
| 2012 | Cabinet Layout Optimization of Supercomputer Topologies for Shorter Cable LengthabstractAs the scales of supercomputers increase total cable length becomes enormous, e.g., up to thousands of kilometers. Recent high-radix switches with dozens of ports make switch layout and system packaging more complex. In this study, we study the optimization of the physical layout of topologies of switches on a machine room floor with the goal of reducing cable length. For a given topology, using graph clustering algorithms, we group switches logically into cabinets so that the number of inter-cabinet cables is small. Then, we map the cabinets onto a physical floor space so as to minimize total cable length. This is done by modeling and optimizing the mapping problem as a facility location problem. Our evaluation results show that, when compared to standard clustering/mapping approaches and for popular network topologies, our clustering approach can reduce the number of inter-cabinet cables by up to 40.3% and our mapping approach can reduce the inter-rack cable length by up to 39.6%. Ikki Fujiwara, Michihiro Koibuchi, Henri Casanova |
PDCAT | 3 |
| 2012 | Energy-aware service allocation
Damien Borgetto, Henri Casanova, Georges Da Costa, Jean-Marc Pierson |
Future Gener. Comput. Syst. | 2 |
| 2012 | Dynamic Fractional Resource Scheduling versus Batch SchedulingabstractWe propose a novel job scheduling approach for homogeneous cluster computing platforms. Its key feature is the use of virtual machine technology to share fractional node resources in a precise and controlled manner. Other VM-based scheduling approaches have focused primarily on technical issues or extensions to existing batch scheduling systems, while we take a more aggressive approach and seek to find heuristics that maximize an objective metric correlated with job performance. We derive absolute performance bounds and develop algorithms for the online nonclairvoyant version of our scheduling problem. We further evaluate these algorithms in simulation against both synthetic and real-world HPC workloads and compare our algorithms to standard batch scheduling approaches. We find that our approach improves over batch scheduling by orders of magnitude in terms of job stretch, while leading to comparable or better resource utilization. Our results demonstrate that virtualization technology coupled with lightweight online scheduling strategies can afford dramatic improvements in performance for executing HPC workloads. Mark Stillwell, Frédéric Vivien, Henri Casanova |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2011 | On the Utility of DVFS for Power-Aware Job Placement in Clusters
Jean-Marc Pierson, Henri Casanova |
Euro-Par (1) | 2 |
| 2011 | Scheduling Parallel Iterative Applications on Volatile ResourcesabstractIn this paper we study the execution of iterative applications on volatile processors such as those found on desktop grids. We develop master-worker scheduling schemes that attempt to achieve good trade-offs between worker speed and worker availability. A key feature of our approach is that we consider a communication model where the bandwidth capacity of the master for sending application data to workers is limited. This limitation makes the scheduling problem more difficult both in a theoretical sense and in a practical sense. Furthermore, we consider that a processor can be in one of three states: available, down, or temporarily preempted by its owner. This preempted state also complicates the scheduling problem. In practical settings, e.g., desktop grids, master bandwidth is limited and processors are temporarily reclaimed. Consequently, addressing the aforementioned difficulties is necessary for successfully deploying master-worker applications on volatile platforms. Our first contribution is to determine the complexity of the scheduling problem in its off-line version, i.e., when processor availability behaviors are known in advance. Even with this knowledge, the problem is NP-hard, and cannot be approximated within a factor $8/7$. Our second contribution is a closed-form formula for the expectation of the time needed by a worker to complete a set of tasks. This formula relies on a Markovian assumption for the temporal availability of processors, and is at the heart of some heuristics that aim at favoring "reliable'' processors in a sensible manner. Our third contribution is a set of heuristics, which we evaluate in simulation. Our results provide guidance to selecting the best strategy as a function of processor state availability versus average task duration. Henri Casanova, Fanny Dufossé, Yves Robert, Frédéric Vivien |
IPDPS | 1 |
| 2011 | Single Node On-Line Simulation of MPI Applications with SMPIabstractSimulation is a popular approach for predicting the performance of MPI applications for platforms that are not at one's disposal. It is also a way to teach the principles of parallel programming and high-performance computing to students without access to a parallel computer. In this work we present SMPI, a simulator for MPI applications that uses on-line simulation, i.e., the application is executed but part of the execution takes place within a simulation component. SMPI simulations account for network contention in a fast and scalable manner. SMPI also implements an original and validated piece-wise linear model for data transfer times between cluster nodes. Finally SMPI simulations of large-scale applications on large-scale platforms can be executed on a single node thanks to techniques to reduce the simulation's compute time and memory footprint. These contributions are validated via a large set of experiments in which SMPI is compared to popular MPI implementations with a view to assess its accuracy, scalability, and speed. Pierre-Nicolas Clauss, Mark Stillwell, Stéphane Genaud, Frédéric Suter, Henri Casanova, Martin Quinson |
IPDPS | 5 |
| 2011 | Checkpointing strategies for parallel jobsabstractThis work provides an analysis of checkpointing strategies for minimizing expected job execution times in an environment that is subject to processor failures. In the case of both sequential and parallel jobs, we give the optimal solution for exponentially distributed failure inter-arrival times, which, to the best of our knowledge, is the first rigorous proof that periodic checkpointing is optimal. For non-exponentially distributed failures, we develop a dynamic programming algorithm to maximize the amount of work completed before the next failure, which provides a good heuristic for minimizing the expected execution time. Our work considers various models of job parallelism and of parallel checkpointing overhead. We first perform extensive simulation experiments assuming that failures follow Exponential or Weibull distributions, the latter being more representative of real-world systems. The obtained results not only corroborate our theoretical findings, but also show that our dynamic programming algorithm significantly outperforms previously proposed solutions in the case of Weibull failures. We then discuss results from simulation experiments that use failure logs from production clusters. These results confirm that our dynamic programming algorithm significantly outperforms existing solutions for real-world clusters. Marin Bougeret, Henri Casanova, Mikaël Rabie, Yves Robert, Frédéric Vivien |
SC | 2 |
| 2011 | Resource allocation for multiple concurrent in-network stream-processing applications
Anne Benoit, Henri Casanova, Veronika Rehn-Sonigo, Yves Robert |
Parallel Comput. | 2 |
| 2010 | Non-clairvoyant Scheduling of Multiple Bag-of-Tasks Applications
Henri Casanova, Matthieu Gallet, Frédéric Vivien |
Euro-Par (1) | 1 |
| 2010 | Fast and scalable simulation of volunteer computing systems using SimGridabstractAdvances in internetworking technology and the decreasing cost-performance ratio of commodity computing components have enabled Volunteer Computing (VC). VC platforms aggregate tens or hundreds of thousands of hosts. These hosts are typically volatile, which raises difficult research questions. Most research in this area relies on simulation. The main issue when developing VC simulators is scalability: How to perform simulations of large-scale VC platforms with reasonable amounts of memory and reasonably fast? To achieve scalability, state-of-the-art VC simulators employ simplistic simulation models and/or target on narrow platform and application scenarios. In this paper we enable VC simulations using the general-purpose SimGrid simulation framework, which provides significantly more realistic and flexible simulation capabilities than the aforementioned simulators. Our key contribution is a set of improvements to SimGrid so that it brings these benefits to VC simulations while achieving good scalability. Bruno Donassolo, Henri Casanova, Arnaud Legrand, Pedro Velho |
HPDC | 2 |
| 2010 | Checkpointing vs. Migration for Post-Petascale SupercomputersabstractAn alternative to classical fault-tolerant approaches for large-scale clusters is failure avoidance, by which the occurrence of a fault is predicted and a preventive measure is taken. We develop analytical performance models for two types of preventive measures: preventive checkpointing and preventive migration. We also develop an analytical model of the performance of a standard periodic checkpoint fault-tolerant approach. We instantiate these models for platform scenarios representative of current and future technology trends. We find that preventive migration is the better approach in the short term by orders of magnitude. However, in the longer term, both approaches have comparable merit with a marginal advantage for preventive checkpointing. We also find that standard non-prediction-based fault tolerance achieves poor scaling when compared to prediction-based failure avoidance, thereby demonstrating the importance of failure prediction capabilities. Finally, our results show that achieving good utilization in truly large-scale machines (e.g., 220nodes) for parallel workloads will require more than the failure avoidance techniques evaluated in this work. Franck Cappello, Henri Casanova, Yves Robert |
ICPP | 2 |
| 2010 | Minimizing Stretch and Makespan of Multiple Parallel Task Graphs via Malleable AllocationsabstractMany scientific applications can be structured as Parallel Task Graphs (PTGs), i.e., graphs of data-parallel tasks. Adding data-parallelism to a task-parallel application provides opportunities for higher performance and scalability, but poses scheduling challenges. We study the off-line scheduling of multiple PTGs on a single, homogeneous cluster. The objective is to optimize performance and fairness. We propose a novel algorithm that first computes perfectly fair PTG completion times assuming that each PTG is an ideal malleable job. These completion times are then relaxed so that the schedule is organized as a sequence of periods and is still close to the perfectly fair schedule. Finally, since PTGs are not perfectly malleable, the algorithm increases the execution time of all PTGs uniformly until it can successfully schedule each task in a period. Our evaluation in simulation, using both synthetic and real-world application configurations, shows that our algorithm outperforms previously proposed algorithms when considering two different performance metrics and one fairness metric. Henri Casanova, Frédéric Desprez, Frédéric Suter |
ICPP | 1 |
| 2010 | Dynamic fractional resource scheduling for HPC workloadsabstractWe propose a novel job scheduling approach for homogeneous cluster computing platforms. Its key feature is the use of virtual machine technology for sharing resources in a precise and controlled manner. We justify our approach and propose several job scheduling algorithms. We present results obtained in simulations for synthetic and real-world High Performance Computing (HPC) workloads, in which we compare our proposed algorithms with standard batch scheduling algorithms. We find that our approach widely outperforms batch scheduling. We also identify a few promising algorithms that perform well across most experimental scenarios. Our results demonstrate that virtualization technology coupled with lightweight scheduling strategies affords dramatic improvements in performance for HPC workloads. Mark Stillwell, Frédéric Vivien, Henri Casanova |
IPDPS | 3 |
| 2010 | Characterizing fault tolerance in genetic programming
Daniel Lombraña Gonzalez, Francisco Fernández de Vega, Henri Casanova |
Future Gener. Comput. Syst. | 3 |
| 2010 | On cluster resource allocation for multiple parallel task graphs
Henri Casanova, Frédéric Desprez, Frédéric Suter |
J. Parallel Distributed Comput. | 1 |
| 2010 | Resource allocation algorithms for virtualized service hosting platforms
Mark Stillwell, David Schanzenbach, Frédéric Vivien, Henri Casanova |
J. Parallel Distributed Comput. | 4 |
| 2009 | Resource Allocation Using Virtual ClustersabstractWe propose a novel approach for sharing cluster resources among competing jobs. The key advantage of our approach over current solutions is that it increases cluster utilization while optimizing a user-centric metric that captures both notions of performance and fairness. We motivate and formalize the corresponding resource allocation problem, determine its complexity, and propose several algorithms to solve it in the case of a static workload that consists of sequential jobs. Via extensive simulation experiments we identify an algorithm that runs quickly, that is always on par with or better than its competitors, and that produces resource allocations that are close to optimal. We find that the extension of our approach to parallel jobs leads to similarly good results. Finally, we explain how to extend our work to dynamic workloads. Mark Stillwell, David Schanzenbach, Frédéric Vivien, Henri Casanova |
CCGRID | 4 |
| 2009 | Resource allocation strategies for constructive in-network stream processingabstractWe consider the operator mapping problem for in-network stream processing, i.e., the application of a tree of operators in steady-state to multiple data objects that are continuously updated at various locations in a network. Examples of in-network stream processing include the processing of data in a sensor network, or of continuous queries on distributed relational databases. Our aim is to provide the user a set of processors that should be bought or rented in order to ensure that the application achieves a minimum steady-state throughput, and with the objective of minimizing platform cost. We prove that even the simplest variant of the problem is NP-hard, and we design several polynomial time heuristics, which are evaluated via extensive simulations and compared to theoretical bounds. Anne Benoit, Henri Casanova, Veronika Rehn-Sonigo, Yves Robert |
IPDPS | 2 |
| 2009 | Scheduling Parallel Task Graphs on (Almost) Homogeneous Multicluster PlatformsabstractApplications structured as parallel task graphs exhibit both data and task parallelism and arise in many domains. Scheduling these applications efficiently on parallel platforms has been a long-standing challenge. In the case of a single homogeneous platform, such as a cluster, results have been obtained both in theory, i.e., guaranteed algorithms, and, in practice, i.e., pragmatic heuristics. Due to task parallelism, these applications are well suited for execution on distributed platforms that span multiple clusters possibly in multiple institutions. However, the only available results in this context are nonguaranteed heuristics. In this paper, we develop a scheduling algorithm, MCGAS, which is applicable to multicluster platforms that are almost homogeneous. Such platforms are often found as large subsets of multicluster platforms. Our novel contribution is that MCGAS computes task allocations so that a (tunable) performance guarantee is provided. Since a performance guarantee does not necessarily imply good average performance in practice, we also compare MCGAS with a recently proposed nonguaranteed algorithm. Using simulation over a wide range of experimental scenarios, we find that MCGAS leads to better average application makespans than its competitor. Pierre-François Dutot, Tchimou N'Takpé, Frédéric Suter, Henri Casanova |
IEEE Trans. Parallel Distributed Syst. | 4 |
| 2008 | Scheduling mixed-parallel applications with advance reservationsabstractThis paper investigates the scheduling of mixed-parallel applications, which exhibit both task and data parallelism, in advance reservations settings. Both the problem of minimizing application turn-around time and that of meeting a deadline are studied. For each several scheduling algorithms are proposed, some of which borrow ideas from previously published work in non-reservation settings. Algorithms are compared in simulation over a wide range of application and reservation scenarios. The main finding is that schedules computed using the previously published CPA algorithm can be adapted to advance reservation settings, notably resulting in low resource consumption andthus high efficiency. Kento Aida, Henri Casanova |
HPDC | 2 |
| 2008 | Probabilistic allocation of tasks on desktop gridsabstractWhile desktop grids are attractive platforms for executing parallel applications, their volatile nature has often limited their use to so-called "high-throughput" applications. Checkpointing techniques can enable a broader class of applications. Unfortunately, a volatile host can delay the entire execution for a long period of time. Allocating redundant copies of each task to hosts can alleviate this problem by increasing the likelihood that at least one instance of each application task completes successfully. In this paper we demonstrate that it is possible to use statistical characterizations of host availability to make sound task replication decisions. We find that strategies that exploit such statistical characterizations are effective when compared to alternate approaches. We show that this result holds for real-world host availability data, in spite of only imperfect statistical characterizations. Joshua Wingstrom, Henri Casanova |
IPDPS | 2 |
| 2007 | Topic 3 Scheduling and Load-Balancing
Henri Casanova, Olivier Beaumont, Uwe Schwiegelshohn, Marek S. Tudruj |
Euro-Par | 1 |
| 2007 | Generating grid resource requirement specificationsabstractNo abstract available. Richard Y. Huang, Andrew A. Chien, Henri Casanova |
HPDC | 3 |
| 2007 | A Comparison of Scheduling Approaches for Mixed-Parallel Applications on Heterogeneous PlatformsabstractMixed-parallel applications can take advantage of large-scale computing platforms but scheduling them efficiently on such platforms is challenging. In this paper we compare the two main proposed approaches for solving this scheduling problem on a heterogeneous set of homogeneous clusters. We first modify previously proposed algorithms for both approaches and show that our modifications lead to significant improvements. We then perform a comparison of the modified algorithms in simulation over a wide range of application and platform conditions. We find that although both approaches have advantages, one of them is most likely the most appropriate for the majority of users. Tchimou N'Takpé, Frédéric Suter, Henri Casanova |
ISPDC | 3 |
| 2007 | Automatic resource specification generation for resource selectionabstractWith an increasing number of available resources in large-scale distributed environments, a key challenge is resource selection. Fortunately, several middleware systems provide resource selection services. However, a user is still faced with a difficult question: "What should I ask for?" Since most users end up using naïve and suboptimal resource specifications, we propose an automated way to answer this question. We present an empirical model that given a workflow application (DAG-structured) generates an appropriate resource specification, including number of resources, the range of clock rates among the resources, and network connectivity. The model employs application structure information as well as an optional utility function that trades off cost and performance. With extensive simulation experiments for different types of applications, resource conditions, and scheduling heuristics, we show that our model leads consistently to close to optimal application performance and often reduces resource usage. Richard Y. Huang, Henri Casanova, Andrew A. Chien |
SC | 2 |
| 2007 | Characterizing resource availability in enterprise desktop grids
Derrick Kondo, Gilles Fedak, Franck Cappello, Andrew A. Chien, Henri Casanova |
Future Gener. Comput. Syst. | 5 |
| 2007 | Benefits and Drawbacks of Redundant Batch Requests
Henri Casanova |
J. Grid Comput. | 1 |
| 2007 | Scheduling Task Parallel Applications for Rapid Turnaround on Enterprise Desktop Grids
Derrick Kondo, Andrew A. Chien, Henri Casanova |
J. Grid Comput. | 3 |
| 2006 | Scalable Grid Application Scheduling via Decoupled Resource Selection and SchedulingabstractOver the past years grid infrastructures have been deployed at larger and larger scales, with envisioned deployments incorporating tens of thousands of resources. Therefore, application scheduling algorithms can become unscalable (albeit polynomial) and thus unusable in large-scale environments. One reason for unscalability is that these algorithms perform implicit resource selection. One can achieve better scalability by performing explicit resource selection independently from scheduling in a "decoupled' approach. Furthermore, we hypothesize that one can achieve similar or even better performance as with the non-decoupled approach, which we call the "one step" approach, by selecting resources judiciously. Leveraging the Virtual Grid abstraction, we demonstrate that the decoupled approach is indeed both scalable and effective in large-scale and highly heterogeneous resource environments. Anirban Mandal, Henri Casanova, Andrew A. Chien, Yang-Suk Kee, Ken Kennedy, Charles Koelbel |
CCGRID | 3 |
| 2006 | On Resource Volatility in Enterprise Desktop GridsabstractDesktop grids, which use the idle cycles of many desktop PC's, are currently one of the largest distributed systems in the world. Despite the popularity and success of many desk-top grid projects, the volatility of hosts within desktop grids has been poorly understood. Yet, such host characterization is essential for accurate simulation and modelling of such platforms. In this paper, we present application-level traces of four enterprise desktop grids with a wide range of user bases. We then describe aggregate and per host statistics that reflect the volatility of desktop grid resources. Further, we determine the correlation of volatility between resources, and investigate the correlation of volatility and other host characteristics. Finally, we detail a number of implications of these findings with respect to application performance. Derrick Kondo, Gilles Fedak, Franck Cappello, Andrew A. Chien, Henri Casanova |
e-Science | 5 |
| 2006 | On the Harmfulness of Redundant Batch RequestsabstractMost parallel computing resources are controlled by batch schedulers that place requests for computation in a queue until access to compute nodes are granted. Queue waiting times are notoriously hard to predict, making it difficult for users not only to estimate when their applications may start, but also to pick among multiple batch-scheduled resources the one that produce the shortest turnaround time. As a result, an increasing number of users resort to "redundant requests": several requests are simultaneously submitted to multiple batch schedulers on behalf of a single job; once one of these requests is granted access to compute nodes, the others are canceled. Using simulation as well as experiments with a production batch scheduler we investigate whether redundant requests are harmful in terms of (i) schedule performance and fairness, (ii) system load, and (iii) system predictability. We find that two main issues with redundant requests are load on the middleware and unfairness towards users who do not use redundant requests, which both depend on the number of users who use redundant requests and on the amount of request redundancy these users employ Henri Casanova |
HPDC | 1 |
| 2006 | Robust Resource Allocation for Large-scale Distributed Shared Resource EnvironmentsabstractThis paper presents a new formulation of the resource selection and binding problem and proposes a new algorithm called integrated selection and binding to solve this problem. Our insight is that a resource selection algorithm should consider binding failures. Consequently, the key idea of the integrated selection and binding approach is to decompose a resource collection request into components that can be bound and composed independently and to select multiple sets of resources for each component. The integrated approach is more efficient and effective than the separate approach for competitive access to federated resources Yang-Suk Kee, Ken Yocum, Andrew A. Chien, Henri Casanova |
HPDC | 4 |
| 2006 | The SIMGRID Project Simulation and Deployment of Distributed ApplicationsabstractThis paper presents the SlMGRlD software architecture that comprises of four main components: SURF, MSG, GRAS, and SMPI. The last three components provide APIs for implementing, simulating and/or deploying distributed applications. The first component, SURF, is a fast and accurate simulation engine. The article describes all four components in terms of their goals, their usage, and their functionality, including experimental validation results when applicable. We review each components below Arnaud Legrand, Martin Quinson, Henri Casanova, Kayo Fujiwara |
HPDC | 3 |
| 2006 | Using virtual grids to simplify application schedulingabstractUsers and developers of grid applications have access to increasing numbers of resources. While more resources generally mean higher capabilities for an application, they also raise the issue of application scheduling scalability. First, even polynomial time scheduling heuristics may take a prohibitively long time to compute a schedule. Second, and perhaps more critical, it may not be possible to gather all the resource information needed by a scheduling algorithm in a scalable manner. Our application focus is scientific workflows, which can be represented as directed acyclic graphs (DAGs). Our claim is that, in future resource-rich environments, simple scheduling algorithms may be sufficient to achieve good workflow performances. We introduce a scalable scheduling approach that uses a resource abstraction called a virtual grid (VG). Our simulations of a range of typical DAG structures and resources demonstrate that a simple greedy scheduling heuristic combined with the virtual grid abstraction is as effective and more scalable than more complex heuristic DAG scheduling algorithms on large-scale platforms Richard Y. Huang, Henri Casanova, Andrew A. Chien |
IPDPS | 2 |
| 2006 | Grid allocation and reservation - Improving grid resource allocation via integrated selection and bindingabstractDiscovering and acquiring appropriate, complex resource collections in large-scale distributed computing environments is a fundamental challenge and is critical to application performance. This paper presents a new formulation of the resource selection problem and a new solution to the resource selection and binding problem called integrated selection and binding. Composition operators in our resource description language and efficient data organization enable our approach to allocate complex resource collections efficiently and effectively even in the presence of competition for resources. Our empirical evaluation shows that the integrated approach can produce solutions of significantly higher quality at higher success rate and lower cost than the traditional separate approach. The success rate of the integrated approach can tolerate as much as 15%-60% lower resource availability than the separate approach. Moreover, most requests have at least the 98th percentile rank and can be served in 6 seconds with a population of 1 million hosts. Yang-Suk Kee, Ken Yocum, Andrew A. Chien, Henri Casanova |
SC | 4 |
| 2006 | Guest Editorial: Special Section on Algorithm Design and Scheduling Techniques (Realistic Platform Models) for Heterogeneous ClustersabstractTHE last decade has seen a dramatic increase in the deployment of heterogeneous distributed computing platforms, in particular, those consisting of heterogeneous clusters, and multiple heterogeneous collections of clusters aggregated over wide-area networks into grids. The software infrastructures and mechanisms to deploy such platforms have been well studied and implementations are already used in production, so that heterogeneous platforms represent a significant, and growing, fraction of the computational power delivered by parallel platforms today. In spite of these successes, many research challenges remain, including those pertaining to distributed algorithms and scheduling algorithms, which are critical for ensuring that these platforms are used effectively. In this context, the goal of this special section on “Algorithm Design and Scheduling Techniques (Realistic Platform Models) for Heterogeneous Clusters” is to gather papers that further our understanding of the impact of platform heterogeneity on the design and evaluation of new such algorithms. In the paper entitled “Allocating Non-Real-Time and Soft Real-Time Jobs in Multiclusters,” Ligang He, Stephen A. Jarvis, Daniel P. Spooner, Hong Jiang, Donna N. Dillenberger, and Graham R. Nudd introduce two workload allocation strategies for large-scale heterogeneous platforms. The first strategy achieves an optimized mean response time for jobs having no real-time requirements. The second strategy obtains an optimized mean miss rate for jobs having soft real-time requirements (i.e., a fraction of jobs are permitted to miss the real-time constraints). Both strategies take into account average system behaviors (such as the mean arrival rate of jobs) to calculate the workload proportions for individual clusters, and update on-the-fly the workload allocation when the change in the mean arrival rate reaches a certain threshold. The allocation schemes are combined with two job dispatching strategies (weighted random and weighted round-robin) to generate new job scheduling algorithms for multicluster environments. In their paper “On the Distribution of Sequential Jobs in Random Brokering for Heterogeneous Computational Grids,” Vandy Berten, Joel Goossens, and Emmanuel Jeannot study resource brokering for scheduling sequential jobs onto a grid platform that consists of heterogeneous sets of homogeneous processors, such as a set of clusters. Resources in each cluster are managed by a local scheduler that maintains a job queue. The paper studies a centralized “metascheduler” that uses a randomized strategy to share available resources among competing jobs. This research considers two cases depending on whether the platform is heavily loaded or lightly loaded. For each case, it obtains both analytical and experimental characterizations of the queue lengths at each local scheduler, CPU utilization, and average job slowdowns. Furthermore, the paper presents a discussion of the system’s behavior when it transitions between a heavily loaded state and a lightly loaded one. All presented theoretical results are corroborated by simulations and provide a thorough description of randomized resource brokering. The research in “Multiple Job Scheduling in a Connection-Limited Data Parallel System” presents a new method for scheduling jobs in a distributed system where the critical resource is the bandwidth to access the stored data. The authors, Alessandro Amoroso and Keith Marzullo, describe an approach that supports the master-worker scheme and can be applied to data parallel computation. They consider a typical wide-area data grid that is comprised of a set of sites, where each site has one or more local area networks. The platform model used is based on the Nile data grid. This paper uses a set of synthetic jobs to compare three schedulers: Greedy, Maxfow, and Hybrid. They tested their new approach under various circumstances and measured its performance by means of several metrics. The new Hybrid scheduler is never worse than either of the other two schedulers, and in 20 percent of the simulated runs, it produced runs that were at least 20 percent better. The paper entitled “Capacity-Aware Multicast Algorithms on Heterogeneous Overlay Networks,” coauthored by Zhan Zhang, Shigang Chen, Yibei Ling, and Randy Chow, addresses the problem of multicast for group IEEE TRANSACTIONS ON PARALLEL AND DISTRIBUTED SYSTEMS, VOL. 17, NO. 2, FEBRUARY 2006 97 Henri Casanova, Yves Robert, Howard Jay Siegel |
IEEE Trans. Parallel Distributed Syst. | 1 |
| 2005 | Efficient resource description and high quality selection for virtual gridsabstractSimple resource specification, resource selection, and effective binding are critical capabilities for grid middleware. We describe the virtual grid, an abstraction for providing these capabilities complex resource environments. Elements of the virtual grid include a novel resource description language (vgDL) and a resource selection and binding component (vgFAB), which accepts a vgDL specification and returns a virtual grid, that is, a set of selected and bound resources. The goals of vgFAB are efficiency, scalability, robustness to high resource contention, and the ability to produce results with quantifiable high quality. We present the design of vgDL, showing how it captures application-level resource abstractions using resource aggregates and connectivity amongst them. We present and evaluate a prototype implementation of vgFAB. Our results show that resource selection and binding for virtual grids of 10,000's of resources can scale up to grids with millions of resources, identifying good matches in less than one second. Further, these matches have quantifiable quality, enabling applications to have high confidence in the results. We demonstrate the effectiveness of our combined selection and binding approach in the presence of resource contention, showing that robust selection and binding can be achieved at moderate cost. Yang-Suk Kee, Dionysios Logothetis, Richard Y. Huang, Henri Casanova, Andrew A. Chien |
CCGRID | 4 |
| 2005 | Scheduling Divisible Loads on Star and Tree Networks: Results and Open ProblemsabstractMany applications in scientific and engineering domains are structured as large numbers of independent tasks with low granularity. These applications are thus amenable to straightforward parallelization, typically in master-worker fashion, provided that efficient scheduling strategies are available. Such applications have been called divisible-loads because a scheduler may divide the computation among worker processes arbitrarily, both in terms of number of tasks and of task sizes. Divisible load scheduling has been an active area of research for the last 15 years. A vast literature offers results and scheduling algorithms for various models of the underlying distributed computing platform. Broad surveys are available that report on, accomplishments in the field. By contrast, We propose a unified theoretical perspective that synthesizes previously published results, several novel results, and open questions, in a view to foster hover divisible load scheduling research. Specifically, we discuss both one-round and multiround algorithms, and we restrict our scope to the popular star and tree network topologies, which we study with both linear and affine cost models for communication and computation. Olivier Beaumont, Henri Casanova, Arnaud Legrand, Yves Robert |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2005 | Multiround Algorithms for Scheduling Divisible LoadsabstractDivisible load applications occur in many fields of science and engineering and can be easily parallelized in a master-worker fashion, but pose several scheduling challenges. While a number of approaches have been proposed that allocate load to workers in a single round, using multiple rounds improves overlap of computation with communication. Unfortunately, multiround algorithms are difficult to analyze and have thus received only limited attention. In this paper, we answer three open questions in the multiround divisible load scheduling area: 1) how to account for latencies, 2) how to account for heterogeneous platforms, and 3) how many rounds should be used. To answer 1), we derive the first closed-form optimal schedule for a homogeneous platform with both computation and communication latencies, for a given number of rounds. To answer 2) and 3), we present a novel algorithm, UMR. We evaluate UMR in a variety of realistic scenarios. Krijn van der Raadt, Henri Casanova |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2004 | From Heterogeneous Task Scheduling to Heterogeneous Mixed Parallel Scheduling
Frédéric Suter, Frédéric Desprez, Henri Casanova |
Euro-Par | 3 |
| 2004 | Modeling Large-Scale Platforms for the Analysis and the Simulation of Scheduling StrategiesabstractSummary form only given. An important trend in scientific computing is the establishment of computing platforms that span multiple institutions to support applications at unprecedented scales and levels of performance. A key issue for achieving high performance is the scheduling of application components onto available resources, which has been an active area of research for several decades. However, most of the platform models traditionally used in scheduling research, and in particular the network models, break down for platforms spanning multiple (wide-area) networks. In this paper we examine modeling issues for large-scale platforms. More specifically, we discuss network latency, bandwidth sharing, and network topology. Our discussion is from the perspective of scheduling research and the main challenge we address is to develop models that are sophisticated enough to be realistic, but simple enough that they are amenable to analysis. Finally, while the models we propose can be used to study scheduling problems directly, they also form a good basis for realistic simulation, which is often the method of choice for comparing scheduling strategies. Henri Casanova |
IPDPS | 1 |
| 2004 | Benchmark Probes for Grid AssessmentabstractSummary form only given. Like all computing platforms, grids are in need of a suite of benchmarks by which they can be evaluated, compared and characterized. As a first step towards this goal, we have developed a set of probes that exercise basic grid operations with the goal of measuring the performance and the performance variability of basic grid operations, as well as the failure rates of these operations. We present measurement data obtained by running our probes on a grid testbed that spans 5 clusters in 3 institutions. These measurements quantify compute times, network transfer times, and Globus middleware overhead. Our results help provide insight into the stability, robustness, and performance of our testbed, and lead us to make some recommendations for future grid development. Greg Chun, Holly Dail, Henri Casanova, Allan Snavely |
IPDPS | 3 |
| 2004 | Characterizing and Evaluating Desktop Grids: An Empirical StudyabstractSummary form only given. Desktop resources are attractive for running compute-intensive distributed applications. Several systems that aggregate these resources in desktop grids have been developed. While these systems have been successfully used for many high throughput applications there has been little insight into the detailed temporal structure of CPU availability of desktop grid resources. Yet, this structure is critical to characterize the utility of desktop grid platforms for both task parallel and even data parallel applications. We address the following questions: (i) What are the temporal characteristics of desktop CPU availability in an enterprise setting? (ii) How do these characteristics affect the utility of desktop grids? (iii) Based on these characteristics, can we construct a model of server "equivalents" for the desktop grids, which can be used to predict application performance? We present measurements of an enterprise desktop grid with over 220 hosts running the Entropia commercial desktop grid software. We utilize these measurements to characterize CPU availability and develop a performance model for desktop grid applications for various task granularities, showing that there is an optimal task size. We then use a cluster equivalence metric to quantify the utility of the desktop grid relative to that of a dedicated cluster. Derrick Kondo, Michela Taufer, Charles L. Brooks III, Henri Casanova, Andrew A. Chien |
IPDPS | 4 |
| 2004 | On the Interference of Communication on Computation in JavaabstractSummary form only given. Overlapping communication with computation is a well-known technique to increase application performance. While it is commonly assumed that communication and computation can be overlapped at no cost, in reality, they do contend for resources and thus interfere with each other. Here we present an empirical quantification of the interference rate of communication on computation. We measure this rate on a single processor communicating with both local and remote processors via Java sockets. Among other results we find that the computation rate can suffer by as much as 50%, and that the reduction is approximately proportional to the communication rate. We conclude that interference deserves further study. Barbara Kreaseck, Larry Carter, Henri Casanova, Jeanne Ferrante |
IPDPS | 3 |
| 2004 | Realistic Modeling and Svnthesis of Resources for Computational GridsabstractUnderstanding large Grid platform configurations and generating representative synthetic configurations is critical for Grid computing research. This paper presents an analysis of existing resource configurations and proposes a Grid platform generator that synthesizes realistic configurations of both computing and communication resources. Our key contributions include the development of statistical models for currently deployed resources and using these estimates for modeling the characteristics of future systems. Through the analysis of the configurations of 114 clusters and over 10,000 processors, we identify appropriate distributions for resource configuration parameters in many typical clusters. Using well-established statistical tests, we validate our models against a second resource collection of 191 clusters and over 10,000 processors, and show that our models effectively capture the resource characteristics found in real world resource infrastructures. These models are realized in a resource generator, which can be easily recalibrated by running it on a training sample set. Yang-Suk Kee, Henri Casanova, Andrew A. Chien |
SC | 2 |
| 2004 | Resource Management for Rapid Application Turnaround on Enterprise Desktop GridsabstractDesktop grids are popular platforms for high throughput applications, but due their inherent resource volatility it is difficult to exploit them for applications that require rapid turnaround. Efficient desktop grid execution of short-lived applications is an attractive proposition and we claim that it is achievable via intelligent resource selection. We propose three general techniques for resource selection: resource prioritization, resource exclusion, and task duplication. We use these techniques to instantiate several scheduling heuristics. We evaluate these heuristics through trace-driven simulations of four representative desktop grid configurations. We find that ranking desk-top resources according to their clock rates, without taking into account their availability history, is surprisingly effective in practice. Our main result is that a heuristic that uses the appropriate combination of resource prioritization, resource exclusion, and task replication achieves performance within a factor of 1.7 of optimal. Derrick Kondo, Andrew A. Chien, Henri Casanova |
SC | 3 |
| 2003 | Clustering Hosts in P2P and Global Computing PlatformsabstractBeing able to identify clusters of nearby hosts among Internet clients provides very useful information for a number of internet and p2p applications. Examples of such applications include web applications, request routing in peer-to-peer overlay network, and distributed computing applications. In this paper, we present and formulate the internet host clustering problem. Leveraging previous work on internet host distance measurement, we propose two hierarchical clustering techniques to solve this problem. The first technique is a marker based hierarchical partitioning approach. The second technique is based on the well known K-means clustering algorithm. We evaluated these two approaches in simulation using a representative Internet topology generated with the GT ITM generator for over 1,000 hosts. Our simulation results demonstrate that our algorithmic clustering approaches effectively identify clusters with arbitrary diameters. Our conclusion is that by leveraging previous work on internet host distance estimation, it is possible to cluster Internet hosts to benefit various applications with various requirements. Abhishek Agrawal, Henri Casanova |
CCGRID | 2 |
| 2003 | Scheduling Distributed Applications: the SimGrid Simulation FrameworkabstractSince the advent of distributed computer systems an active field of research has been the investigation of scheduling strategies for parallel applications. The common approach is to employ scheduling heuristics that approximate an optimal schedule. Unfortunately, it is often impossible to obtain analytical results to compare the efficacy of these heuristics. One possibility is to conducts large numbers of back-to-back experiments on real platforms. While this is possible on tightly-coupled platforms, it is infeasible on modern distributed platforms (i.e. Grids) as it is labor-intensive and does not enable repeatable results. The solution is to resort to simulations. Simulations not only enables repeatable results but also make it possible to explore wide ranges of platform and application scenarios. In this paper we present the SimGrid framework which enables the simulation of distributed applications in distributed computing environments for the specific purpose of developing and evaluating scheduling algorithms. This paper focuses on SimGrid v2, which greatly improves on the first version of the software with more realistic network models and topologies. SimGrid v2 also enables the simulation of distributed scheduling agents, which has become critical for current scheduling research in large-scale platforms. After describing and validating these features, we present a case study by which we demonstrate the usefulness of SimGrid for conducting scheduling research. Arnaud Legrand, Loris Marchal, Henri Casanova |
CCGRID | 3 |
| 2003 | Topic Introduction
Yves Robert, Henri Casanova, Arjan J. C. van Gemund, Dieter Kranzlmüller |
Euro-Par | 2 |
| 2003 | Policies for Swapping MPI ProcessesabstractDespite the enormous amount of research and development work in the area of parallel computing, it is a common observation that simultaneous performance and ease-of-use are elusive. We believe that ease-of-use is critical for many end users, and thus seek performance enhancing techniques that can be easily retrofitted to existing parallel applications. In a precious paper we have presented MPI (message passing interface) process swapping, a simple add-on to the MPI programming environment that can improve performance in shared computing environments. MPI process swapping requires as few as three lines of source code change to an existing application. In this paper we explore a question that we had left open in our previous work: based on which policies should processes be swapped for best performance? Our results show that, with adequate swapping policies, MPI process swapping can provide substantial performance benefits with very limited implementation effort. Otto Sievert, Henri Casanova |
HPDC | 2 |
| 2003 | RUMR: Robust Scheduling for Divisible WorkloadsabstractDivisible workload applications arise in many fields of science and engineering. They can be parallelized in master-worker fashion and relevant scheduling strategies have been proposed to reduce application markspan. Our goal is to developed a practical divisible workload scheduling strategy. This requires that previous work be revisited as several usual assumptions about the computing platform do not hold in practice. We have partially addressed this concern in a previous paper via an algorithm that achieves high performance with realistic resource latency models. In this paper we extend our approach to account for performance prediction errors, which are expected for most real-world performance and applications. In essence, we combine ideas from multiround divisible workload scheduling, for performance, and from factoring-based scheduling, for robustness. We present simulation results to quantify the benefits of our approach compared to our original algorithm and to other previously proposed algorithms. Henri Casanova |
HPDC | 2 |
| 2003 | Using TOP-C and AMPIC to port large parallel applications to the Computational Grid
Gene Cooperman, Henri Casanova, Jim Hayes, Thomas Witzel |
Future Gener. Comput. Syst. | 2 |
| 2003 | A decoupled scheduling approach for Grid application development environments
Holly Dail, Francine Berman, Henri Casanova |
J. Parallel Distributed Comput. | 3 |
| 2003 | Adaptive Computing on the Grid Using AppLeSabstractEnsembles of distributed, heterogeneous resources, also known as computational grids, have emerged as critical platforms for high-performance and resource-intensive applications. Such platforms provide the potential for applications to aggregate enormous bandwidth, computational power, memory, secondary storage, and other resources during a single execution. However, achieving this performance potential in dynamic, heterogeneous environments is challenging. Recent experience with distributed applications indicates that adaptivity is fundamental to achieving application performance in dynamic grid environments. The AppLeS (Application Level Scheduling) project provides a methodology, application software, and software environments for adaptively scheduling and deploying applications in heterogeneous, multiuser grid environments. We discuss the AppLeS project and outline our findings. Francine Berman, Richard Wolski, Henri Casanova, Walfredo Cirne, Holly Dail, Marcio Faerman, Silvia M. Figueira, Jim Hayes, Graziano Obertelli, Jennifer M. Schopf, Gary Shao, Shava Smallen, Neil Spring, Alan Su 0001, Dmitrii Zagorodnov |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2002 | Using TOP-C and AMPIC to Port Large Parallel Applications to the Computational GridabstractPorting large applications to distributed computing platforms is a challenging task from a software engineering perspective. The Computational Grid has gained tremendous popularity as it aggregates unprecedented amounts of compute and storage resources by means of increasingly high performance network technology. The primary aim of this paper is to demonstrate how the development time to port very large applications to this environment can be significantly reduced. TOP-C and AMPIC are software packages that have each seen successful application in their respective domains of parallel computing and process creation/communication. We combine them to implement and deploy a master-worker model of parallel computing over the Computational Grid. To demonstrate the benefit of our approach, we ported the 1,000,000 line Geant4 sequential code in three man-weeks by using our TOP-C/AMPIC integration. This paper evaluates the benefits of our approach from a software engineering perspective, and presents experimental results obtained with the new implementation of Geant4 on a Grid testbed. Gene Cooperman, Henri Casanova, Jim Hayes, Thomas Witzel |
CCGRID | 2 |
| 2002 | A decoupled scheduling approach for the GrADS program development environmentabstractProgram development environments are instrumental in providing users with easy and efficient access to parallel computing platforms. While a number of such environments have been widely accepted and used for traditional HPC systems, there are currently no widely used environments for Grid programming. The goal of the Grid Application Development Software (GrADS) project is to develop a coordinated set of tools, libraries and run-time execution facilities for Grid program development. In this paper, we describe a Grid scheduler component that is integrated as part of the GrADS software system. Traditionally, application-level schedulers (e.g. AppLeS) have been tightly integrated with the application itself and were not easily applied to other applications. Our design is generic: we decouple the scheduler core (the search procedure) from the application-specific (e.g. application performance models) and platform-specific (e.g. collection of resource information) components used by the search procedure. We provide experimental validation of our approach for two representative regular, iterative parallel programs in a variety of real-world Grid testbeds. Our scheduler consistently outperforms static and user-driven scheduling methods. Holly Dail, Henri Casanova, Francine Berman |
SC | 2 |
| 2002 | Innovations of the NetSolve Grid Computing SystemabstractAbstract The NetSolve Grid Computing System was first developed in the mid 1990s to provide users with seamless access to remote computational hardware and software resources. Since then, the system has benefitted from many enhancements like security services, data management faculties and distributed storage infrastructures. This article is meant to provide the reader with details regarding the present state of the project, describing the current architecture of the system, its latest innovations and other systems that make use of the NetSolve infrastructure. Copyright © 2002 John Wiley & Sons, Ltd. Dorian C. Arnold, Henri Casanova, Jack J. Dongarra |
Concurr. Comput. Pract. Exp. | 2 |
| 2002 | Middleware for the use of storage in communication
Micah D. Beck, Dorian C. Arnold, Alessandro Bassi, Francine Berman, Henri Casanova, Jack J. Dongarra, Terry Moore, Graziano Obertelli, James S. Plank, D. Martin Swany, Sathish S. Vadhiyar, Richard Wolski |
Parallel Comput. | 5 |
| 2001 | Simgrid: A Toolkit for the Simulation of Application SchedulingabstractAdvances in hardware and software technologies have made it possible to deploy parallel applications over increasingly large sets of distributed resources. Consequently, the study of scheduling algorithms for such applications has been an active area of research. Given the nature of most scheduling problems one must resort to simulation to effectively evaluate and compare their efficacy over a wide range of scenarios. It has thus become necessary to simulate those algorithms for increasingly complex distributed dynamic, heterogeneous environments. We present Simgrid a simulation toolkit for the study of scheduling algorithms for distributed application. We give the main concepts and models behind Simgrid, describe its API and highlight current implementation issues. We also give some experimental results and describe work that builds on Simgrid's functionalities. Henri Casanova |
CCGRID | 1 |
| 2001 | Topic 03: Scheduling and Load Balancing
Ishfaq Ahmad 0001, Henri Casanova, Rupert W. Ford, Yves Robert |
Euro-Par | 2 |
| 2001 | A Study of Deadline Scheduling for Client-Server Systems on the Computational GridabstractThe Computational Grid is a promising platform for the deployment of various high-performance computing applications. A number of projects have addressed the idea of software as a service on the network. These systems usually implement client-server architectures with many servers running on distributed Grid resources and have commonly been referred to as network-enabled servers (NES). An important question is that of scheduling in this multi-client multi-server scenario. Note that in this context most requests are computationally intensive as they are generated by high-performance computing applications. The Bricks simulation framework has been developed and extensively used to evaluate scheduling strategies for NES systems. The authors first present recent developments and extensions to the Bricks simulation models. They discuss a deadline scheduling strategy that is appropriate for the multi-client multi-server case, and augment it with "Load Correction" and "Fallback" mechanisms which could improve the performance of the algorithm. We then give Bricks simulation results. The results show that future NES systems should use deadline scheduling with multiple fallbacks and it is possible to allow users to make a trade-off between failure-rate and cost by adjusting the level of conservatism of deadline scheduling algorithms. Atsuko Takefusa, Satoshi Matsuoka, Henri Casanova, Francine Berman |
HPDC | 3 |
| 2001 | Applying scheduling and tuning to on-line parallel tomographyabstractTomography is a popular technique to reconstruct the three-dimensional structure of an object from a series of two-dimensional projections. Tomography is resource-intensive and deployment of a parallel implementation onto Computational Grid platforms has been studied in previous work. In this work, we address on-line execution of the application where computation is performed as data is collected from an on-line instrument. The goal is to compute incremental 3-D reconstructions that provide quasi-real-time feedback to the user.We model on-line parallel tomography as a tunable application: trade-offs between resolution of the reconstruction and frequency of feedback can be used to accommodate various resource availabilities. We demonstrate that application scheduling/tuning can be framed as multiple constrained optimization problems and evaluate our methodology in simulation. Our results show that prediction of dynamic network performance is key to efficient scheduling and that tunability allows for production runs of on-line parallel tomography in Computational Grid environments. Shava Smallen, Henri Casanova, Francine Berman |
SC | 2 |
| 2000 | The AppLeS Parameter Sweep Template: User-Level Middleware for the GridabstractThe Computational Grid is a promising platform for the efficient execution of parameter sweep applications over large parameter spaces. To achieve performance on the Grid, such applications must be scheduled so that shared data files are strategically placed to maximize reuse, and so that the application execution can adapt to the deliverable performance potential of target heterogeneous, distributed and shared resources. Parameter sweep applications are an important class of applications and would greatly benefit from the development of Grid middleware that embeds a scheduler for performance and targets Grid resources transparently. In this paper we describe a user-level Grid middleware project, the AppLeS Parameter Sweep Template (APST), that uses application-level scheduling techniques [1] and various Grid technologies to allow the efficient deployment of parameter sweep applications over the Grid. We discuss several possible scheduling algorithms and detail our software design. We then describe our current implementation of APST using systems like Globus [2], NetSolve [3] and the Network Weather Service [4], and present experimental results. Henri Casanova, Graziano Obertelli, Francine Berman, Richard Wolski |
SC | 1 |
| 1999 | Adaptive Scheduling for Task Farming with Grid Middleware
Henri Casanova, James S. Plank, Jack J. Dongarra |
Euro-Par | 1 |
| 1999 | Logistical quality of service in NetSolve
Micah D. Beck, Henri Casanova, Jack J. Dongarra, Terry Moore, James S. Plank, Francine Berman, Richard Wolski |
Comput. Commun. | 2 |
| 1999 | Deploying fault tolerance and taks migration with NetSolve
James S. Plank, Henri Casanova, Micah D. Beck, Jack J. Dongarra |
Future Gener. Comput. Syst. | 2 |
| 1999 | Stochastic Performance Prediction for Iterative Algorithms in Distributed Environments
Henri Casanova, Michael G. Thomason, Jack J. Dongarra |
J. Parallel Distributed Comput. | 1 |
| 1998 | Using Agent-Based Software for Scientific Computing in the NetSolve System
Henri Casanova, Jack J. Dongarra |
Parallel Comput. | 1 |
| 1997 | Java Access to Numerical LibrariesabstractIt is a common and somewhat erroneous belief that Java will always be ‘too slow’ for scientific computing. Two projects under way at the University of Tennessee are addressing the question of scientific computing via Java: NetSolve and f2j. The approaches taken by these two projects are radically different. NetSolve allows users to access pre-installed computational resources, such as hardware and software, distributed across the network. Using these resources, the user can easily perform scientific computing tasks without having any computing resource installed on his or her computer. NetSolve features a Graphical User Interface written in Java as well as a Java Application Programming Interface. The f2j (Fortran to Java) project will provide the numerical subroutines translated from their Fortran source into class files suitable for use by Java programmers. This makes it possible for a Java application or applet to use established legacy numerical code that was originally written in Fortran. This article describes the research issues involved in these two projects and their current limitations. We also explain how, although using two different paradigms and addressing somewhat different classes of users and applications, NetSolve and f2j achieve a common goal: to provide efficient, reliable and portable access to standard numerical libraries via Java. © 1997 John Wiley & Sons, Ltd. Henri Casanova, Jack J. Dongarra, David M. Doolin |
Concurr. Pract. Exp. | 1 |
| 1996 | NetSovle: A Network Server for Solving Computational Science ProblemsabstractThis paper presents a new system, called NetSolve, that allows users to access computational resources, such as hardware and software, distributed across the network. The development of NetSolve was motivated by the need for an easy-to-use, efficient mechanism for using computational resources remotely. Ease of use is obtained as a result of different interfaces, some of which require no programming effort from the user. Good performance is ensured by a load-balancing policy that enables NetSolve to use the computational resources available as efficiently as possible. NetSolve offers the ability to look for computational resources on a network, choose the best one availab le, solve a problem (with retry for fault-tolerance), and return the answer to the user. Henri Casanova, Jack J. Dongarra |
SC | 1 |