VLDB 2026 Research / reviewers in the wild / expert
Henri E. Bal
dblp:b/HenriEBal
· DBLP profile ↗
152ranked-venue papers
21as first author
16since 2021 · last 2025
0000-0001-9827-4461ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 102 · 16 first-author · 4 since 2021Software engineering, systems software and programming languages · 15 · 5 first-authorDatabases, data management, data science and information retrieval · 13Applied, interdisciplinary, general and emerging computing · 10 · 2 since 2021Computer networks · 8 · 5 since 2021Artificial intelligence and machine learning · 5 · 3 since 2021Graphics, computer vision, multimedia, augmented reality and games · 4Security and privacy · 1Human-computer interaction and ubiquitous computing · 1
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2025 | Uirapuru: Timely Video Analytics for High-Resolution Steerable Cameras on Edge DevicesabstractReal-time video analytics on high-resolution cameras has become a popular technology for various intelligent services like traffic control and crowd monitoring. While extensive work has been done on improving analytics accuracy with timing guarantees, virtually all of them target static viewpoint cameras. In this paper, we present Uirapuru, a novel framework for real-time, edge-based video analytics on high-resolution steerable cameras. The actuation performed by those cameras brings significant dynamism to the scene, presenting a critical challenge to existing popular approaches such as frame tiling. To address this problem, Uirapuru incorporates a comprehensive understanding of camera actuation into the system design paired with fast adaptive tiling at a per-frame level. We evaluate Uirapuru on a high-resolution video dataset, augmented by pan-tilt-zoom (PTZ) movements typical for steerable cameras and on real-world videos collected from an actual PTZ camera. Our experimental results show that Uirapuru provides up to 1.45× improvement in accuracy while respecting specified latency budgets or reaches up to 4.53× inference speedup with on-par accuracy compared to state-of-the-art static camera approaches. Guilherme Henrique Apostolo, Pablo Bauszat, Vinod Nigade, Henri E. Bal, Lin Wang 0015 |
MobiCom | 4 |
| 2024 | Secure and Private Vertical Federated Learning for Predicting Personalized CVA Outcomes
Corinne G. Allaart, Marc X. Makkes, Lea Dijksman, Paul van der Nat, Douwe Biesma, Henri E. Bal, Aart van Halteren |
AIME (1) | 6 |
| 2024 | A Little Certainty is All We Need: Discovery and Synchronization Acceleration in Battery-Free IoTabstractThe vision of sustainable IoT constructed from battery-free devices has attracted ample interest in the research community. Yet, efficient device discovery and synchronization—a fundamental problem in IoT systems—remains a critical challenge mainly due to the uncertain ambient energy availability across battery-free devices. We argue that bringing in a small level of certainty is necessary for facilitating communication in battery-free IoT. We propose Pulsar where we introduce a small number of battery-powered devices, serving as the communication coordinator for a large number of battery-free devices. We develop two communication schemes, namely one-to-one, and all-to-all, for Pulsar. Our results based on simulations and prototype-based experiments show that Pulsar achieves consistently good performance across different scenarios while requiring no special hardware or environmental conditions. Gaosheng Liu, Vinod Nigade, Henri E. Bal, Lin Wang 0015 |
APNet | 3 |
| 2024 | CAPSlog: Scalable Memory-Centric Partitioning for Pipeline ParallelismabstractPipeline-parallel training has emerged as a popular method to train large Deep Neural Networks (DNNs), as it allows the use of the combined compute power and memory capacity of multiple Graphics Processing Units (GPUs). However, with the sustaining increase in Deep Learning (DL) model sizes, pipeline parallelism provides only a partial solution to the memory bottleneck in large-scale DNN training. Careful partitioning of the DL model over the available GPUs based on memory usage is required to further alleviate the memory bottleneck and train larger DNNs. mCAP is such a memory-oriented partitioning approach for pipeline parallel systems, but it does not scale to models with many layers and very large hardware setups, as it requires extensive profiling and fails to efficiently navigate the partitioning space to find the most memory-friendly partitioning. In this work, we propose CAPSlog, a scalable memory-centric partitioning approach that can recommend model partitionings for larger and more heterogeneous DL models and for larger hardware setups than existing approaches. CAPSlog introduces a new profiling method and a new, much more scalable algorithm for recommending memory-efficient partitionings. CAPSlog re-duces the profiling time by 67 % compared to existing approaches, searches the partitioning space for the optimal solution orders of magnitude faster and can train significantly larger models. Henk Dreuning, Anna Badia Liokouras, Xiaowei Ouyang, Henri E. Bal, Rob van Nieuwpoort |
PDP | 4 |
| 2024 | NetCL: A Unified Programming Framework for In-Network ComputingabstractThe emergence of programmable data planes (PDPs) has paved the way for in-network computing (INC), a paradigm wherein networking devices actively participate in distributed computations. However, PDPs are still a niche technology, mostly available to network operators, and rely on packet-processing DSLs like P4. This necessitates great networking expertise from INC programmers to articulate computational tasks in networking terms and reason about their code. To lift this barrier to INC we propose a unified compute interface for the data plane. We introduce $\mathrm{C} / \mathrm{C}++$ extensions that allow INC to be expressed as kernel functions processing in-flight messages, and APIs for establishing INC-aware communication. We develop a compiler that translates kernels into P4, and thin runtimes that handle the required network plumbing, shielding INC programmers from low-level networking details. We evaluate our system using common INC applications from the literature. George Karlos, Henri E. Bal, Lin Wang 0015 |
SC | 2 |
| 2024 | Inference serving with end-to-end latency SLOs over dynamic edge networksabstractAbstract While high accuracy is of paramount importance for deep learning (DL) inference, serving inference requests on time is equally critical but has not been carefully studied especially when the request has to be served over a dynamic wireless network at the edge. In this paper, we propose Jellyfish—a novel edge DL inference serving system that achieves soft guarantees for end-to-end inference latency service-level objectives (SLO). Jellyfish handles the network variability by utilizing both data and deep neural network (DNN) adaptation to conduct tradeoffs between accuracy and latency. Jellyfish features a new design that enables collective adaptation policies where the decisions for data and DNN adaptations are aligned and coordinated among multiple users with varying network conditions. We propose efficient algorithms to continuously map users and adapt DNNs at runtime, so that we fulfill latency SLOs while maximizing the overall inference accuracy. We further investigate dynamic DNNs, i.e., DNNs that encompass multiple architecture variants, and demonstrate their potential benefit through preliminary experiments. Our experiments based on a prototype implementation and real-world WiFi and LTE network traces show that Jellyfish can meet latency SLOs at around the 99th percentile while maintaining high accuracy. Vinod Nigade, Pablo Bauszat, Henri E. Bal, Lin Wang 0015 |
Real Time Syst. | 3 |
| 2023 | CAPTURE: Memory-Centric Partitioning for Distributed DNN Training with Hybrid ParallelismabstractDeep Learning (DL) model sizes are increasing at a rapid pace, as larger models typically offer better statistical performance. Modern Large Language Models (LLMs) and image processing models contain billions of trainable parameters. Training such massive neural networks incurs significant memory requirements and financial cost. Hybrid-parallel training approaches have emerged that combine pipelining with data and tensor parallelism to facilitate the training of large DL models on distributed hardware setups. However, existing approaches to design a hybrid-parallel partitioning and parallelization plan for DL models focus on achieving high throughput and not on minimizing memory usage and financial cost. We introduce CAPTURE, a partitioning and parallelization approach for hybrid parallelism that minimizes peak memory usage. CAPTURE combines a profiling-based approach with statistical modeling to recommend a partitioning and parallelization plan that minimizes the peak memory usage across all the Graphics Processing Units (GPUs) in the hardware setup. Our results show a reduction in memory usage of up to 43.9% compared to partitioners in state-of-the-art hybrid-parallel training systems. The reduced memory footprint enables the training of larger DL models on the same hardware resources and training with larger batch sizes. CAPTURE can also train a given model on a smaller hardware setup than other approaches, reducing the financial cost of training massive DL models. Henk Dreuning, Kees Verstoep, Henri E. Bal, Rob van Nieuwpoort |
HiPC | 3 |
| 2023 | Channel-Adaptive Early Exiting Using Reinforcement Learning for Multivariate Time Series ClassificationabstractAs machine and deep learning solutions are deployed on edge devices to tackle real-world classification problems, the approach of early classification during inference is becoming increasingly popular. This approach entails performing the classification after having observed only part of the input and it is motivated by goals such as faster results, computation and communication reduction, and energy conservation, especially in resource-constrained edge intelligence environments. Early classification in the time series analysis domain has been extensively studied. However, when it comes to multivariate time series problems, current solutions only consider early exiting across the temporal dimension, treating the multiple input channels as a singular entity. In this work, we propose a flexible early-exit framework, which can consider input channels in a more fine-grained manner and exit at different points for different channel groups. We implement this framework using reinforcement learning methods and we set heuristics to make our objective tractable while maintaining its practicality. We verify the expected behavior of our framework on synthetic data and evaluate its performance on 26 real datasets. Our extensive experiments demonstrate that, depending on the use case, it can add value to the early classification paradigm, achieving better accuracy for equal input savings. Leonardos Pantiskas, Kees Verstoep, Mark Hoogendoorn, Henri E. Bal |
ICMLA | 4 |
| 2022 | Taking ROCKET on an Efficiency Mission: Multivariate Time Series Classification with LightWaveSabstractNowadays, with the rising number of sensor signals in sectors such as healthcare and industry, the problem of multivariate time series classification (MTSC) is getting increasingly relevant and is a prime target for machine and deep learning approaches. Their expanding adoption in real-world environments is causing a shift in focus from the pursuit of everhigher prediction accuracy with complex models towards practical, deployable solutions that balance accuracy and parameters such as prediction speed. An MTSC model that has attracted attention recently is ROCKET, based on random convolutional kernels, both because of its very fast training process and its state-of-the-art accuracy. However, the large number of features it utilizes may be detrimental to inference time. Examining its theoretical background and limitations enables us to address potential drawbacks and present LightWaveS: a framework for accurate MTSC, which is fast both during training and inference. We show that LightWaveS achieves accuracy comparable to recent MTSC models and speedup ranging from 9x to 53x compared to ROCKET during inference on an edge device, on datasets with comparable accuracy. Leonardos Pantiskas, Kees Verstoep, Mark Hoogendoorn, Henri E. Bal |
DCOSS | 4 |
| 2022 | mCAP: Memory-Centric Partitioning for Large-Scale Pipeline-Parallel DNN Training
Henk Dreuning, Henri E. Bal, Rob van Nieuwpoort |
Euro-Par | 2 |
| 2022 | An Empirical Evaluation of Multivariate Time Series Classification with Input Transformation across Different DimensionsabstractIn current research, machine and deep learning solutions for the classification of temporal data are shifting from single-channel datasets (univariate) to problems with multiple channels of information (multivariate). The majority of these works are focused on the method novelty and architecture, and the format of the input data is often treated implicitly. Particularly, multivariate datasets are often treated as a stack of univariate time series in terms of input preprocessing, with scaling methods applied across each channel separately. In this evaluation, we aim to demonstrate that the additional channel dimension is far from trivial and different approaches to scaling can lead to significantly different results in the accuracy of a solution. To that end, we test seven different data transformation methods on four different temporal dimensions and study their effect on the classification accuracy of five recent methods. We show that, for the large majority of tested datasets, the best transformation-dimension configuration leads to an increase in the accuracy compared to the result of each model with the same hyperparameters and no scaling, ranging from 0.16 to 76.79 percentage points. We also show that if we keep the transformation method constant, there is a statistically significant difference in accuracy results when applying it across different dimensions, with accuracy differences ranging from 0.23 to 47.79 percentage points. Finally, we explore the relation of the transformation methods and dimensions to the classifiers, and we conclude that there is no prominent general trend, and the optimal configuration is dataset- and classifier-specific. Leonardos Pantiskas, Kees Verstoep, Mark Hoogendoorn, Henri E. Bal |
ICMLA | 4 |
| 2022 | Vertical Split Learning - an exploration of predictive performance in medical and other use casesabstractIn healthcare and other fields, data of an individual is often vertically partitioned across multiple organizations. Creating a centralized data store for AI algorithm development is cumbersome in such cases because of concerns like privacy and data ownership. Methods of distributed learning over vertically partitioned data could offer a solution here. While several studies have evaluated the feasibility, privacy and efficiency of such methods, an extensive evaluation of their impact on predictive performance compared to a centralized approach is missing. Vertical Split Learning (VSL) aims to provide vertical distributed learning through distributed neural network architectures. Our study adapts and applies VSL to 8 datasets, both in medicine and beyond, evaluating the impact of different network and (vertical) feature distributions on predictive performance. In most configurations VSL yields comparable predictive performance to its centralized counterparts. However, certain data and network distributions give an unexpected and severe loss of performance. Based on our findings we give some initial recommendations under which conditions VSL can be applied as a suitable alternative for data centralization. Corinne G. Allaart, Björn Keyser, Henri E. Bal, Aart van Halteren |
IJCNN | 3 |
| 2022 | Jellyfish: Timely Inference Serving for Dynamic Edge NetworksabstractWhile high accuracy is of paramount importance for deep learning (DL) inference, serving inference requests on time is equally critical but has not been carefully studied especially when the request has to be served over a dynamic wireless network at the edge. In this paper, we propose Jellyfish—a novel edge DL inference serving system that achieves soft guarantees on end-to-end inference latency often specified as a service-level objective (SLO). To handle the network variability, Jellyfish exploits both data and deep neural network (DNN) adaptation to conduct tradeoffs between accuracy and latency. Jellyfish features a new design that enables collective adaptation policies where the decisions for data and DNN adaptations are aligned and coordinated among multiple users with varying network conditions. We propose efficient algorithms to dynamically adapt DNNs and map users, so that we fulfill latency SLOs while maximizing the overall inference accuracy. Our experiments based on a prototype implementation and real-world WiFi and LTE network traces show that Jellyfish can meet latency SLOs at around the 99th percentile while maintaining high accuracy. Vinod Nigade, Pablo Bauszat, Henri E. Bal, Lin Wang 0015 |
RTSS | 3 |
| 2021 | Don't You Worry 'Bout a Packet: Unified Programming for In-Network ComputingabstractIn-network computing is gaining momentum as programmable switches are increasingly employed for compute acceleration. Designed for packet processing, data plane programming languages force developers to express compute in networking terms, resulting in a complex, error-prone practice. We envision the unification of switch and host programming and propose the Net Compute Language (NCL), a C/C++ extension for expressing computational kernels for switches to execute. NCL implements Compute Centric Communication (C3), our proposed programming model for INC under which, point-to-point primitives are augmented to carry out computations. We motivate our approach with real-world use cases and discuss the technical challenges for its realization. George Karlos, Henri E. Bal, Lin Wang 0015 |
HotNets | 2 |
| 2021 | Better Never Than Late: Timely Edge Video Analytics Over the AirabstractEdge video analytics based on deep learning has become an important building block for many modern intelligent applications such as mobile augmented reality and autonomous driving. Various mechanisms have been developed to handle dynamic wireless networks, compute resource availability, and achieve high analytics accuracy via filtering, DNN compression, pruning, and adaptation. So far, limited attention has been paid to timeliness---providing strict service-level objectives (SLO) for edge video analytics pipelines, which is essential for the usability of user-interactive and mission-critical intelligent applications. In this paper, we analyze the challenges in achieving SLO for edge video analytics and present a system design for timely edge video analytics over the air leveraging a simple yet effective idea---feedback control. Our preliminary evaluation based on a system prototype and real-world network traces shows the potential of our design. We also discuss the limitations, calling for future work. Vinod Nigade, Ramon Winder, Henri E. Bal, Lin Wang 0015 |
SenSys | 3 |
| 2021 | Service Placement for Collaborative Edge ApplicationsabstractEdge computing is emerging as a promising computing paradigm for supporting next-generation applications that rely on low-latency network connections in the Internet-of-Things (IoT) era. Many edge applications, such as multi-player augmented reality (AR) gaming and federated machine learning, require that distributed clients work collaboratively for a common goal through message exchanges. Given an edge network, it is an open problem how to deploy such collaborative edge applications to achieve the best overall system performance. This paper presents a formal study of this problem. We first provide a mix of cost models to capture the system. Based on a thorough formulation, we propose an iterative algorithm dubbed ITEM, where in each iteration, we construct a graph to encode all the costs and convert the cost optimization problem into a graph cut problem. By obtaining the minimum s-t cut via existing max-flow algorithms, we address the original problem via solving a series of graph cuts. We rigorously prove that ITEM has a parameterized constant approximation ratio. Inspired by the optimal stopping theory, we further design an online algorithm called OPTS, based on optimally alternating between partial and full placement updates. Our evaluations with real-world data traces demonstrate that ITEM performs close to the optimum (within 5%) and converges fast. OPTS achieves a bounded performance as expected while reducing full updates by more than 67% of the time. Lin Wang 0015, Lei Jiao 0002, Ting He 0001, Jun Li 0001, Henri E. Bal |
IEEE/ACM Trans. Netw. | 5 |
| 2020 | Handling Impossible Derivations During Stream Reasoning
Hamid R. Bazoobandi, Henri E. Bal, Frank van Harmelen, Jacopo Urbani |
ESWC | 2 |
| 2020 | Accelerating Overlapping Community Detection: Performance Tuning a Stochastic Gradient Markov Chain Monte Carlo Algorithm
Ismail El-Helw, Rutger F. H. Hofman, Henri E. Bal |
Euro-Par | 3 |
| 2020 | Clownfish: Edge and Cloud Symbiosis for Video Stream AnalyticsabstractDeep learning (DL) has shown promising results on complex computer vision tasks for video stream analytics recently. However, DL-based analytics typically requires intensive computation, which imposes challenges to the current computing infrastructure. In particular, cloud-only solutions struggle to maintain stable real-time performance due to the streaming over the best-effort Internet, while edge-only solutions require the DL model to be optimized (e.g., pruned or quantized) carefully to fit on resource-constrained devices, affecting the analytics quality. In this paper, we propose Clownfish, a framework for efficient video stream analytics that achieves symbiosis of the edge and the cloud. Clownfish deploys a lightweight optimized DL model at the edge for fast response and a complete DL model at the cloud for high accuracy. By exploiting the temporal correlation in video content, Clownfish sends only a subset of video frames intermittently to the cloud and enhances the analytics quality by fusing the results from the cloud model with these from the edge model. Our evaluation based on a system prototype shows that Clownfish always runs in real time and is able to achieve analytics quality comparable to that of cloud-only solutions, even under highly variable network conditions. Clownfish is generally applicable to all video stream analytics tasks that can leverage temporal correlations. Vinod Nigade, Lin Wang 0015, Henri E. Bal |
SEC | 3 |
| 2020 | Rocket: efficient and scalable all-pairs computations on heterogeneous platformsabstractAll-pairs compute problems apply a user-defined function to each combination of two items of a given data set. Although these problems present an abundance of parallelism, data reuse must be exploited to achieve good performance. Several researchers considered this problem, either resorting to partial replication with static work distribution or dynamic scheduling with full replication. In contrast, we present a solution that relies on hierarchical multi-level software-based caches to maximize data reuse at each level in the distributed memory hierarchy, combined with a divide-and-conquer approach to exploit data locality, hierarchical work-stealing to dynamically balance the workload, and asynchronous processing to maximize resource utilization. We evaluate our solution using three real-world applications, from digital forensics, localization microscopy, and bioinformatics, on different platforms, from desktop machine to a supercomputer. Results shows excellent efficiency and scalability when scaling to 96 GPUs, even obtaining super-linear speedups due to a distributed cache. Stijn Heldens, Pieter Hijma, Ben van Werkhoven, Jason Maassen, Henri E. Bal, Rob van Nieuwpoort |
SC | 5 |
| 2020 | Parallel and Distributed Machine Learning Algorithms for Scalable Big Data Analytics
Henri E. Bal, Arindam Pal 0001 |
Future Gener. Comput. Syst. | 1 |
| 2019 | Aves: A Decision Engine for Energy-efficient Stream Analytics across Low-power DevicesabstractToday's low-power devices, such as smartphones and wearables, form a very heterogeneous ecosystem. Applications in such a system typically follow a reactive pattern based on stream analytics, i.e., sensing, processing, and actuating. Despite the simplicity of this pattern, deciding where to place the processing tasks of an application to achieve energy efficiency is non-trivial in a heterogeneous system since application components are distributed across multiple devices. In this paper, we present Aves - a decision-making engine based on a holistic energy-prediction model, with which the processing tasks of applications can be placed automatically in an energy-efficient manner without programmer/user intervention. We validate the effectiveness of the model and reveal several counter-intuitive placement decisions. Our decision engine's improvements are typically 10-30%, with up to a factor 14 in the most extreme cases. We also show that Aves gives an accurate decision in comparison with real energy measurements for two sensor-based applications. Roshan Bharath Das, Marc X. Makkes, Alexandru Uta, Lin Wang 0015, Henri E. Bal |
IEEE BigData | 5 |
| 2019 | A Programming Framework for Heterogeneous Stream AnalyticsabstractSensor-based applications using Big Data are of increasing importance in various fields. A typical example of such use cases is building health-care applications [1], [2]. A typical scenario is where a patient's heart rate is monitored by a smartwatch. A smartphone can then analyze the gathered data and identify patterns in the patient's heart rate. However, if the data analysis is too complex to be performed on a smartphone, the computation could be offloaded to a nearby cloudlet or a remote cloud. A decision usually follows the analysis, and actuation is performed accordingly (e.g., a message is sent to either the patient or the doctor). Developing such an application is intrinsically complex, as the programmer needs to reconcile different APIs specific to different platforms. Roshan Bharath Das, Marc X. Makkes, Alexandru Uta, Lin Wang 0015, Henri E. Bal |
IEEE BigData | 5 |
| 2018 | RideMatcher: Peer-to-Peer Matching of Passengers for Efficient RidesharingabstractThe daily home-office commute of millions of people in crowded cities puts a strain on air quality, traveling time and noise pollution. This is especially problematic in western cities, where cars and taxis have low occupancy with daily commuters. To reduce these issues, authorities often encourage commuters to share their rides, also known as carpooling or ridesharing. To increase the ridesharing usage it is essential that commuters are efficiently matched. In this paper we present RideMatcher, a novel peer-to-peer system for matching car rides based on their routes and travel times. Unlike other ridesharing systems, RideMatcher is completely decentralized, which makes it possible to deploy it on distributed infrastructures, using fog and edge computing. Despite being decentralized, our system is able to efficiently match ridesharing users in near real-time. Our evaluations performed on a dataset with 34,837 real taxi trips from New York show that RideMatcher is able to reduce the number of taxi trips by up to 65%, the distance traveled by taxi cabs by up to 64%, and the cost of the trips by up to 66%. Nicolae Vladimir Bozdog, Marc X. Makkes, Aart van Halteren, Henri E. Bal |
CCGrid | 4 |
| 2018 | On Optimising Cost and Value in eScience: Case Studies in Radio AstronomabstractLarge-scale science instruments, such as the LHC and recent distributed radio telescopes such as LOFAR, show that we are in an era of data-intensive scientific discovery. All of these instruments rely critically on significant eScience resources, both hardware and software, to do science. Considering limited science budgets, and the small fraction of these that can be dedicated to compute hardware and software, there is a strong and obvious desire for low-cost computing. However, optimizing for cost is only half of the equation, the value potential over the lifetime of the instrument should also be taken into account. Using a tangible example, compute hardware, we introduce a conceptual model to approximate the lifetime relative science merit of such a system. With a number of case studies, focused on eScience applications in radio astronomy past, present and future, we show that the hardware-based analysis can be applied more broadly. While the introduced model is not intended to result in a numeric value for merit, it does enumerate some components that define this metric. P. Chris Broekema, Verity L. Allan, Henri E. Bal |
eScience | 3 |
| 2017 | Kea: A Computation Offloading System for Smartphone Sensor DataabstractNowadays smartphones are equipped with many sensors which applications can continuously invoke to acquire real-time sensor information, such as GPS tracking. Due to the resource-constrained nature of the smartphones, it is often beneficial if the processing of the sensor data is offloaded to a remote resource. However, the decision to offload the computation depends on a multitude of factors such as the hardware capabilities of the phone, the communication energy and latency and the characteristics of the stream computations, e.g., window size, sensor frequency and operational complexity.In this paper we introduce Kea, a profiling-based computation offloading system that automatically decides whether offloading is beneficial for smartphones. The decision making is based on two criteria: the power consumption of the application and the elapsed time for processing the sensor data. Our evaluation results show that unexpected factors such as CPU frequency scaling and the network state also influence the decision-making process. In addition, we show that Kea's profiling overhead is negligible. Roshan Bharath Das, Nicolae Vladimir Bozdog, Marc X. Makkes, Henri E. Bal |
CloudCom | 4 |
| 2017 | P^2-SWAN: Real-Time Privacy Preserving Computation for IoT EcosystemsabstractSensitive personal user-data collected by Internet-of-Things (IoT) devices is vulnerable to information leaks when uploaded to third-party cloud computing infrastructures. Even though data is encrypted before being sent, to perform analyses on the received data, the computing infrastructure typically decrypts the data, and then performs computation. Therefore, during computation, data can be leaked by means of honest-but-curious system administrators. To overcome this vulnerability, homomorphic encryption enables "blind" computation directly on encrypted data, thus rendering obsolete any data leaks. However, homomorphic encryption is highly resource demanding, as it performs many compute-intensive operations during encryption, while also increasing the size of ciphertexts, which makes it unsuitable for low-powered (mobile) IoT devices. For similar reasons, performing operations on encrypted data is also a challenging task, especially when real-time decision-making is needed. In such scenarios, efficient solutions must be augmented by placing computation close to the data: at the network edge. In this paper, we introduce privacy preserving SWAN (P2-SWAN), a homomorphic-encryption enabled mobile computing framework. Even though such encryption adds significant computational overhead, our evaluation shows that it is feasible on low-powered (mobile) devices. The overhead induced on such devices is minimized due to our carefully crafted implementation. We show that performing encrypted operations achieves excellent scalability, thus only modest numbers of computing servers can handle the load for data generated by millions of devices. Furthermore, our proposed approach achieves real-time computation not only for encrypting data on mobile devices, but also for performing encrypted computation. Marc X. Makkes, Alexandru Uta, Roshan Bharath Das, Nicolae Vladimir Bozdog, Henri E. Bal |
ICFEC | 5 |
| 2017 | An Empirical Study on How the Distribution of Ontologies Affects Reasoning on the Web
Hamid R. Bazoobandi, Jacopo Urbani, Frank van Harmelen, Henri E. Bal |
ISWC (1) | 4 |
| 2017 | On the complexities of utilizing large-scale lightpath-connected distributed cyberinfrastructureabstractSummary In Autumn 2013, we—an international team of climate scientists, computer scientists, eScience researchers, and e‐Infrastructure specialists—participated in the enlighten your research global competition, organized to showcase advanced lightpath technologies in support of state‐of‐the‐art research questions. As one of the winning entries, our enlighten your research global team embarked on a very ambitious project to run an extremely high resolution climate model on a collection of supercomputers distributed over two continents and connected using an advanced 10 G lightpath networking infrastructure. Although good progress was made, we were not able to perform all desired experiments due to a varying combination of technical problems, configuration issues, policy limitations and lack of (budget for) human resources to solve these issues. In this paper, we describe our goals, the technical and non‐technical barriers, we encountered and provide recommendations on how these barriers can be removed so future project of this kind may succeed. Copyright © 2016 John Wiley & Sons, Ltd. Jason Maassen, Ben van Werkhoven, Maarten A. J. van Meersbergen, Henri E. Bal, Michael Kliphuis, Sandra E. Brunnabend, Henk A. Dijkstra, Gerben van Malenstein, Migiel de Vos, Sylvia Kuijpers, Sander Boele, Jules Wolfrat, Nick Hill, David Wallom, Christian Grimm, Dieter Kranzlmüller, Dinesh Ganpathi, Shantenu Jha, Yaakoub El Khamra, Frank O. Bryan, Benjamin Kirtman, Frank J. Seinstra |
Concurr. Comput. Pract. Exp. | 4 |
| 2016 | Towards Fast Overlapping Community DetectionabstractAccelerating sequential algorithms in order to achieve high performance is often a nontrivial task. However, there are certain properties that can exacerbate this process and make it particularly daunting. For example, building an efficient parallel solution for a data-intensive algorithm requires a deep analysis of the memory access patterns and data reuse potential. Attempting to scale out the computations on clusters of machines introduces further complications due to network speed limitations. In this context, the optimization landscape can be extremely complex owing to the large number of trade-off decisions. In this paper, we discuss our experience designing two parallel implementations of an existing data-intensive machine learning algorithm that detects overlapping communities in graphs. The first design uses a single GPU to accelerate the computations of small data sets. We employed a code generation strategy in order to test and identify the best performing combination of optimizations. The second design uses a cluster of machines to scale out the computations for larger problem sizes. We used a mixture of MPI, RDMA and pipelining in order to circumvent networking overhead. Both these efforts bring us closer to understanding the complex relationships hidden within networks of entities. Ismail El-Helw, Rutger F. H. Hofman, Henri E. Bal |
CCGrid | 3 |
| 2015 | Finding Pulsars in Real-TimeabstractFinding new pulsars has always been a challenging problem, but this challenge is nowadays exacerbated by the increasing data rates of modern radio telescopes. Because of these increased data rates, traditional approaches to searching, based on storing data for off-line processing, are becoming unfeasible. Therefore, we propose a new pulsar searching pipeline that, by exploiting high-performance computing techniques, is able to process observational data in real-time. To achieve the real-time goal we parallelized all the steps of the pipeline to run on many-core accelerators, and used auto-tuning to adapt and optimize the pipeline for different platforms, telescopes, and searching parameters. In this paper, we test our pipeline on three different platforms: two Graphics Processing Units from AMD and NVIDIA, and an Intel Xeon Phi. Furthermore, we test it on three different scenarios, based on the operational parameters of three state-of-the-art telescopes. Results show that our pipeline can adapt to all tested platforms and scenarios, and achieves real-time performance and linear scalability. Because power consumption is a main concern for radio telescopes, and will be the main bottleneck for the construction of the Square Kilometer Array, we also measure the power consumed by our pipeline. By comparing the results obtained on many-core accelerators with the results obtained using a traditional multi-core CPU, we conclude that the accelerators can provide up to a factor 8 improvement in execution time, and up to a factor 6 reduction in power consumption. Alessio Sclocco, Henri E. Bal, Rob van Nieuwpoort |
e-Science | 2 |
| 2015 | A Compact In-Memory Dictionary for RDF Data
Hamid R. Bazoobandi, Steven de Rooij, Jacopo Urbani, Annette ten Teije, Frank van Harmelen, Henri E. Bal |
ESWC | 6 |
| 2015 | Cashmere: Heterogeneous Many-Core ComputingabstractNew generations of many-core hardware become available frequently and are typically attractive extensions for data-centers because of power-consumption and performance benefits. As a result, supercomputers and clusters are becoming heterogeneous and start to contain a variety of many-core devices. Obtaining performance from a homogeneous cluster-computer is already challenging, but achieving it from a heterogeneous cluster is even more demanding. Related work primarily focuses on homogeneous many-core clusters. In this paper we present Cashmere, a programming system for heterogeneous many-core clusters. Cashmere is a tight integration of two existing systems: Satin is a programming system that provides a divide- and-conquer programming model with automatic load-balancing and latency-hiding, while Many-Core Levels is a programming system that provides a powerful methodology to optimize computational kernels for varying types of many-core hardware. We evaluate our system with several classes of applications and show that Cashmere achieves high performance and good scalability. The efficiency of heterogeneous executions is comparable to the homogeneous runs and is >90% in three out of four applications. Pieter Hijma, Ceriel J. H. Jacobs, Rob van Nieuwpoort, Henri E. Bal |
IPDPS | 4 |
| 2015 | Stepwise-refinement for performance: a methodology for many-core programmingabstractSummary Many‐core hardware is targeted specifically at obtaining high performance, but reaching high performance is often challenging because hardware‐specific details have to be taken into account. Although there are many programming systems that try to alleviate many‐core programming, some providing a high‐level language, others providing a low‐level language for control, none of these systems have a clear and systematic methodology as a foundation. In this article, we proposestepwise‐refinement for performance: a novel, clear, and structured methodology for obtaining high performance on many‐cores. We present a system that supports this methodology, offers multiple levels of abstraction to provide programmers a trade‐off between high‐level and low‐level programming, and provides programmers detailed performance feedback. We evaluate our methodology with several widely varying compute kernels on two different many‐core architectures: a Graphical Processing Unit (GPU) and the Xeon Phi. We show that our methodology gives insight in the performance, and that in almost all cases, we gain a substantial performance improvement using our methodology. Copyright © 2015 John Wiley & Sons, Ltd. Pieter Hijma, Rob van Nieuwpoort, Ceriel J. H. Jacobs, Henri E. Bal |
Concurr. Comput. Pract. Exp. | 4 |
| 2014 | Performance Models for CPU-GPU Data TransfersabstractMany GPU applications perform data transfers to and from GPU memory at regular intervals. For example because the data does not fit into GPU memory or because of internode communication at the end of each time step. Overlapping GPU computation with CPU-GPU communication can reduce the costs of moving data. Several different techniques exist for transferring data to and from GPU memory and for overlapping those transfers with GPU computation. It is currently not known when to apply which method. Implementing and benchmarking each method is often a large programming effort and not feasible. To solve these issues and to provide insight in the performance of GPU applications, we propose an analytical performance model that includes PCIe transfers and overlapping computation and communication. Our evaluation shows that the performance models are capable of correctly classifying the relative performance of the different implementations. Ben van Werkhoven, Jason Maassen, Frank J. Seinstra, Henri E. Bal |
CCGRID | 4 |
| 2014 | Property Specification Made Easy: Harnessing the Power of Model Checking in UML Designs
Daniela Remenska, Tim A. C. Willemse, Jeffrey Templon, Kees Verstoep, Henri E. Bal |
FORTE | 5 |
| 2014 | A detailed GPU cache model based on reuse distance theoryabstractAs modern GPUs rely partly on their on-chip memories to counter the imminent off-chip memory wall, the efficient use of their caches has become important for performance and energy. However, optimising cache locality system-atically requires insight into and prediction of cache behaviour. On sequential processors, stack distance or reuse distance theory is a well-known means to model cache behaviour. However, it is not straightforward to apply this theory to GPUs, mainly because of the parallel execution model and fine-grained multi-threading. This work extends reuse distance to GPUs by modelling: 1) the GPU's hierarchy of threads, warps, threadblocks, and sets of active threads, 2) conditional and non-uniform latencies, 3) cache associativity, 4) miss-status holding-registers, and 5) warp divergence. We implement the model in C++ and extend the Ocelot GPU emulator to extract lists of memory addresses. We compare our model with measured cache miss rates for the Parboil and PolyBench/GPU benchmark suites, showing a mean absolute error of 6% and 8% for two cache configurations. We show that our model is faster and even more accurate compared to the GPGPU-Sim simulator. Cedric Nugteren, Gert-Jan van den Braak, Henk Corporaal, Henri E. Bal |
HPCA | 4 |
| 2014 | Glasswing: accelerating mapreduce on multi-core and many-core clustersabstractThe impact and significance of parallel computing techniques is continuously increasing given the current trend of incorporating more cores in new processor designs. However, many Big Data systems fail to exploit the abundant computational power of multi-core CPUs and GPUs to their full potential. We present Glasswing, a scalable MapReduce framework that employs a configurable mixture of coarse- and fine-grained parallelism to achieve high performance on multi-core CPUs and GPUs. We experimentally evaluated the performance of five MapReduce applications and show that Glasswing outperforms Hadoop on a 64-node multi-core CPU cluster by a factor between 1.8 and 4, and by a factor from 20 to 30 on a 16-node GPU cluster. Ismail El-Helw, Rutger F. H. Hofman, Henri E. Bal |
HPDC | 3 |
| 2014 | AJIRA: A Lightweight Distributed Middleware for MapReduce and Stream ProcessingabstractCurrently, MapReduce is the most popular programming model for large-scale data processing and this motivated the research community to improve its efficiency either with new extensions, algorithmic optimizations, or hardware. In this paper we address two main limitations of MapReduce: one relates to the model's limited expressiveness, which prevents the implementation of complex programs that require multiple steps or iterations. The other relates to the efficiency of its most popular implementations (e.g., Hadoop), which provide good resource utilization only for massive volumes of input, operating sub optimally for smaller or rapidly changing input. To address these limitations, we present AJIRA, a new middleware designed for efficient and generic data processing. At a conceptual level, AJIRA replaces the traditional map/reduce primitives by generic operators that can be dynamically allocated, allowing the execution of more complex batch and stream processing jobs. At a more technical level, AJIRA adopts a distributed, multi-threaded architecture that strives at minimizing overhead for non-critical functionality. These characteristics allow AJIRA to be used as a single programming model for both batch and stream processing. To this end, we evaluated its performance against Hadoop, Spark, Esper, and Storm, which are state of the art systems for both batch and stream processing. Our evaluation shows that AJIRA is competitive in a wide range of scenarios both in terms of processing time and scalability, making it an ideal choice where flexibility, extensibility, and the processing of both large and dynamic data with a single programming model are either desirable or even mandatory requirements. Jacopo Urbani, Alessandro Margara, Ceriel J. H. Jacobs, Spyros Voulgaris, Henri E. Bal |
ICDCS | 5 |
| 2014 | Auto-Tuning Dedispersion for Many-Core AcceleratorsabstractDedispersion is a basic algorithm to reconstruct impulsive astrophysical signals. It is used in high sampling-rate radio astronomy to counteract temporal smearing by intervening interstellar medium. To counteract this smearing, the received signal train must be dedispersed for thousands of trial distances, after which the transformed signals are further analyzed. This process is expensive on both computing and data handling. This challenge is exacerbated in future, and even some current, radio telescopes which routinely produce hundreds of such data streams in parallel. There, the compute requirements for dedispersion are high (petascale), while the data intensity is extreme. Yet, the dedispersion algorithm remains a basic component of every radio telescope, and a fundamental step in searching the sky for radio pulsars and other transient astrophysical objects. In this paper, we study the parallelization of the dedispersion algorithm on many-core accelerators, including GPUs from AMD and NVIDIA, and the Intel Xeon Phi. An important contribution is the computational analysis of the algorithm, from which we conclude that dedispersion is inherently memory-bound in any realistic scenario, in contrast to earlier reports. We also provide empirical proof that, even in unrealistic scenarios, hardware limitations keep the arithmetic intensity low, thus limiting performance. We exploit auto-tuning to adapt the algorithm, not only to different accelerators, but also to different observations, and even telescopes. Our experiments show how the algorithm is tuned automatically for different scenarios and how it exploits and highlights the underlying specificities of the hardware: in some observations, the tuner automatically optimizes device occupancy, while in others it optimizes memory bandwidth. We quantitatively analyze the problem space, and by comparing the results of optimal auto-tuned versions against the best performing fixed codes, we show the impact that auto-tuning has on performance, and conclude that it is statistically relevant. Alessio Sclocco, Henri E. Bal, Jason W. T. Hessels, Joeri van Leeuwen, Rob van Nieuwpoort |
IPDPS | 2 |
| 2014 | Scaling MapReduce Vertically and HorizontallyabstractGlass wing is a MapReduce framework that uses OpenCL to exploit multi-core CPUs and accelerators. However, compute device capabilities may vary significantly and require targeted optimization. Similarly, the availability of resources such as memory, storage and interconnects can severely impact overall job performance. In this paper, we present and analyze how MapReduce applications can improve their horizontal and vertical scalability by using a well controlled mixture of coarse- and fine-grained parallelism. Specifically, we discuss the Glass wing pipeline and its ability to overlap computation, communication, memory transfers and disk access. Additionally, we show how Glass wing can adapt to the distinct capabilities of a variety of compute devices by employing fine-grained parallelism. We experimentally evaluated the performance of five MapReduce applications and show that Glass wing outperforms Hadoop on a 64-node multi-core CPU cluster by factors between 1.2 and 4, and factors from 20 to 30 on a 23-node GPU cluster. Similarly, we show that Glass wing is at least 1.5 times faster than GPMR on the GPU cluster. Ismail El-Helw, Rutger F. H. Hofman, Henri E. Bal |
SC | 3 |
| 2014 | Optimizing convolution operations on GPUs using adaptive tiling
Ben van Werkhoven, Jason Maassen, Henri E. Bal, Frank J. Seinstra |
Future Gener. Comput. Syst. | 3 |
| 2014 | Streaming the Web: Reasoning over dynamic data
Alessandro Margara, Jacopo Urbani, Frank van Harmelen, Henri E. Bal |
J. Web Semant. | 4 |
| 2013 | DynamiTE: Parallel Materialization of Dynamic RDF Data
Jacopo Urbani, Alessandro Margara, Ceriel J. H. Jacobs, Frank van Harmelen, Henri E. Bal |
ISWC (1) | 5 |
| 2013 | Scalable RDF data compression with MapReduceabstractSUMMARY The Semantic Web contains many billions of statements, which are released using the resource description framework (RDF) data model. To better handle these large amounts of data, high performance RDF applications must apply a compression technique. Unfortunately, because of the large input size, even this compression is challenging. In this paper, we propose a set of distributed MapReduce algorithms to efficiently compress and decompress a large amount of RDF data. Our approach uses a dictionary encoding technique that maintains the structure of the data. We highlight the problems of distributed data compression and describe the solutions that we propose. We have implemented a prototype using the Hadoop framework, and evaluate its performance. We show that our approach is able to efficiently compress a large amount of data and scales linearly on both input size and number of nodes. Copyright © 2012 John Wiley & Sons, Ltd. Jacopo Urbani, Jason Maassen, Niels Drost, Frank J. Seinstra, Henri E. Bal |
Concurr. Comput. Pract. Exp. | 5 |
| 2013 | User transparent data and task parallel multimedia computing with Pyxis-DT
Timo van Kessel, Ben van Werkhoven, Niels Drost, Jason Maassen, Henri E. Bal, Frank J. Seinstra |
Future Gener. Comput. Syst. | 5 |
| 2013 | Using model checking to analyze the system behavior of the LHC production grid
Daniela Remenska, Tim A. C. Willemse, Kees Verstoep, Jeffrey Templon, Henri E. Bal |
Future Gener. Comput. Syst. | 5 |
| 2012 | User Transparent Data and Task Parallel Multimedia Computing with Pyxis-DTabstractThe research area of Multimedia Content Analysis (MMCA) considers all aspects of the automated extraction of knowledge from multimedia archives and data streams. To satisfy the increasing computational demands of emerging MMCA problems, there is an urgent need to apply High Performance Computing (HPC) techniques. However, as most MMCA researchers are not also HPC experts, in the field there is a demand~for~programming models and tools that are both efficient and easy~to~use. Today several user transparent library-based parallelization tools exist that aim to satisfy both these requirements. Such tools generally use a data parallel approach in which data structures (e.g. video frames) are scattered among the available nodes in a compute cluster. However, for certain MMCA applications a data parallel approach induces intensive communication, which significantly decreases performance. In these situations, we can benefit from applying alternative approaches. This paper presents Pyxis-DT: a user transparent parallel programming model for MMCA applications that employs both data and task parallelism. Hybrid parallel execution is obtained by run-time construction and execution of a task graph consisting of strictly defined building block operations. Each of these building block operations can be executed in data parallel fashion. Results show that for realistic MMCA applications the concurrent use of data and task parallelism can significantly improve performance compared to using either approach in isolation. Timo van Kessel, Niels Drost, Jason Maassen, Henri E. Bal, Frank J. Seinstra |
CCGRID | 4 |
| 2012 | Using Model Checking to Analyze the System Behavior of the LHC Production GridabstractDIRAC (Distributed Infrastructure with Remote Agent Control) is the grid solution designed to support production activities as well as user data analysis for the Large Hadron Collider "beauty" experiment. It consists of cooperating distributed services and a plethora of light-weight agents delivering the workload to the grid resources. Services accept requests from agents and running jobs, while agents actively fulfill specific goals. Services maintain database back-ends to store dynamic state information of entities such as jobs, queues, or requests for data transfer. Agents continuously check for changes in the service states, and react to these accordingly. The logic of each agent is rather simple, the main source of complexity lies in their cooperation. These agents run concurrently, and communicate using the services' databases as a shared memory for synchronizing the state transitions. Despite the effort invested in making DIRAC reliable, entities occasionally get into inconsistent states. Tracing and fixing such behaviors is difficult, given the inherent parallelism among the distributed components and the size of the implementation. In this paper we present an analysis of DIRAC with mCRL2, process algebra with data. We have reverse engineered two critical and related DIRAC subsystems, and subsequently modeled their behavior with the mCRL2 toolset. This enabled us to easily locate race conditions and live locks which were confirmed to occur in the real system. We further formalized and verified several behavioral properties of the two modeled subsystems. Daniela Remenska, Tim A. C. Willemse, Kees Verstoep, Wan J. Fokkink, Jeffrey Templon, Henri E. Bal |
CCGRID | 6 |
| 2012 | Resource optimization in distributed real-time multimedia applicationsabstractThe research area of multimedia content analysis (MMCA) considers all aspects of the automated extraction of knowledge from multimedia archives and data streams. To adhere to strict time constraints, large-scale multimedia applications typically are being executed on distributed systems consisting of large collections of compute clusters. In a distributed scenario, it is first essential to determine the optimal number of compute nodes used by each cluster, properly balancing the complex tradeoff between computation and communication. This issue is referred as the “resource utilization” (RU) problem. Next, it is important to tune the transmission of newly generated data sent to each cluster, so as to obtain the highest service utilization , while minimizing the need for buffering. This latter issue is referred as the problem of “just-in-time” (JIT) communication. In this paper, we first present a simple and easy-to-implement method for the RU problem, which is based on the classical binary search method. Second, we address the JIT problem by introducing a smart adaptive control method that properly reacts to the continuously changing circumstances in distributed systems. Extensive experimental validation of the two approaches on a real distributed system shows that our optimization approaches are indeed highly effective. Robert D. van der Mei, Dennis Roubos, Frank J. Seinstra, Henri E. Bal |
Multim. Tools Appl. | 5 |
| 2012 | Generating synchronization statements in divide-and-conquer programs
Pieter Hijma, Rob van Nieuwpoort, Ceriel J. H. Jacobs, Henri E. Bal |
Parallel Comput. | 4 |
| 2012 | WebPIE: A Web-scale Parallel Inference Engine using MapReduce
Jacopo Urbani, Spyros Kotoulas, Jason Maassen, Frank van Harmelen, Henri E. Bal |
J. Web Semant. | 5 |
| 2012 | Reply to comment on "WebPIE: A Web-scale parallel inference engine using MapReduce"
Jacopo Urbani, Spyros Kotoulas, Jason Maassen, Frank van Harmelen, Henri E. Bal |
J. Web Semant. | 5 |
| 2012 | Corrigendum to "WebPIE: A Web-scale Parallel Inference Engine using MapReduce" [Web Semant. Sci. Serv. Agents World Wide Web 10 (2012) 59-75]
Jacopo Urbani, Spyros Kotoulas, Jason Maassen, Frank van Harmelen, Henri E. Bal |
J. Web Semant. | 5 |
| 2011 | Profiling Energy Consumption of VMs for Green Cloud ComputingabstractThe Green Clouds project in the Netherlands investigates a system-level approach towards greening High-Performance Computing (HPC) infrastructures and clouds. In this paper we present our initial results in profiling virtual machines with respect to three power metrics, i.e. power, power efficiency and energy, under different high performance computing workloads. We built a linear power model that represents the behavior of a single work node and includes the contribution from individual components, i.e. CPU, memory and HDD, to the total power consumption of a single work node. Our results could be part of a power characterization module integrated into clusters' monitoring systems, future Green Clouds energy-savvy scheduler would use this monitoring system to support system-level optimization. Qingwen Chen, Paola Grosso, Karel van der Veldt, Cees T. A. M. de Laat, Rutger F. H. Hofman, Henri E. Bal |
DASC | 6 |
| 2011 | Towards Collaborative Editing of Structured Data on Mobile DevicesabstractThe age of collaborative editing applications on mobile devices is upon us. However, such applications traditionally rely on centralized servers and thus do not operate in fully decentralized environments. This is a problem on mobile devices where network partitions are the norm due to mobility. Furthermore, such systems typically use either a fixed schema for the data, which makes them inflexible to change and leads to abuse of structured data fields for purposes other than the original intent, or else are based on XML, which limits data to document oriented data stores or else makes querying much more cumbersome for developers and users. In contrast, our Interdroid Versioned Database system provides distributed, fully decentralized, compact, relational, versioned databases for Android powered mobile phones. It offers application designers a unique set of tools for easily building decentralized collaborative applications on Android powered mobile devices using familiar Content Provider and SQL like interfaces. Unfortunately, this system requires that the structure of the database and the user interface (UI) used to edit records in the database to be written at compile time. What users would like is to beable to define and adapt the structure of the data at runtime. What developers would like is a system which makes building structured applications, including editing UI even easier than it is with our prior work. In this paper we present an extension to our Interdroid Versioned Database system which adds the ability to define a Content Provider using an Avro schema, as well as a generic editing interface for instances of that schema. We demonstrate how this system allows us to create powerful data oriented applications at either compile or runtime using an example "ToDo" application, and detail how this work will serve as a basis for our future work on merging shared structured data. Nicholas Palmer, Emilian Miron, Roelof Kemp, Thilo Kielmann, Henri E. Bal |
Mobile Data Management (1) | 5 |
| 2011 | QueryPIE: Backward Reasoning for OWL Horst over Very Large Knowledge Bases
Jacopo Urbani, Frank van Harmelen, Stefan Schlobach, Henri E. Bal |
ISWC (1) | 4 |
| 2011 | JEL: unified resource tracking for parallel and distributed applicationsabstractAbstract When parallel applications are run in large‐scale distributed environments, such as grids, peer‐to‐peer (P2P) systems, and clouds, the set of resources used can change dynamically as machines crash, reservations end, and new resources become available. It is vital for applications to respond to these changes. Therefore, it is necessary to keep track of the available resources—a problem which is known to be notoriously difficult. In this article we argue that resource tracking must be provided as the standard functionality in the lower parts of the software stack. We propose a general solution to resource tracking: the Join–Elect–Leave (JEL) model. JEL provides unified resource tracking for parallel and distributed applications across environments. JEL is a simple yet powerful model based on notifying when resources have Joined or Left the computation. We demonstrate that JEL is suitable for resource tracking in a wide variety of programming models, ranging from the fixed resource sets traditionally used in MPI‐1 to flexible grid‐oriented programming models. We compare several JEL implementations, and show these to perform and scale well in several real‐world scenarios involving grids, clouds and P2P systems applied concurrently, and wide‐area systems with failing resources. Using JEL, we have won the first prize in a number of international distributed computing competitions. Copyright © 2010 John Wiley & Sons, Ltd. Niels Drost, Rob van Nieuwpoort, Jason Maassen, Frank J. Seinstra, Henri E. Bal |
Concurr. Comput. Pract. Exp. | 5 |
| 2011 | Zorilla: a peer-to-peer middleware for real-world distributed systemsabstractAbstract The inherent complex nature of current distributed computing architectures hinders the widespread adoption of these systems for mainstream use. In general, users have access to a highly heterogeneous set of compute resources, which may include clusters, grids, desktop grids, clouds, and other compute platforms. This heterogeneity is especially problematic when running parallel and distributed applications. Software is needed which easily combines as many resources as possible into one coherent computing platform. In this paper, we introduce Zorilla: peer‐to‐peer (P2P) middleware that creates a single distributed environment from any available set of compute resources. Zorilla imposes minimal requirements on the resource used, is platform independent, and does not rely on central components. In addition to providing functionality on bare resources, Zorilla can exploit locally available middleware. Zorilla explicitly supports distributed and parallel applications, and allows resources from multiple sites to cooperate in a single computation. Zorilla makes extensive use of both virtualization and P2P techniques. We will demonstrate how virtualization and P2P combine into a simple design, while enhancing functionality and ease of use. Together, these techniques bring our goal a step closer: transparent, easy use of resources, even on very heterogeneous distributed systems. Copyright © 2011 John Wiley & Sons, Ltd. Niels Drost, Rob van Nieuwpoort, Jason Maassen, Frank J. Seinstra, Henri E. Bal |
Concurr. Comput. Pract. Exp. | 5 |
| 2011 | Application-Tailored I/O with StreamlineabstractStreamline is a stream-based OS communication subsystem that spans from peripheral hardware to userspace processes. It improves performance of I/O-bound applications (such as webservers and streaming media applications) by constructing tailor-made I/O paths through the operating system for each application at runtime. Path optimization removes unnecessary copying, context switching and cache replacement and integrates specialized hardware. Streamline automates optimization and only presents users a clear, concise job control language based on Unix pipelines. For backward compatibility Streamline also presents well known files, pipes and sockets abstractions. Observed throughput improvement over Linux 2.6.24 for networking applications is up to 30-fold, but two-fold is more typical. Willem de Bruijn, Herbert Bos, Henri E. Bal |
ACM Trans. Comput. Syst. | 3 |
| 2010 | OWL Reasoning with WebPIE: Calculating the Closure of 100 Billion Triples
Jacopo Urbani, Spyros Kotoulas, Jason Maassen, Frank van Harmelen, Henri E. Bal |
ESWC (1) | 5 |
| 2010 | Massive Semantic Web data compression with MapReduceabstractThe Semantic Web consists of many billions of statements made of terms that are either URIs or literals. Since these terms usually consist of long sequences of characters, an effective compression technique must be used to reduce the data size and increase the application performance. One of the best known techniques for data compression is dictionary encoding. In this paper we propose a MapReduce algorithm that efficiently compresses and decompresses a large amount of Semantic Web data. We have implemented a prototype using the Hadoop framework and we report an evaluation of the performance. The evaluation shows that our approach is able to efficiently compress a large amount of data and that it scales linearly regarding the input size and number of nodes. Jacopo Urbani, Jason Maassen, Henri E. Bal |
HPDC | 3 |
| 2010 | Satin: A high-level and efficient grid programming modelabstractComputational grids have an enormous potential to provide compute power. However, this power remains largely unexploited today for most applications, except trivially parallel programs. Developing parallel grid applications simply is too difficult. Grids introduce several problems not encountered before, mainly due to the highly heterogeneous and dynamic computing and networking environment. Furthermore, failures occur frequently, and resources may be claimed by higher-priority jobs at any time. In this article, we solve these problems for an important class of applications: divide-and-conquer. We introduce a system called Satin that simplifies the development of parallel grid applications by providing a rich high-level programming model that completely hides communication. All grid issues are transparently handled in the runtime system, not by the programmer. Satin's programming model is based on Java, features spawn-sync primitives and shared objects, and uses asynchronous exceptions and an abort mechanism to support speculative parallelism. To allow an efficient implementation, Satin consistently exploits the idea that grids are hierarchically structured. Dynamic load-balancing is done with a novel cluster-aware scheduling algorithm that hides the long wide-area latencies by overlapping them with useful local work. Satin's shared object model lets the application define the consistency model it needs. If an application needs only loose consistency, it does not have to pay high performance penalties for wide-area communication and synchronization. We demonstrate how grid problems such as resource changes and failures can be handled transparently and efficiently. Finally, we show that adaptivity is important in grids. Satin can increase performance considerably by adding and removing compute resources automatically, based on the application's requirements and the utilization of the machines and networks in the grid. Using an extensive evaluation on real grids with up to 960 cores, we demonstrate that it is possible to provide a simple high-level programming model for divide-and-conquer applications, while achieving excellent performance on grids. At the same time, we show that the divide-and-conquer model scales better on large systems than the master-worker approach, since it has no single central bottleneck. Rob van Nieuwpoort, Gosia Wrzesinska, Ceriel J. H. Jacobs, Henri E. Bal |
ACM Trans. Program. Lang. Syst. | 4 |
| 2009 | Ibis: A Programming System for Real-World Distributed Computing
Henri E. Bal |
Euro-Par | 1 |
| 2009 | Mapping and Synchronizing Streaming Applications on Cell Processors
Maik Nijhuis, Herbert Bos, Henri E. Bal, Cédric Augonnet |
HiPEAC | 3 |
| 2009 | Ibis: Real-world problem solving using real-world gridsabstractIbis is an open source software framework that drastically simplifies the process of programming and deploying large-scale parallel and distributed grid applications. Ibis supports a range of programming models that yield efficient implementations, even on distributed sets of heterogeneous resources. Also, Ibis is specifically designed to run in hostile grid environments that are inherently dynamic and faulty, and that suffer from connectivity problems. Recently, Ibis has been put to the test in two competitions organized by the IEEE Technical Committee on Scalable Computing, as part of the CCGrid 2008 and Cluster/Grid 2008 international conferences. Each of the competitions' categories focused either on the aspect of scalability, efficiency, or fault-tolerance. Our Ibis-based applications have won the first prize in all of these categories. In this paper we give an overview of Ibis, and - to exemplify its power and flexibility - we discuss our contributions to the competitions, and present an overview of our lessons learned. Henri E. Bal, Niels Drost, Roelof Kemp, Jason Maassen, Rob van Nieuwpoort, C. van Reeuwijk, Frank J. Seinstra |
IPDPS | 1 |
| 2009 | Assessing the impact of future reconfigurable optical networks on application performanceabstractThe introduction of optical private networks (lightpaths) has significantly improved the capacity of long distance network links, making it feasible to run large parallel applications in a distributed fashion on multiple sites of a computational grid. Besides offering bandwidths of 10 Gbit/s or more, lightpaths also allow network connections to be dynamically reconfigured. This paper describes our experiences with running data-intensive applications on a grid that offers a (manually) reconfigurable optical wide-area network. We show that the flexibility offered by such a network is useful for applications and that it is often possible to estimate the necessary network configuration in advance. Jason Maassen, Kees Verstoep, Henri E. Bal, Paola Grosso, Cees T. A. M. de Laat |
IPDPS | 3 |
| 2009 | Efficient large-scale model checkingabstractModel checking is a popular technique to systematically and automatically verify system properties. Unfortunately, the well-known state explosion problem often limits the extent to which it can be applied to realistic specifications, due to the huge resulting memory requirements. Distributed-memory model checkers exist, but have thus far only been evaluated on small-scale clusters, with mixed results. We examine one well-known distributed model checker, DiVinE, in detail, and show how a number of additional optimizations in its runtime system enable it to efficiently check very demanding problem instances on a large-scale, multi-core compute cluster. We analyze the impact of the distributed algorithms employed, the problem instance characteristics and network overhead. Finally, we show that the model checker can even obtain good performance in a high-bandwidth computational grid environment. Kees Verstoep, Henri E. Bal, Jiri Barnat, Lubos Brim |
IPDPS | 2 |
| 2009 | eyeDentify: Multimedia Cyber Foraging from a SmartphoneabstractThe recent introduction of smartphones has resulted in an explosion of innovative mobile applications. The computational requirements of many of these applications, however, can not be met by the smartphone itself. The compute power of the smartphone can be enhanced by distributing the application over other compute resources. Existing solutions comprise of a light weight client running on the smartphone and a heavy weight compute server running on, for example, a cloud. This places the user in a dependent position, however, because the user only controls the client application. In this paper, we follow a different model, called cyber foraging, that gives users full control over all parts of the application. We have implemented the model using the Ibis middleware. We evaluate the model using an innovative application in the domain of multimedia computing, and show that cyber foraging increases the application's responsiveness and accuracy whilst decreasing its energy usage. Roelof Kemp, Nicholas Palmer, Thilo Kielmann, Frank J. Seinstra, Niels Drost, Jason Maassen, Henri E. Bal |
ISM | 7 |
| 2009 | Executing multicellular differentiation: quantitative predictive modelling of C.elegans vulval developmentabstractMOTIVATION: Understanding the processes involved in multi-cellular pattern formation is a central problem of developmental biology, hopefully leading to many new insights, e.g. in the treatment of various diseases. Defining suitable computational techniques for development modelling, able to perform in silico simulation experiments, is an open and challenging problem. RESULTS: Previously, we proposed a coarse-grained, quantitative approach based on the basic Petri net formalism, to mimic the behaviour of the biological processes during multicellular differentiation. Here, we apply our modelling approach to the well-studied process of Caenorhabditis elegans vulval development. We show that our model correctly reproduces a large set of in vivo experiments with statistical accuracy. It also generates gene expression time series in accordance with recent biological evidence. Finally, we modelled the role of microRNA mir-61 during vulval development and predict its contribution in stabilizing cell pattern formation. Nicola Bonzanni, Elzbieta Krepska, K. Anton Feenstra, Wan J. Fokkink, Thilo Kielmann, Henri E. Bal, Jaap Heringa |
Bioinform. | 6 |
| 2009 | Executing multicellular differentiation: quantitative predictive modelling of C.elegans vulval developmentabstractBioinformatics 25(16), 2049–2056 We regret that Figure 5 on page 4 of this paper was incorrect and should appear as below. Nicola Bonzanni, Elzbieta Krepska, K. Anton Feenstra, Wan J. Fokkink, Thilo Kielmann, Henri E. Bal, Jaap Heringa |
Bioinform. | 6 |
| 2008 | Experiences with Fine-Grained Distributed Supercomputing on a 10G TestbedabstractThis paper shows how lightpath-based networks can allow challenging, fine-grained parallel supercomputing applications to be run on a grid, using parallel retrograde analysis on DAS-3 as a case study. Detailed performance analysis shows that several problems arise that are not present on tightly-coupled systems like clusters. In particular, flow control, asynchronous communication, and host- level communication overheads become new obstacles. By optimizing these aspects, however, a 10 G grid can obtain high performance for this type of communication-intensive application. The class of large-scale distributed applications suitable for running on a grid is therefore larger than previously thought realistic. Kees Verstoep, Jason Maassen, Henri E. Bal, John W. Romein |
CCGRID | 3 |
| 2008 | Modeling "Just-in-Time" Communication in Distributed Real-Time Multimedia ApplicationsabstractThe research area of multimedia content analysis (MMCA) considers all aspects of the automated extraction of new knowledge from large multimedia data streams and archives. In recent years, there has been a tremendous growth (in data and computational demands) in the MMCA domain, and this growth is likely to continue in the near future. Multimedia applications operating in real-time environments must run under very strict time constraints, e.g., to analyze video frames at the same rate as a camera produces them. To adhere to such constraints, large-scale multimedia applications typically are being executed on Grid systems consisting of large collections of compute clusters. In services-based scenarios, where video content analysis is being performed by a set of remote multimedia servers, results on a particular video frame are obtained quickest if a server is unoccupied (i.e., not working on previously submitted frames). Keeping a server unoccupied, however, is a waste of available compute resources. Therefore, it is important to tune the transmission of newly generated video frames to the occupation of remote servers. However, due to variations in transmission latencies, it is difficult to accurately tune the sending of video frames such that resource utilization is optimized. In this paper we refer to this issue as the problem of "just-in-time" communication. In this paper we address this issue by introducing an adaptive control method that reacts to the continuously changing circumstances in Grid systems so as to obtain the highest service utilization possible, and to minimize service response time for individual video frames. Extensive experimental validation on a real distributed system, in combination with a trace-driven simulation, show that our control method indeed is highly effective. Robert D. van der Mei, Dennis Roubos, Frank J. Seinstra, Ger Koole, Henri E. Bal |
CCGRID | 6 |
| 2008 | HITP: A Transmission Protocol for Scalable High-Performance Distributed Storage
Pierpaolo Giacomin, Alessandro Bassi, Frank J. Seinstra, Thilo Kielmann, Henri E. Bal |
Euro-Par | 5 |
| 2008 | Resource tracking in parallel and distributed applicationsabstractwww.cs.vu.nl/ibis In this paper, we introduce the Join-Elect-Leave (JEL) model, a simple yet powerful model for tracking the resources participating in an application. This model is based on the concept of signaling, i.e., notifying the application when resources have Joined or Left the computation. In addition, the model includes Elections, which can be used to select resources with a special role. JEL supports several consistency models and is suitable for resource coordination of a wide variety of applications, ranging from the traditional fixed resource sets used in MPI, to flexible grid-oriented programming models. Categories and Subject Descriptors: C.2.1 [Computer-Communication Networks]: [Distributed Niels Drost, Rob van Nieuwpoort, Jason Maassen, Henri E. Bal |
HPDC | 4 |
| 2007 | Toward an International "Computer Science Grid"abstractThe computer science discipline, especially in large scale distributed systems like grids and P2P systems and in high performance computing areas, tends to address issues related to increasingly complex systems, gathering thousands to millions of non trivial components. Theoretical analysis, simulation and even emulation are reaching their limits. Like in other scientific disciplines such as physics, chemistry and life science, there is a need to develop, run and maintain generations of scientific instruments for the observation of complex distributed systems running at real scale and under reproducible experimental conditions. Grid'5000 and DAS3 are two large scale systems designed as scientific instruments for researchers in the domains of grid, P2P and networking. More than testbeds, Grid'5000 and DAS3 have been designed as "computer science grids", where researchers share experimental resources spanning over large geographical distances, are able to reserve resources, configure them, run their experiments, realize precise measurements and replay the same experiments with the same experimental conditions. Computer scientists use these two platforms to address issues in the different software layers between the hardware and the users: networking protocols, OS, middleware, parallel and distributed application runtimes, and applications. In this paper, we will present two computer science grids: Grid '5000 and DAS3. We will describe the motivations, design and current status of these two systems. We will also present some of their key results, not only in terms of scientific results in computer science, but also the impact of these two systems as research tools. The success of the Grid'5000 and DAS platforms is the basis of an international initiative, having the objective to deploy a European level "computer science grid". Franck Cappello, Henri E. Bal |
CCGRID | 2 |
| 2007 | Persistent Fault-Tolerance for Divide-and-Conquer Applications on the Grid
Gosia Wrzesinska, Ana-Maria Oprescu, Thilo Kielmann, Henri E. Bal |
Euro-Par | 4 |
| 2007 | ARRG: real-world gossipingabstractGossiping is an effective way of disseminating information in large dynamic systems. Until now, most gossiping algorithms have been designed and evaluated using simulations. However, these algorithms often cannot cope with several real-world problems that tend to be overlooked in simulations, such as node failures, message loss, non-atomicity ofinformation exchange, and firewalls. Niels Drost, Elth Ogston, Rob van Nieuwpoort, Henri E. Bal |
HPDC | 4 |
| 2007 | Smartsockets: solving the connectivity problems in grid computingabstractTightly coupled parallel applications are increasingly run in Grid environments. Unfortunately, on many Grid sites the ability of machines to create or accept network connections is severely limited by ?rewalls, network address translation (NAT)or non-routed networks. Multi homing further complicates connection setup and machine identi?cation. Although ad-hoc solutions exist for some of these problems, it is usually up to the application's user to discover the cause of the connectivity problems and ?nd a solution. In this paper we describe SmartSockets1 a communication library that lifts this burden by automatically discovering the connectivity problems and solving them with as little support from the user as possible. Jason Maassen, Henri E. Bal |
HPDC | 2 |
| 2007 | A Component-based Coordination Language for Efficient Reconfigurable Streaming ApplicationsabstractConsumer electronics applications are becoming increasingly complex because of increased functionality requirements, such as watching multiple compressed video streams on a single screen. We address this complexity by allowing a programmer to specify the application in terms of independent components. Components interact using streaming communication and by sending and receiving events. From this component specification, the executable is generated. We use the Hinch run time system and the SpaceCAKE architecture to validate the effectiveness of our approach. Because the specification language is generic, the application can easily be ported to different platforms. Maik Nijhuis, Herbert Bos, Henri E. Bal |
ICPP | 3 |
| 2007 | Self-adaptive applications on the gridabstractGrids are inherently heterogeneous and dynamic. One important problemin grid computing is resource selection, that is, finding anappropriate resource set for the application. Another problem is adaptation to the changing characteristics of the grid environment. Existing solutions to these two problems require that a performance model for an application is known. However, constructing such models is a complex task. In this paper, we investigate an approach that does not require performance models. We start an application on any set of resources. During the application run, we periodically collect the statistics about the application run and deduce application requirements from these statistics. Then, we adjustthe resource set to better fit the application needs. This approach allows us to avoid performance bottlenecks, such as overloaded WAN links or very slow processors, and therefore can yield significant performance improvements. We evaluate our approach in a number of scenarios typical for the Grid. Gosia Wrzesinska, Jason Maassen, Henri E. Bal |
PPoPP | 3 |
| 2007 | User-friendly and reliable grid computing based on imperfect middlewareabstractWriting grid applications is hard. First, interfaces to existing grid middleware often are too low-level for application programmers who are domain experts rather than computer scientists. Second, grid APIs tend to evolve too quickly for applications to follow. Third, failures and configuration incompatibilities require applications to use different solutions to the same problem, depending on the actual sites in use. Rob van Nieuwpoort, Thilo Kielmann, Henri E. Bal |
SC | 3 |
| 2007 | Special Issue Featuring Selected Papers from HPDC-15
Richard Wolski, Henri E. Bal |
J. Grid Comput. | 2 |
| 2007 | Lessons learned from building and calibrating the ICWall, a stereo tiled displayabstractAbstract Implementation of stereo tiled displays is a rather demanding task. In this article we want to share the lessons we have learned during the design and construction of the ICWall tiled display. This large display, used in a classroom setting, is a high‐resolution stereo tiled display (2 × 8 tiles), built from low‐cost commodity components. The overall image is produced by an array of projectors. When building such a system, a key challenge is to align the projector images. We describe our automated approach for alignment/calibration of the left‐ and right‐eye stereo images. We provide measurements that show accuracy of this procedure. We explain and compare two calibration approaches: a single‐pass and a two‐pass rendering method to align the tiled images. We explain how to provide seamless image on the tiled display and which issues have to be solved. We also discuss the depth perception issues on the ICWall for the large audiences. Another important aspect, is the architecture of the software used for PC‐cluster‐based rendering. We describe Aura, the parallel scene graph API that is used for rendering on our tiled display. Copyright © 2007 John Wiley & Sons, Ltd. Tom van der Schaaf, Desmond Germans, Henri E. Bal, Michal Koutek |
Comput. Animat. Virtual Worlds | 3 |
| 2006 | Simple Locality-Aware Co-allocation in Peer-to-Peer Supercomputing
Niels Drost, Rob van Nieuwpoort, Henri E. Bal |
CCGRID | 3 |
| 2006 | Satin++: Divide-and-Share on the GridabstractDivide-and-conquer is a popular and effective paradigm for writing grid-enabled applications. I t has been shown to perform well i n environments wtth high network latencies and dynamically changing numbers of processors. However, an important disadvantage of the divide-and-conquer paradigm is its limited applicability due to the lack of a shared data abstraction. W e propose a divide-and-share model: the divide-andconquer model extended with shared objects. Shared objects implement a relaxed consistency model called guard consistency. W e have implemented Satin++: a framework for writing divide-and-share applications. With Satin++ we implemented a number of applications including VLSI routing, N-body simulation and a S A T solver. W e evaluate the performance of our model on a cluster supercomputer and on the heterogeneous, wide-area Grid15000 testbed and demonstrate that our applications can achieve high eficiencies on the Grid. Gosia Wrzesinska, Jason Maassen, Kees Verstoep, Henri E. Bal |
e-Science | 4 |
| 2006 | Supporting Reconfigurable Parallel Multimedia Applications
Maik Nijhuis, Herbert Bos, Henri E. Bal |
Euro-Par | 3 |
| 2006 | A Problem Solving Environment for interactive modelling of multiway dataabstractAbstract A prototype Problem Solving Environment (PSE) is presented for problems in interactive modelling of multiway data. Multiway data result from measurements as a function of two or more independent variables. The PSE comprises a parameter estimation loop and a model adjustment loop. The model can be specified hierarchically using mathematically described building blocks which encapsulate the model assumptions. A typical case study of three‐way data illustrates the need for interactive model adjustment. Requirements for interactive problem solving are discussed. Copyright © 2005 John Wiley & Sons, Ltd. Ivo H. M. van Stokkum, Henri E. Bal |
Concurr. Comput. Pract. Exp. | 2 |
| 2005 | NETIBIS: an efficient and dynamic communication system for heterogeneous gridsabstractGrids are more heterogeneous and dynamic than traditional parallel or distributed systems, both in terms of processors and of interconnects. A grid communication system must handle many issues: first, it must run on networks that are not yet determined when the application is launched, including user-space interconnects; second, it must transparently run on different networks at the same time; third, it should yield performance close to that of specialized communication systems. In this paper, we present NETIBIS, a new Java communication system that provides a uniform interface for any underlying inter-cluster or intracluster network. NETIBIS solves the heterogeneity issues posed by grid computing by dynamically constructing network protocol stacks out of drivers, self-contained building blocks for flexible configuration, with limited functionality per driver. We describe the design and implementation of the major NETIBIS drivers for serialization, multicast, reliability, and various underlying networks. We also describe various optimizations for performance, like layer collapsing for the GMdriver. We evaluate the performance of NETIBIS on several platforms, including a European grid. Olivier Aumage, Rutger F. H. Hofman, Henri E. Bal |
CCGRID | 3 |
| 2005 | Developing Java Grid Applications with Ibis
Kees van Reeuwijk, Rob van Nieuwpoort, Henri E. Bal |
Euro-Par | 3 |
| 2005 | Balanced Multicasting: High-throughput Communication for Grid ApplicationsabstractMany grid applications need to transfer large amounts of data between the geographically distributed sites of a grid environment. Network heterogeneity between these sites makes throughput optimization of data transfers to multiple sites (multicast) hard or even impossible. We present a technique called balanced multicasting that uses monitoring information for both bandwidth capacity and achievable bandwidth to compute balanced multicast trees at runtime that use application-level traffic shaping at the sender side to avoid self-induced congestion. Our experimental evaluation shows that our approach outperforms existing multicast strategies by large margins. Mathijs den Burger, Thilo Kielmann, Henri E. Bal |
SC | 3 |
| 2005 | Robust Distributed Systems Achieving Self-Management through InferenceabstractSelf-management has often been proposed as a means to reduce the growing complexity of administration in distributed systems. We argue that this can be achieved through aggressive automation of management tasks. To reach a high level of automation, we propose to take an inference-based approach: codify best practices so that they can be reasoned about and adapted at runtime. Concerns specific to distributed systems are dealt with by the innate support for knowledge sharing. We introduce the methodology along with a reference architecture. The method's validity is tested by applying a preliminary implementation to a handful of practical problems. Willem de Bruijn, Herbert Bos, Henri E. Bal |
WOWMOM | 3 |
| 2005 | Ibis: a flexible and efficient Java-based Grid programming environmentabstractAbstract In computational Grids, performance‐hungry applications need to simultaneously tap the computational power of multiple, dynamically available sites. The crux of designing Grid programming environments stems exactly from the dynamic availability of compute cycles: Grid programming environments (a) need to beportableto run on as many sites as possible, (b) they need to beflexibleto cope with different network protocols and dynamically changing groups of compute nodes, while (c) they need to provideefficient(local) communication that enables high‐performance computing in the first place. Existing programming environments are either portable (Java), or flexible (Jini, Java Remote Method Invocation or (RMI)), or they are highly efficient (Message Passing Interface). No system combines all three properties that are necessary for Grid computing. In this paper, we present Ibis, a new programming environment that combines Java's ‘run everywhere’ portability both with flexible treatment of dynamically available networks and processor pools, and with highly efficient, object‐based communication. Ibis can transfer Java objects very efficiently by combining streaming object serialization with a zero‐copy protocol. Using RMI as a simple test case, we show that Ibis outperforms existing RMI implementations, achieving up to nine times higher throughputs with trees of objects. Copyright © 2005 John Wiley & Sons, Ltd. Rob van Nieuwpoort, Jason Maassen, Gosia Wrzesinska, Rutger F. H. Hofman, Ceriel J. H. Jacobs, Thilo Kielmann, Henri E. Bal |
Concurr. Pract. Exp. | 7 |
| 2005 | Object combining: a new aggressive optimization for object intensive programsabstractAbstract Object combining tries to put objects together that have roughly the same life times in order to reduce strain on the memory manager and to reduce the number of pointer indirections during a program's execution. Object combining works by appending the fields of one object to another, allowing allocation and freeing of multiple objects with a single heap (de)allocation. Unlike object inlining, which will only optimize objects where one has a (unique) pointer to another, our optimization also works if there is no such relation. Object inlining also directly replaces the pointer by the inlined object's fields. Object combining leaves the pointer in place to allow more combining. Elimination of the pointer accesses is implemented in a separate compiler optimization pass. Unlike previous object inlining systems, reference field overwrites are allowed and handled, resulting in much more aggressive optimization. Our object combining heuristics also allow unrelated objects to be combined, for example, those allocated inside a loop; recursive data structures (linked lists, trees) can be allocated several at a time and objects that are always used together can be combined. As Java explicitly permits code to be loaded at runtime and allows the new code to contribute to a running computation, we do not require a closed‐world assumption to enable these optimizations (but it will increase performance). The main focus of object combining in this paper is on reducing object (de)allocation overhead, by reducing both garbage collection work and the number of object allocations. Reduction of memory management overhead causes execution time to be reduced by up to 35%. Indirection removal further reduces execution time by up to 6%. Copyright © 2005 John Wiley & Sons, Ltd. Ronald Veldema, Ceriel J. H. Jacobs, Rutger F. H. Hofman, Henri E. Bal |
Concurr. Pract. Exp. | 4 |
| 2004 | An simple and efficient fault tolerance mechanism for divide-and-conquer systemsabstractSummary form only given. We study if fault tolerance can be made simpler and more efficient by exploiting the structure of the application. More specifically, we study divide-and-conquer parallelism, which is a popular and effective paradigm for writing parallel Grid applications. We have designed a novel fault tolerance mechanism for divide-and-conquer applications that reduces the amount of redundant computation by storing results of the discarded in a global (replicated) table. These results can later be reused, thereby minimizing the amount of work lost as a result of a crash. The execution time overhead of our mechanism is close to zero. Our mechanism can handle crashes of multiple processors or entire clusters at the same time.. It can also handle crashes of the root node that initially started the parallel computation. We have incorporated our fault tolerance mechanism in Satin, which is a Java-based divide-and-conquer system. Satin is implemented on top of the Ibis communication library. The core of Ibis is implemented in pure Java, without using any native libraries. The Satin runtime system and our fault tolerance extension also are written entirely in Java. The resulting system therefore is highly portable allowing the software to run unmodified on a heterogeneous Grid. We evaluated the performance of our fault tolerance scheme on a cluster of the Distributed ASCI Supercomputer 2 (DAS-2). In the first part of our tests, we show that the execution time overhead of our mechanism is close to zero. The results of the second part of our tests show that our algorithm salvages most of the work done by alive processors. Finally, we carried out tests on the European GridLab testbed. We ran one of our applications on a set of six heterogeneous parallel machines (four different operating systems, four different architectures) located in four different European countries. After manually killing one of the sites, the program recovered and finished normally. Gosia Wrzesinska, Rob van Nieuwpoort, Jason Maassen, Henri E. Bal |
CCGRID | 4 |
| 2004 | Topic 9: Distributed Systems and Algorithms
Henri E. Bal, Andrzej M. Goscinski, Eric Jul, Giuseppe Prencipe |
Euro-Par | 1 |
| 2004 | Wide-Area Communication for Grids: An Integrated Solution to Connectivity, Performance and Security Problems
Alexandre Denis 0001, Olivier Aumage, Rutger F. H. Hofman, Kees Verstoep, Thilo Kielmann, Henri E. Bal |
HPDC | 6 |
| 2004 | A High Performance Java Middleware with a Real ApplicationabstractPrevious experiments with high-performance Java were initially disappointing. After several years of optimization, this paper investigates the current suitability of such object-oriented middle-ware for High-Performance and Grid programming. Using a middleware o®ering high level abstractions (ProActive), we have replaced the standard Java RMI layer with the optimized Ibis RMI interface. Ibis is a grid programming environment featuring e±cient communications. Using a 3D electromagnetic application (an object-oriented time domain ¯nite volume solver for 3D Maxwell equations) we have ¯rst conducted benchmarks on single clusters, including comparisons with the same application in Fortran MPI. Finally, Grid experiments have been conducted simultaneously on up to 5 di®erent clusters. Overall, the paper reports extremely promising results. For instance, a speed up of 12 on 16 machines (vs. 13.8 for Fortran), a speedup of 100 on 150 machines on a Grid. 1 Fabrice Huet, Denis Caromel, Henri E. Bal |
SC | 3 |
| 2004 | Cluster communication protocols for parallel-programming systemsabstractClusters of workstations are a popular platform for high-performance computing. For many parallel applications, efficient use of a fast interconnection network is essential for good performance. Several modern System Area Networks include programmable network interfaces that can be tailored to perform protocol tasks that otherwise would need to be done by the host processors. Finding the right trade-off between protocol processing at the host and the network interface is difficult in general. In this work, we systematically evaluate the performance of different implementations of a single, user-level communication interface. The implementations make different architectural assumptions about the reliability of the network and the capabilities of the network interface. The implementations differ accordingly in their division of protocol tasks between host software, network-interface firmware, and network hardware. Also, we investigate the effects of alternative data-transfer methods and multicast implementations, and we evaluate the influence of packet size. Using microbenchmarks, parallel-programming systems, and parallel applications, we assess the performance of the different implementations at multiple levels. We use two hardware platforms with different performance characteristics to validate our conclusions. We show how moving protocol tasks to a relatively slow network interface can yield both performance advantages and disadvantages, depending on specific characteristics of the application and the underlying parallel-programming system. Kees Verstoep, Raoul Bhoedjang, Tim Rühl, Henri E. Bal, Rutger F. H. Hofman |
ACM Trans. Comput. Syst. | 4 |
| 2003 | A Java-Based Grid Programming Environment
Henri E. Bal |
Euro-Par | 1 |
| 2003 | Topic Introduction
Henri E. Bal, Domenico Laforenza, Thierry Priol, Péter Kacsuk |
Euro-Par | 1 |
| 2003 | A Million-Fold Speed Improvement in Genomic Repeats DetectionabstractThis paper presents a novel, parallel algorithm for generating top alignments. Top alignments are used for finding internal repeats in biological sequences like proteins and genes. Our algorithm replaces an older, sequential algorithm (Repro), which was prohibitively slow for sequence lengths higher than 2000. The new algorithm is an order of magnitude faster (O n 3 rather than O n 4). The paper presents a three-level parallel implementation of the algorithm: using SIMD multimedia extensions found on present-day processors (a novel technique that can be used to parallelize any application that performs many sequence alignments), using shared-memory parallelism, and using distributed-memory parallelism. It allows processing the longest known proteins (nearly 35000 amino acids). We show exceptionally high speed improvements: between 500 and 831 on a cluster of 64 dualprocessor machines, compared to the new sequential algorithm. Especially for long sequences, extreme speed improvements over the old algorithm are obtained. 1 John W. Romein, Jaap Heringa, Henri E. Bal |
SC | 3 |
| 2003 | CCJ: object-based message passing and collective communication in JavaabstractAbstract CCJ is a communication library that adds MPI‐like message passing and collective operations to Java. Rather than trying to adhere to the precise MPI syntax, CCJ aims at a clean integration of communication into Java's object‐oriented framework. For example, CCJ uses thread groups to support Java's multithreading model and it allows any data structure (not just arrays) to be communicated. CCJ is implemented entirely in Java, on top of RMI, so it can be used with any Java virtual machine. The paper discusses three parallel Java applications that use collective communication. It compares the performance (on top of a Myrinet cluster) of CCJ, RMI and mpiJava versions of these applications and also compares their code complexity. A detailed performance comparison between CCJ and mpiJava is given using the Java Grande Forum MPJ benchmark suite. The results show that neither CCJ's object‐oriented design nor its implementation on top of RMI impose a performance penalty on applications compared to their mpiJava counterparts. The source of CCJ is available from our Web site http://www.cs.vu.nl/manta. Copyright © 2003 John Wiley & Sons, Ltd. Arnold Nelisse, Jason Maassen, Thilo Kielmann, Henri E. Bal |
Concurr. Comput. Pract. Exp. | 4 |
| 2003 | Run-time optimizations for a Java DSM implementationabstractAbstract Jackal is a fine‐grained distributed shared memory implementation of the Java programming language. Jackal implements Java's memory model and allows multithreaded Java programs to run unmodified on distributed‐memory systems. This paper focuses on Jackal's run‐time system, which implements a multiple‐writer, home‐based consistency protocol. Protocol actions are triggered by software access checks that Jackal's compiler inserts before object and array references. To reduce access‐check overhead, the compiler exploits source‐level information and performs extensive static analysis to optimize, lift, and remove access checks. We describe optimizations for Jackal's run‐time system, which mainly consists of discovering opportunities to dispense with flushing of cached data. We give performance results for different run‐time optimizations, and compare their impact with the impact of one compiler optimization. We find that our run‐time optimizations are necessary for good Jackal performance, but only in conjunction with the Jackal compiler optimizations described by Veldema et al. As a yardstick, we compare the performance of Java applications run on Jackal with the performance of equivalent applications that use a fast implementation of Java's Remote Method Invocation instead of shared memory. Copyright © 2003 John Wiley & Sons, Ltd. Ronald Veldema, Rutger F. H. Hofman, Raoul Bhoedjang, Henri E. Bal |
Concurr. Comput. Pract. Exp. | 4 |
| 2003 | Editorial
Henri E. Bal, Klaus-Peter Löhr, Alexander Reinefeld, Craig A. Lee |
Future Gener. Comput. Syst. | 1 |
| 2003 | Griz: experience with remote visualization over an optical grid
Luc Renambot, Tom van der Schaaf, Henri E. Bal, Desmond Germans, Hans J. W. Spoelder |
Future Gener. Comput. Syst. | 3 |
| 2002 | The Polder Computing Environment: a system for interactive distributed simulationabstractAbstract The paper provides an overview of an experimental, Grid‐like computing environment, Polder, and its components. Polder offers high‐performance computing and interactive simulation facilities to computational science. It was successfully implemented on a wide‐area cluster system, the Distributed ASCI Supercomputer. An important issue is an efficient management of resources, in particular multi‐level scheduling and migration of tasks that use PVM or sockets. The system can be applied to interactive simulation, where a cluster is used for high‐performance computations, while a dedicated immersive interactive environment (CAVE) offers visualization and user interaction. Design considerations for the construction of dynamic exploration environments using such a system are discussed, in particular the use of intelligent agents for coordination. A case study of simulatedabdominal vascular reconstruction is subsequently presented: the results of computed tomography or magnetic resonance imaging of a patient are displayed in CAVE, and a surgeon can evaluate the possible treatments by performing the surgeries virtually and analysing the resulting blood flow which is simulated using the lattice‐Boltzmann method. Copyright © 2002 John Wiley & Sons, Ltd. Kamil Iskra, Robert G. Belleman, G. Dick van Albada, J. Santoso, Peter M. A. Sloot, Henri E. Bal, Hans J. W. Spoelder, Marian Bubak |
Concurr. Comput. Pract. Exp. | 6 |
| 2002 | Programming environments for high-performance Grid computing: the Albatross project
Thilo Kielmann, Henri E. Bal, Jason Maassen, Rob van Nieuwpoort, Lionel Eyraud-Dubois, Rutger F. H. Hofman, Kees Verstoep |
Future Gener. Comput. Syst. | 2 |
| 2002 | Analysis of Transposition-Table-Driven Work Scheduling in Distributed SearchabstractThis paper discusses a new work-scheduling algorithm for parallel search of single-agent state spaces, called transposition-table-driven work scheduling, that places the transposition table at the heart of the parallel work scheduling. The scheme results in less synchronization overhead, less processor idle time, and less redundant search effort. Measurements on a 128-processor parallel machine show that the scheme achieves close-to-linear speedups; for large problems the speedups are even superlinear due to better memory usage. On the same machine, the algorithm is 1.6 to 12.9 times faster than traditional work-stealing-based schemes. John W. Romein, Henri E. Bal, Jonathan Schaeffer 0001, Aske Plaat |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2001 | Wide-Area Transposition-Driven SchedulingabstractThe distributed searching of state spaces containing cycles is a challenging task and has been studied for several years. Traditional parallel search algorithms either ignore the cyclic nature of the state space and waste much time in duplicated search effort, or they rely on heavy communication to reduce duplicate work, resulting in a large communication overhead. Both methods perform poorly, even when using a fast, local interconnection. A recently-developed task distribution scheme, called transposition-driven scheduling (TDS), performs much better, since it communicates asynchronously and efficiently suppresses duplicate search effort. TDS, however, requires bandwidths of megabytes per second per processor. In this paper, we investigate how cyclic state spaces can be searched efficiently on a meta-computing system containing multiple clusters, connected by high-latency, low-bandwidth wide-area links. This is quite a challenge, because the wide-area links provide neither the bandwidth required for TDS nor the latency required for traditional distributed search algorithms. We propose a scheme that strongly reduces communication between clusters at the expense of some duplicate search effort. Performance measurements for several applications show that the new scheme outperforms traditional schemes by a wide margin. John W. Romein, Henri E. Bal |
HPDC | 2 |
| 2001 | Efficient load balancing for wide-area divide-and-conquer applicationsabstractDivide-and-conquer programs are easily parallelized by letting the programmer annotate potential parallelism in the form of spawn and sync constructs. To achieve efficient program execution, the generated work load has to be balanced evenly among the available CPUs. For single cluster systems, Random Stealing (RS) is known to achieve optimal load balancing. However, RS is inefficient when applied to hierarchical wide-area systems where multiple clusters are connected via wide-area networks (WANs) with high latency and low bandwidth. Rob van Nieuwpoort, Thilo Kielmann, Henri E. Bal |
PPoPP | 3 |
| 2001 | Source-level global optimizations for fine-grain distributed shared memory systemsabstractThis paper describes and evaluates the use of aggressive static analysis in Jackal, a fine-grain Distributed Shared Memory (DSM) system for Java. Jackal uses an optimizing, source-level compiler rather than the binary rewriting techniques employed by most other fine-grain DSM systems. Source-level analysis makes existing access-check optimizations (e.g., access-check batching) more effective and enables two novel fine-grain DSM optimizations: object-graph aggregation and automatic computation migration. Ronald Veldema, Rutger F. H. Hofman, Raoul Bhoedjang, Ceriel J. H. Jacobs, Henri E. Bal |
PPoPP | 5 |
| 2001 | Parallel application experience with replicated method invocationabstractAbstract We describe and evaluate a new approach to object replication in Java, aimed at improving the performance of parallel programs. Our programming model allows the programmer to define groups of objects that can be replicated and updated as a whole, using reliable, totally‐ordered broadcast to send update methods to all machines containing a copy. The model has been implemented in the Manta highperformance Java system. We evaluate system performance both with microbenchmarks and with a set of five parallel applications. For the applications, we also evaluate ease of programming, compared to RMI implementations. We present performance results for a Myrinet‐based workstation cluster as well as for a wide‐area distributed system consisting of four such clusters. The microbenchmarks show that updating a replicated object on 64 machines only takes about three times the RMI latency in Manta. Applications using Manta's object replication mechanism perform at least as fast as manually optimized versions based on RMI, while keeping the application code as simple as with naive versions that use shared objects without taking locality into account. Using a replication mechanism in Manta's runtime system enables several unmodified applications to run efficiently even on the wide‐area system. Copyright © 2001 John Wiley & Sons, Ltd. Jason Maassen, Thilo Kielmann, Henri E. Bal |
Concurr. Comput. Pract. Exp. | 3 |
| 2001 | Sensitivity of parallel applications to large differences in bandwidth and latency in two-layer interconnects
Aske Plaat, Henri E. Bal, Rutger F. H. Hofman, Thilo Kielmann |
Future Gener. Comput. Syst. | 2 |
| 2001 | Network performance-aware collective communication for clustered wide-area systems
Thilo Kielmann, Henri E. Bal, Sergei Gorlatch, Kees Verstoep, Rutger F. H. Hofman |
Parallel Comput. | 2 |
| 2001 | Efficient Java RMI for parallel programmingabstractJava offers interesting opportunities for parallel computing. In particular, Java Remote Method Invocation (RMI) provides a flexible kind of remote procedure call (RPC) that supports polymorphism. Sun's RMI implementation achieves this kind of flexibility at the cost of a major runtime overhead. The goal of this article is to show that RMI can be implemented efficiently, while still supporting polymorphism and allowing interoperability with Java Virtual Machines (JVMs). We study a new approach for implementing RMI, using a compiler-based Java system called Manta. Manta uses a native (static) compiler instead of a just-in-time compiler. To implement RMI efficiently, Manta exploits compile-time type information for generating specialized serializers. Also, it uses an efficient RMI protocol and fast low-level communication protocols.A difficult problem with this approach is how to support polymorphism and interoperability. One of the consequences of polymorphism is that an RMI implementation must be able to download remote classes into an application during runtime. Manta solves this problem by using a dynamic bytecode compiler, which is capable of compiling and linking bytecode into a running application. To allow interoperability with JVMs, Manta also implements the Sun RMI protocol (i.e., the standard RMI protocol), in addition to its own protocol.We evaluate the performance of Manta using benchmarks and applications that run on a 32-node Myrinet cluster. The time for a null-RMI (without parameters or a return value) of Manta is 35 times lower than for the Sun JDK 1.2, and only slightly higher than for a C-based RPC protocol. This high performance is accomplished by pushing almost all of the runtime overhead of RMI to compile time. We study the performance differences between the Manta and the Sun RMI protocols in detail. The poor performance of the Sun RMI protocol is in part due to an inefficient implementation of the protocol. To allow a fair comparison, we compiled the applications and the Sun RMI protocol with the native Manta compiler. The results show that Manta's null-RMI latency is still eight times lower than for the compiled Sun RMI protocol and that Manta's efficient RMI protocol results in 1.8 to 3.4 times higher speedups for four out of six applications. Jason Maassen, Rob van Nieuwpoort, Ronald Veldema, Henri E. Bal, Thilo Kielmann, Ceriel J. H. Jacobs, Rutger F. H. Hofman |
ACM Trans. Program. Lang. Syst. | 4 |
| 2000 | Evaluating Design Alternatives for Reliable Communication on High-Speed NetworksabstractWe systematically evaluate the performance of five implementations of a single, user-level communication interface. Each implementation makes different architectural assumptions about the reliability of the network hardware and the capabilities of the network interface. The implementations differ accordingly in their division of protocol tasks between host software, network-interface firmware, and network hardware. Using microbenchmarks, parallel-programming systems, and parallel applications, we assess the performance impact of different protocol decompositions. We show how moving protocol tasks to a relatively slow network interface yields both performance advantages and disadvantages, depending on the characteristics of the application and the underlying parallel-programming system. In particular, we show that a communication system that assumes highly reliable network hardware and that uses network-interface support to process multicast traffic performs best for all applications. Raoul Bhoedjang, Kees Verstoep, Tim Rühl, Henri E. Bal, Rutger F. H. Hofman |
ASPLOS | 4 |
| 2000 | Programming Support for Distributed Clustercomputing
Henri E. Bal |
CLUSTER | 1 |
| 2000 | Satin: Efficient Parallel Divide-and-Conquer in Java
Rob van Nieuwpoort, Thilo Kielmann, Henri E. Bal |
Euro-Par | 3 |
| 2000 | CAVEStudy: An Infrastructure for Computational Steering in Virtual Reality EnvironmentsabstractWe present the CAVEStudy system that enables scientists to interactively steer a simulation from a virtual reality (VR) environment. No modification to the source code is necessary. CAVEStudy allows interactive and immersive analysis of a simulation running on a remote computer. Using a high-level description of the simulation, the system generates the communication layer (based on CAVERN-Soft) needed to control the execution and to gather data at runtime. We describe three case-studies implemented with CAVEStudy: soccer simulation, diode laser simulation and molecular dynamics. Luc Renambot, Henri E. Bal, Desmond Germans, Hans J. W. Spoelder |
HPDC | 2 |
| 2000 | Bandwidth-Efficient Collective Communication for Clustered Wide Area SystemsabstractMetacomputing infrastructures couple multiple clusters (or MPPs) via wide-area networks. A major problem in programming parallel applications for such platforms is their hierarchical network structure: latency and bandwidth of WANs often are orders of magnitude worse than those of local networks. Our goal is to optimize MPI's collective operations for such platforms. In this paper we focus on optimized utilization of the (scarce) wide-area bandwidth. We use two techniques: selecting suitable communication graph shapes, and splitting messages into multiple segments that are sent in parallel over different WAN links. To determine the best graph shape and segment size, we introduce a performance model called parameterized LogP (P-LogP), a hierarchical extension of the LogP model that covers messages of arbitrary length. With P-LogP, the optimal segment size and the best broadcast tree shape can be determined at runtime. (For conciseness, we restrict our discussion to the broadcast operation). An experimental performance evaluation shows that the new broadcast has significantly improved performance (for large messages) and that there is a close match between the theoretical model and the measured completion times. Thilo Kielmann, Henri E. Bal, Sergei Gorlatch |
IPDPS | 2 |
| 2000 | Man Multi-agent Interaction in VR: A Case Study with RoboCupabstractWe discuss the use of virtual reality (VR) techniques for interaction between humans and a multi-agent system in the context of RoboCup. The goal of RoboCup is to let teams of cooperating autonomous agents play a soccer match, using either robots or simulated players. We use RoboCup to study distributed collaborative applications, which allow multiple users at different geographic locations to cooperate, by interacting in real time through a shared simulation program. Our objective is to construct a VR environment in which humans at different locations can play along with a running RoboCup simulation in a natural way. The simulation system consists of the Soccer Server and a set of processes modeling the players. The server keeps track of the state of the game: provides the players with information on the game, and enforces the rules. The players request state information and autonomously calculate a behavior, sending the server commands that consist of accelerations, turns and kicks. The server discretizes time into slots, only one command is executed per time slot. We have developed a 3D visualization system that allows a user in a CAVE to interact with the soccer simulation software. Hans J. W. Spoelder, Luc Renambot, Desmond Germans, Henri E. Bal, Frans C. A. Groen |
VR | 4 |
| 2000 | Wide-area parallel programming using the remote method invocation modelabstractJava's support for parallel and distributed processing makes the language attractive for metacomputing applications, such as parallel applications that run on geographically distributed (wide-area) systems. To obtain actual experience with a Java-centric approach to metacomputing, we have built and used a high-performance wide-area Java system, called Manta. Manta implements the Java Remote Method Invocation (RMI) model using different communication protocols (active messages and TCP/IP) for different networks. The paper shows how wide-area parallel applications can be expressed and optimized using Java RMI. Also, it presents performance results of several applications on a wide-area system consisting of four Myrinet-based clusters connected by ATM WANs. We finally discuss alternative programming models, namely object replication, JavaSpaces, and MPI for Java. Copyright © 2000 John Wiley & Sons, Ltd. Rob van Nieuwpoort, Jason Maassen, Henri E. Bal, Thilo Kielmann, Ronald Veldema |
Concurr. Pract. Exp. | 3 |
| 1999 | Sensitivity of Parallel Applications to Large Differences in Bandwidth and Latency in Two-Layer InterconnectsabstractThis paper studies application performance on systems with strongly non-uniform remote memory access. In current generation NUMAs the speed difference between the slowest and fastest link in an interconnect-the "NUMA gap"-is typically less than an order of magnitude, and many conventional parallel programs achieve good performance. We study how different NUMA gaps influence application performance, up to and including typical wide-area latencies and bandwidths. We find that for gaps larger than those of current generation NUMAs, performance suffers considerably (for applications that were designed for a uniform access interconnect). For many applications, however, performance can be greatly improved with comparatively simple changes: traffic over slow links can be reduced by making communication patterns hierarchical-like the interconnect. We find that in four out of our six applications the size of the gap can be increased by an order of magnitude or more without severely impacting speedup. We analyze why the improvements are needed, why they work so well, and how much non-uniformity they can mask. Aske Plaat, Henri E. Bal, Rutger F. H. Hofman |
HPCA | 2 |
| 1999 | MagPIe: MPI's Collective Communication Operations for Clustered Wide Area SystemsabstractWriting parallel applications for computational grids is a challenging task. To achieve good performance, algorithms designed for local area networks must be adapted to the differences in link speeds. An important class of algorithms are collective operations, such as broadcast and reduce. We have developed MAGPIE, a library of collective communication operations optimized for wide area systems. MAGPIE's algorithms send the minimal amount of data over the slow wide area links, and only incur a single wide area latency. Using our system, existing MPI applications can be run unmodified on geographically distributed systems. On moderate cluster sizes, using a wide area latency of 10 milliseconds and a bandwidth of 1 MByte/s, MAGPIE executes operations up to 10 times faster than MPICH, a widely used MPI implementation; application kernels improve by up to a factor of 4. Due to the structure of our algorithms, MAGPIE's advantage increases for higher wide area latencies. Thilo Kielmann, Rutger F. H. Hofman, Henri E. Bal, Aske Plaat, Raoul Bhoedjang |
PPoPP | 3 |
| 1999 | An Efficient Implementation of Java's Remote Method InvocationabstractJava offers interesting opportunities for parallel computing. In particular, Java Remote Method Invocation provides an unusually flexible kind of Remote Procedure Call. Unlike RPC, RMI supports polymorphism, which requires the system to be able to download remote classes into a running application. Sun's RMI implementation achieves this kind of flexibility by passing around object type information and processing it at run time, which causes a major run time overhead. Using Sun's JDK 1.1.4 on a Pentium Pro/Myri.net cluster, for example, the latency for a null RMI (without parameters or a return value) is 1228 μsec, which is about a factor of 40 higher than that of a user-level RPC. In this paper, we study an alternative approach for implementing RMI, based on native compilation. This approach allows for better optimization, eliminates the need for processing of type information at run time, and makes a light weight communication protocol possible. We have built a Java system based on a native compiler, which supports both compile time and run time generation of marshallers. We find that almost all of the run time overhead of RMI can be pushed to compile time. With this approach, the latency of a null RMI is reduced to 34 μsec, while still supporting polymorphic RMIs (and allowing interoperability with other JVMs). Jason Maassen, Rob van Nieuwpoort, Ronald Veldema, Henri E. Bal, Aske Plaat |
PPoPP | 4 |
| 1998 | Challenging Applications on Fast NetworksabstractParallel computing on clusters of workstations is attractive because of the low costs in comparison to MPPs, but the speed of the local area network limits the class of applications that can be run efficiently. Fortunately, faster network technology is becoming available for the next generation of workstation clusters. This paper studies the effect of running challenging applications that communicate heavily on three types of modern interconnects: 100 Mbit/s Fast Ethernet, 155 Mbit/s ATM, and 1.28 Gbit/s Myrinet. Experimental results show that even challenging communication-intensive applications can achieve acceptable performance on workstation clusters, but only if the communication software has been designed and tuned for high performance. Koen Langendoen, Rutger F. H. Hofman, Henri E. Bal |
HPCA | 3 |
| 1998 | Efficient Multicast on Myrinet using Link-Level Flow ControlabstractThis paper studies the implementation of efficient multicast protocols for Myrinet, a switched, wormhole-routed Gigabit-per-second network technology. Since Myrinet does not support multicasting in hardware, multicast services must be implemented in software. We present a new, efficient, and reliable software multicast protocol that uses the network interface to efficiently forward multicast traffic. The new protocol is constructed on top of reliable, flow-controlled channels between pairs of network interfaces. We describe the design of the protocol and make a detailed comparison with a previous multicast protocol. We show that our protocol is simpler and scales better than the previous protocol. This claim is supported by extensive performance measurements on a 64-node Myrinet cluster. Raoul Bhoedjang, Tim Rühl, Henri E. Bal |
ICPP | 3 |
| 1998 | Optimizing Distributed Data Structures using Application-Specific Network Interface SoftwareabstractNetwork interfaces that contain a programmable processor offer much flexibility, which so far has mainly been used to optimize message passing libraries. We show that high performance gains can be achieved by implementing support for application-specific shared data structures on the network interface processors. As a case study, we have implemented shared transposition tables on a Myrinet network, using customized software that runs partly on the network processor and partly on the host. The customized software greatly reduces the overhead of interactions between the network interface and the host. Also, the software exploits application semantics to obtain a simple and efficient communication protocol. Performance measurements indicate that applications that run application-specific code on the network interface are up to 2.5 times as fast as those that use generic message-passing software. Raoul Bhoedjang, John W. Romein, Henri E. Bal |
ICPP | 3 |
| 1998 | Parallel simulation of ion recombination in nonpolar liquids
Frank J. Seinstra, Henri E. Bal, Hans J. W. Spoelder |
Future Gener. Comput. Syst. | 2 |
| 1998 | Performance Evaluation of the Orca Shared-Object SystemabstractOrca is a portable, object-based distributed shared memory (DSM) system. This article studies and evaluates the design choices made in the Orca system and compares Orca with other DSMs. The article gives a quantitative analysis of Orca's coherence protocol (based on write-updates with function shipping), the totally ordered group communication protocol, the strategy for object placement, and the all-software, user-space architecture. Performance measurements for 10 parallel applications illustrate the trade-offs made in the design of Orca and show that essentially the right design decisions have been made. A write-update protocol with function shipping is effective for Orca, especially since it is used in combination with techniques that avoid replicating objects that have a low read/write ratio. The overhead of totally ordered group communication on application performance is low. The Orca system is able to make near-optimal decisions for object placement and replication. In addition, the article compares the performance of Orca with that of a page-based DSM (TreadMarks) and another object-based DSM (CRL). It also analyzes the communication overhead of the DSMs for several applications. All performance measurements are done on a 32-node Pentium Pro cluster with Myrinet and Fast Ethernet networks. The results show that Orca programs send fewer messages and less data than the TreadMarks and CRL programs and obtain better speedups. Henri E. Bal, Raoul Bhoedjang, Rutger F. H. Hofman, Ceriel J. H. Jacobs, Koen Langendoen, Tim Rühl |
ACM Trans. Comput. Syst. | 1 |
| 1998 | A Task- and Data-Parallel Programming Language Based on Shared ObjectsabstractMany programming languages support either task parallelism, but few languages provide a uniform framework for writing applications that need both types of parallelism or data parallelism. We present a programming language and system that integrates task and data parallelism using shared objects. Shared objects may be stored on one processor or may be replicated. Objects may also be partitioned and distributed on several processors.Task parallelism is achieved by forking processes remotely and have them communicate and synchronize through objects. Data parallelism is achieved by executing operations on partitioned objects in parallel. Writing task-and data-parallel applications with shared objects has several advantages. Programmers use the objects as if they were stored in a memory common to all processors. On distributed-memory machines, if objects are remote, replicated, or partitioned, the system takes care of many low-level details such as data transfers and consistency semantics. In this article, we show how to write task-and data-parallel programs with our shared object model. We also desribe a portable implementation of the model. To assess the performance of the system, we wrote several applications that use task and data parallelism and excuted them on a collection of Pentium Pros connected by Myrinet. The performance of these applications is also discussed in this article. Saniya Ben Hassen, Henri E. Bal, Ceriel J. H. Jacobs |
ACM Trans. Program. Lang. Syst. | 2 |
| 1997 | Performance of a High-Level Parallel Language on a High-Speed Network
Henri E. Bal, Raoul Bhoedjang, Rutger F. H. Hofman, Ceriel J. H. Jacobs, Koen Langendoen, Tim Rühl, Kees Verstoep |
J. Parallel Distributed Comput. | 1 |
| 1996 | Integrating Task and Data Parallelism Using Shared ObjectsabstractSupporting both task and data parallelism in one programming system is useful, since many applications need both types of parallelism. We present a programming model that integrates task and data parallelism using shared objects. The model is a generalization of shared objects in Orca. Orca is a task parallel language that uses shared objects for communication between processes and for storing shared (possibly replicated) data. Our new model also uses shared objects for partitioning of shared data and for distribution of work in a data parallel way. Data parallelism is introduced by executing operations on a partitioned object in parallel. The paper describes the design of the new model, its implementation, and its usage for parallel applications that use mixed task and data parallelism. 1 Introduction Most parallel programming systems are based either on data parallelism or on task parallelism. The advantage of data parallelism is that it is easy to use. The programmer merely specifi... Saniya Ben Hassen, Henri E. Bal |
International Conference on Supercomputing | 2 |
| 1996 | A Flexible Operation Execution Model for Shared Distributed ObjectsabstractMany parallel and distributed programming models are based on some form of shared objects, which may be represented in various ways (e.g., single-copy, replicated, and partitioned objects). Also, many different operation execution strategies have been designed for each representation. In programming systems that use multiple representations integrated in a single object model, one way to provide multiple execution strategies is to implement each strategy independently from the others. However, this leads to rigid systems and provides little opportunity for code reuse. Instead, we propose a flexible operation execution model that allows the implementation of many different strategies, which can even be changed at runtime. We present the model and a distributed implementation of it. Also, we describe how various execution strategies can be expressed using the model, and we look at applications that benefit from its flexibility. 1 Introduction Shared objects have become a popular model... Saniya Ben Hassen, Irina Athanasiu, Henri E. Bal |
OOPSLA | 3 |
| 1995 | Parallel N-Body Simulation on a Large-Scale Homogeneous Distributed System
John W. Romein, Henri E. Bal |
Euro-Par | 2 |
| 1995 | Comparing Kernel-Space and User-Space Communication Protocols on AmoebaabstractMost distributed systems contain protocols for reliable communication, which are implemented either in the microkernel or in user space. In the latter case, the microkernel provides only low-level, unreliable primitives and the higher-level protocols are implemented as a library in user space. This approach is more flexible but potentially less efficient. We study the impact on performance of this choice for RPC and group communication protocols on Amoeba. An important goal in this paper is to look at overall system performance. For this purpose, we use several (communication-intensive) parallel applications written in Orca. We look at two implementations of Orca on Amoeba, one using Amoeba's kernel-space protocols and one using user-space protocols built on top of Amoeba's low-level FLIP protocol. The results show that comparable performance can be obtained with user-space protocols. Marco Oey, Koen Langendoen, Henri E. Bal |
ICDCS | 3 |
| 1995 | Parallel Retrograde Analysis on a Distributed SystemabstractRetrograde Analysis (RA) is an AI search technique used to compute endgame databases, which contain optimal solutions for part of the search space of a game. RA has been applied successfully to several games, but its usefulness is restricted by the huge amount of CPU time and internal memory it requires. We present a parallel distributed algorithm for RA that addresses these problems. RA is hard to parallelize efficiently, because the communication overhead potentially is enormous. We show that the overhead can be reduced drastically using message combining. We implemented the algorithm on an Ethernet-based distributed system. For one example game (awari), we have computed a large database in 50 minutes on 64 processors, whereas one machine took 40 hours (a speedup of 48). An even larger database (computed in 20 hours) would have required over 600 MByte of internal memory on a uniprocessor and would compute for many weeks. Henri E. Bal, L. Victor Allis |
SC | 1 |
| 1995 | Introduction to the Special Section
Henri E. Bal, Boumediene Belkhouche, Mary Lou Soffa |
IEEE Trans. Software Eng. | 1 |
| 1994 | Object-based approach to programming distributed systemsabstractAbstract Two kinds of parallel computers exist: those with shared memory and those without. The former are difficult to build but easy to program. The latter are easy to build but difficult to program. In this paper we present a hybrid model that combines the best properties of each by simulating a restricted object‐based shared memory on machines that do not share physical memory. In this model, objects can be replicated on multiple machines. An operation that does not change an object can then be done locally, without any network traffic. Update operations can be done using the reliable broadcast protocol described in the paper. We have constructed a prototype system, designed and implemented a new programming language for it, and programmed various applications using it. The model, algorithms, language, applications and performance will be discussed. Andrew S. Tanenbaum, Henri E. Bal, Saniya Ben Hassen, M. Frans Kaashoek |
Concurr. Pract. Exp. | 2 |
| 1993 | Programming a Distributed System Using Shared ObjectsabstractBuilding the hardware for a high-performance distributed computer system is a lot easier than building its software. The authors describe a model for programming distributed systems based on abstract data types that can be replicated on all machines that need them. Read operations are done locally, without requiring network traffic. Writes can be done using a reliable broadcast algorithm if the hardware supports broadcasting; otherwise, a point-to-point protocol is used. The authors have built such a system based on the Amoeba microkernel, and implemented a language, Orca, on top of it. For Orca applications that have a high ratio of reads to writes, they measure good speedups on a system with 16 processors.> Andrew S. Tanenbaum, Henri E. Bal, M. Frans Kaashoek |
HPDC | 2 |
| 1993 | Object Distribution in Orca using Compile-Time and Run-Time TechniquesabstractOrca is a language for parallel programming on distributed systems. Communication in Orca is based on shared data-objects, which is a form of distributed shared memory. The performance of Orca programs depends strongly on how shared dataobjects are distributed among the local physical memories of the processors. This paper studies a new and efficient solution to this problem, based on an integration of compile-time and run-time techniques. The Orca compiler has been extended to determine the access patterns of processes to shared objects. The compiler passes a summary of this information to the run-time system, which uses it to make good decisions about which objects to replicate and where to store nonreplicated objects. Measurements show that the new system gives better overall performance than any previous implementation of Orca. 3333333333333333 1 This research was supported in part by a PIONIER grant from the Netherlands Organization for Scientific Research (N.W.O.). 2 This re... Henri E. Bal, M. Frans Kaashoek |
OOPSLA | 1 |
| 1993 | Evaluation of KL1 and the inference machine
Henri E. Bal |
Future Gener. Comput. Syst. | 1 |
| 1992 | Fault-tolerant parallel programming in ArgusabstractAbstract Fault tolerance is an issue ignored in most parallel languages. The overhead of making parallel, high‐performance programs resilient to processor crashes is often too high, given the low probability of such events. If parallel systems become more large‐scaled, however, processor failures will become likely, so they should be dealt with. Two approaches to this problem are feasible. First, the system can make programs fault‐tolerant transparently. It can log messages, make checkpoints, and so on. Second, the programmer can write explicit code for handling failures in an application‐specific way. The latter approach is potentially more efficient, but also requires more work from the programmer. In this paper, we intend to get some initial insight into how hard and efficient explicit fault‐tolerant parallel programming is. We do so by implementing four parallel applications in Argus, a language supporting parallelism as well as fault tolerance. Our experiences indicate that the extra effort needed for fault tolerance varies much between different applications. Also, trade‐offs can frequently be made between programming effort and efficiency. One lesson we learned is that fault tolerance should not be added as an afterthought, but is best taken into account from the start. As another result, the ability to integrate transparent and explicit mechanisms for fault tolerance would sometimes be highly useful. Henri E. Bal |
Concurr. Pract. Exp. | 1 |
| 1992 | Replication techniques for speeding up parallel applications on distributed systemsabstractAbstract Most methods for programming loosely coupled systems are based on message‐passing. Recently, however, methods have emerged based on ‘virtually’ sharing data. These methods simplify distributed programming, but are hard to implement efficiently, as loosely coupled systems do not contain physical shared memory. We introduce a new model,the shared data‐object model, that eases the implementation of parallel applications on loosely coupled systems, but can still be implemented efficiently. In our model, shared data are encapsulated in passive data‐objects, which are variables of user‐defined abstract data types. To speed up access to shared data, data‐objects are replicated. This ability to replicate objects is a significant difference with other object‐based models (e.g. Emerald and Amber). Also, by replicating logical objects rather than physical pages, our model has many advantages over shared virtual memory systems. This paper discusses the design choices involved in replicating objects and their effect on performance. Important issues are: how to maintain consistency among different copies of an object; how to implement changes to objects; which strategy for object replication to use. We have implemented several options to determine which ones are the most efficient. Henri E. Bal, M. Frans Kaashoek, Andrew S. Tanenbaum, Jack Jansen 0001 |
Concurr. Pract. Exp. | 1 |
| 1992 | A comparative study of five parallel programming languages
Henri E. Bal |
Future Gener. Comput. Syst. | 1 |
| 1992 | A Comparison of Two Paradigms for Distributed Shared MemoryabstractAbstract Two paradigms for distributed shared memory on loosely‐coupled computing systems are compared: the shared data‐object model as used in Orca, a programming language specially designed for loosely‐coupled computing systems, and the shared virtual memory model. For both paradigms two systems are described, one using only point‐to‐point messages, the other using broadcasting as well. The two paradigms and their implementations are described briefly. Their performances are compared on four applications: the travelling‐salesman problem, alpha‐beta search, matrix multiplication and the all‐pairs shortest‐paths problem. Measurements were obtained on a system consisting of 10 MC68020 processors connected by an Ethernet. For comparison purposes, the applications have also been run on a system with physical shared memory. In addition, the paper gives measurements for the first two applications above when remote procedure call is used as the communication mechanism. The measurements show that both paradigms can be used efficiently for programming large‐grain parallel applications, with significant speed‐ups. The structured shared data‐object model achieves the highest speed‐ups and is easiest to program and to debug. M. Frans Kaashoek, Henri E. Bal, Andrew S. Tanenbaum |
Softw. Pract. Exp. | 2 |
| 1992 | Orca: A Language For Parallel Programming of Distributed SystemsabstractA detailed description is given of the Orca language design and the design choices are discussed. Orca is intended for applications programmers rather than systems programmers. This is reflected in its design goals to provide a simple, easy-to-use language that is type-secure and provides clean semantics. Three example parallel applications in Orca, one of which is described in detail, are discussed. One of the existing implementations, which is based on reliable broadcasting, is described. Performance measurements of this system are given for three parallel applications. The measurements show that significant speedups can be obtained for all three applications. The authors compare Orca with several related languages and systems.> Henri E. Bal, M. Frans Kaashoek, Andrew S. Tanenbaum |
IEEE Trans. Software Eng. | 1 |
| 1991 | Distributed Programming with Shared Data
Henri E. Bal, Andrew S. Tanenbaum |
Comput. Lang. | 1 |
| 1991 | The Amoeba distributed operating system - A status report
Andrew S. Tanenbaum, M. Frans Kaashoek, Robbert van Renesse, Henri E. Bal |
Comput. Commun. | 4 |
| 1991 | Heuristic search in PARLOG using replicated worker style parallelism
Henri E. Bal |
Future Gener. Comput. Syst. | 1 |
| 1986 | Language- and Machine-Independent Global Optimization on Intermediate Code
Henri E. Bal, Andrew S. Tanenbaum |
Comput. Lang. | 1 |