EDBT 2026 Demo / reviewers in the wild / expert
Rob van Nieuwpoort
dblp:n/RvNieuwpoort · also Rob V. van Nieuwpoort
· DBLP profile ↗
45ranked-venue papers
7as first author
12since 2021 · last 2026
0000-0002-2947-9444ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 33 · 6 first-author · 6 since 2021Software engineering, systems software and programming languages · 10 · 1 first-author · 5 since 2021Applied, interdisciplinary, general and emerging computing · 6 · 2 since 2021
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | Velvet: Parallel Divide-and-Conquer in Safe Rust
Anna Badia Liokouras, Ben van Werkhoven, Rob van Nieuwpoort |
Euro-Par (1) | 3 |
| 2026 | RSMM: A focus area maturity model for research software projectsabstractContext: Research software is instrumental in producing research results. It plays a special role in advancing scientific discovery through tasks such as data analysis, simulation, and visualisation. However, despite its importance, the organisations that produce research software face challenges in managing research software projects. Existing frameworks for software engineering often overlook the needs of research software, including research software engineering practices and open science principles. Without clear guidance, organisations that produce research software must develop and invent new techniques for research software engineering, which is a slow and costly process. Objective: This work presents RSMM, a maturity model designed to improve organisational practices in research software project management. Methods: The initial version of RSMM was developed through a systematic literature review. Expert interviews were then conducted to evaluate and refine the model. Finally, multiple case studies were carried out to validate RSMM and demonstrate its applicability in real-world settings. Results: The final version of RSMM (v1.0) comprises 79 best practices, 17 capabilities, and 10 maturity levels, organised into 4 focus areas: ‘Software Project Management’ , ‘Research Software Management’ , ‘Community Engagement’ , and ‘Software Adoptability’ . We provide a comprehensive analysis of RSMM v1.0 and demonstrate its practical applicability through two illustrative case studies. Conclusion: The RSMM is designed to help organisations in evaluating and improving their research software project management by assessing a project’s current maturity level and providing best practices across four focus areas to guide its progression. Deekshitha, Rena Bakhshi, Jason Maassen, Carlos Martinez-Ortiz, Rob van Nieuwpoort, Antti Knutas, Slinger Jansen |
Inf. Softw. Technol. | 5 |
| 2025 | Tuning the Tuner: Introducing Hyperparameter Optimization for Auto-TuningabstractAutomatic performance tuning (auto-tuning) is widely used to optimize performance-critical applications across many scientific domains by finding the best program variant among many choices. Efficient optimization algorithms are crucial for navigating the vast and complex search spaces in autotuning. As is well known in the context of machine learning and similar fields, hyperparameters critically shape optimization algorithm efficiency. Yet for auto-tuning frameworks, these hyperparameters are almost never tuned, and their potential performance impact has not been studied.We present a novel method for general hyperparameter tuning of optimization algorithms for auto-tuning, thus "tuning the tuner". In particular, we propose a robust statistical method for evaluating hyperparameter performance across search spaces, publish a FAIR data set and software for reproducibility, and present a simulation mode that replays previously recorded tuning data, lowering the costs of hyperparameter tuning by two orders of magnitude. We show that even limited hyperparameter tuning can improve auto-tuner performance by 94.8% on average, and establish that the hyperparameters themselves can be optimized efficiently with meta-strategies (with an average improvement of 204.7%), demonstrating the often overlooked hyperparameter tuning as a powerful technique for advancing auto-tuning research and practice. Floris-Jan Willemsen, Rob van Nieuwpoort, Ben van Werkhoven |
eScience | 2 |
| 2025 | Efficient Construction of Large Search Spaces for Auto-Tuning
Floris-Jan Willemsen, Rob van Nieuwpoort, Ben van Werkhoven |
ICPP | 2 |
| 2025 | Multi-Strided Access Patterns to Boost Hardware PrefetchingabstractImportant memory-bound kernels, such as linear algebra, convolutions, and stencils, rely on SIMD instructions as well as optimizations targeting improved vectorized data traversal and data re-use to attain satisfactory performance. On contemporary CPU architectures, the hardware prefetcher is of key importance for efficient utilization of the memory hierarchy. In this paper, we demonstrate that transforming a memory access pattern consisting of a single stride to one that concurrently accesses multiple strides, can boost the utilization of the hardware prefetcher, and in turn improves the performance of memory-bound kernels significantly. Using a set of micro-benchmarks, we establish that accessing memory in a multi-strided manner enables more cache lines to be concurrently brought into the cache, resulting in improved cache hit ratios and higher effective memory bandwidth without the introduction of costly software prefetch instructions. Subsequently, we show that multi-strided variants of a collection of six memory-bound dense compute kernels outperform state-of-the-art counterparts on three different micro-architectures. More specifically, for kernels among which Matrix Vector Multiplication, Convolution Stencil and kernels from PolyBench, we achieve significant speedups of up to 12.55x over Polly, 2.99x over MKL, 1.98x over OpenBLAS, 1.08x over Halide and 1.87x over OpenCV. The code transformation to take advantage of multi-strided memory access is a natural extension of the loop unroll and loop interchange techniques, allowing this method to be incorporated into compiler pipelines in the future. Miguel O. Blom, Kristian F. D. Rietveld, Rob van Nieuwpoort |
ICPE | 3 |
| 2024 | In Serverless, OS Scheduler Choice Costs Money: A Hybrid Scheduling Approach for Cheaper FaaSabstractIn Function-as-a-Service (FaaS) serverless, large applications are split into short-lived stateless functions. Deploying functions is mutually profitable: users need not be concerned with resource management, while providers can keep their servers at high utilization rates running thousands of functions concurrently on a single machine. It is exactly this high concurrency that comes at a cost. The standard Linux Completely Fair Scheduler (CFS) switches often between tasks, which leads to prolonged execution times. We present evidence that relying on the default Linux CFS scheduler increases serverless workloads cost by up to 10×. Yuxuan Zhao 0003, Weikang Weng, Rob van Nieuwpoort, Alexandru Uta |
Middleware | 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 | 5 |
| 2024 | A methodology for comparing optimization algorithms for auto-tuningabstractAdapting applications to optimally utilize available hardware is no mean feat: the plethora of choices for optimization techniques are infeasible to maximize manually. To this end, auto-tuning frameworks are used to automate this task, which in turn use optimization algorithms to efficiently search the vast searchspaces. However, there is a lack of comparability in studies presenting advances in auto-tuning frameworks and the optimization algorithms incorporated. As each publication varies in the way experiments are conducted, metrics used, and results reported, comparing the performance of optimization algorithms among publications is infeasible. The auto-tuning community identified this as a key challenge at the 2022 Lorentz Center workshop on auto-tuning. The examination of the current state of the practice in this paper further underlines this. We propose a community-driven methodology composed of four steps regarding experimental setup, tuning budget, dealing with stochasticity, and quantifying performance. This methodology builds upon similar methodologies in other fields while taking into account the constraints and specific characteristics of the auto-tuning field, resulting in novel techniques. The methodology is demonstrated in a simple case study that compares the performance of several optimization algorithms used to auto-tune CUDA kernels on a set of modern GPUs. We provide a software tool to make the application of the methodology easy for authors, and simplifies reproducibility of results. Floris-Jan Willemsen, Richard Schoonhoven, Jiri Filipovic, Jacob Odgård Tørring, Rob van Nieuwpoort, Ben van Werkhoven |
Future Gener. Comput. Syst. | 5 |
| 2023 | FAIRSECO: An Extensible Framework for Impact Measurement of Research SoftwareabstractThe growing usage of research software in the research community has highlighted the need to recognize and acknowledge the contributions made not only by researchers but also by Research Software Engineers. However, the existing methods for crediting research software and Research Software Engineers have proven to be insufficient. In response, we have developed FAIRSECO, an extensible open source framework with the objective of assessing the impact of research software in research through the evaluation of various factors. The FAIRSECO framework addresses two critical information needs: firstly, it provides potential users of research software with metrics related to software quality and FAIRness. Secondly, the framework provides information for those who wish to measure the success of a project by offering impact data. By exploring the quality and impact of research software, our aim is to ensure that Research Software Engineers receive the recognition they deserve for their valuable contributions. Deekshitha, Siamak Farshidi, Jason Maassen, Rena Bakhshi, Rob van Nieuwpoort, Slinger Jansen |
e-Science | 5 |
| 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 | 4 |
| 2022 | mCAP: Memory-Centric Partitioning for Large-Scale Pipeline-Parallel DNN Training
Henk Dreuning, Henri E. Bal, Rob van Nieuwpoort |
Euro-Par | 3 |
| 2022 | Lightning: Scaling the GPU Programming Model Beyond a Single GPUabstractThe GPU programming model is primarily aimed at the development of applications that run one GPU. However, this limits the scalability of GPU code to the capabilities of a single GPU in terms of compute power and memory capacity. To scale GPU applications further, a great engineering effort is typically required: work and data must be divided over multiple GPUs by hand, possibly in multiple nodes, and data must be manually spilled from GPU memory to higher-level memories. We present Lightning: a framework that follows the common GPU programming paradigm but enables scaling to large problems with ease. Lightning supports multi-GPU execution of GPU kernels, even across multiple nodes, and seamlessly spills data to higher-level memories (main memory and disk). Existing CUDA kernels can easily be adapted for use in Lightning, with data access annotations on these kernels allowing Lightning to infer their data requirements and the dependencies between subsequent kernel launches. Lightning efficiently distributes the work/data across GPUs and maximizes efficiency by overlapping scheduling, data movement, and kernel execution when possible. We present the design and implementation of Lightning, as well as experimental results on up to 32 GPUs for eight benchmarks and one real-world application. Evaluation shows excellent performance and scalability, such as a speedup of 57.2 x over the CPU using Lighting with 16 GPUs over 4 nodes and 80 GB of data, far beyond the memory capacity of one GPU. Stijn Heldens, Pieter Hijma, Ben van Werkhoven, Jason Maassen, Rob van Nieuwpoort |
IPDPS | 5 |
| 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 | 6 |
| 2018 | Painting the Picture of Software Impact with the Research Software DirectoryabstractIn this lightning talk we will describe the Research Software Directory; a content management system that is tailored to research software with the goal of enabling a qualitative assessment of software impact and improving software findability. Jurriaan H. Spaaks, Tom Klaver, Stefan Verhoeven, Jason Maassen, Tom Bakker, Atze van der Ploeg, Ben van Werkhoven, Willem Robert van Hage, Rob van Nieuwpoort |
eScience | 9 |
| 2018 | Poster Abstracts eScience 2018 ConferenceabstractThe eScience conference aims to bring together leading international researchers and research software engineers from all disciplines to present and discuss how digital technology impacts scientific research. There were many poster abstracts submitted to the conference this year, and we have also invited several authors of full papers submitted to the conference to submit their work for a poster presentation. We are very happy with the selection of posters that have been confirmed for poster presentation at the conference. The diverse set of topics covered by the posters reflects the broad impact of eScience in various domains, as well as the high quality technical work that is performed by the various teams of researchers. Posters are a great way of presenting work at a conference that really encourage discussions and interactions among the participants. We look forward to the poster sessions at this year’s eScience conference. Ben van Werkhoven, Adriënne Mendrik, Rob van Nieuwpoort |
eScience | 3 |
| 2016 | The landscape of GPGPU performance modeling tools
Souley Madougou, Ana Lucia Varbanescu, Cees T. A. M. de Laat, Rob van Nieuwpoort |
Parallel Comput. | 4 |
| 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 | 3 |
| 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 | 3 |
| 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. | 2 |
| 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 | 5 |
| 2012 | Radio Astronomy Beam Forming on Many-Core ArchitecturesabstractTraditional radio telescopes use large steel dishes to observe radio sources. The largest radio telescope in the world, LOFAR, uses tens of thousands of fixed, omni-directional antennas instead, a novel design that promises ground-breaking research in astronomy. Where traditional tele-scopes use custom-built hardware, LOFAR uses software to do signal processing in real time. This leads to an instrument that is inherently more flexible. However, the enormous data rates and processing requirements (tens to hundreds of teraflops) make this extremely challenging. The next-generation telescope, the SKA, will require exaflops. Unlike traditional instruments, LOFAR and SKA can observe in hundreds of directions simultaneously, using beam forming. This is useful, for example, to search the sky for pulsars (i.e. rapidly rotating highly magnetized neutron stars). Beam forming is an important technique in signal processing: it is also used in WIFI and 4G cellular networks, radar systems, and health-care microwave imaging instruments. We propose the use of many-core architectures, such as 48-core CPU systems and Graphics Processing Units (GPUs), to accelerate beam forming. We use two different frameworks for GPUs, CUDA and Open CL, and present results for hardware from different vendors (i.e. AMD and NVIDIA). Additionally, we implement the LOFAR beam former on multi-core CPUs, using Open MP with SSE vector instructions. We use auto-tuning to support different architectures and implementation frameworks, achieving both platform and performance portability. Finally, we compare our results with the production implementation, written in assembly and running on an IBM Blue Gene/P supercomputer. We compare both computational and power efficiency, since power usage is one of the fundamental challenges modern radio telescopes face. Compared to the production implementation, our auto-tuned beam former is 45-50 times faster on GPUs, and 2-8 times more power efficient. Our experimental results lead to the conclusion that GPUs are an attractive solution to accelerate beam forming. Alessio Sclocco, Ana Lucia Varbanescu, Jan David Mol, Rob van Nieuwpoort |
IPDPS | 4 |
| 2012 | Generating synchronization statements in divide-and-conquer programs
Pieter Hijma, Rob van Nieuwpoort, Ceriel J. H. Jacobs, Henri E. Bal |
Parallel Comput. | 2 |
| 2011 | Introduction
Rosa M. Badia, Fabrice Huet, Rob van Nieuwpoort, Rainer Keller |
Euro-Par (1) | 3 |
| 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. | 2 |
| 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. | 2 |
| 2010 | The LOFAR correlator: implementation and performance analysisabstractLOFAR is the first of a new generation of radio telescopes.Rather than using expensive dishes, it forms a distributed sensor network that combines the signals from many thousands of simple antennas. Its revolutionary design allows observations in a frequency range that has hardly been studied before. John W. Romein, P. Chris Broekema, Jan David Mol, Rob van Nieuwpoort |
PPoPP | 4 |
| 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. | 1 |
| 2009 | Using many-core hardware to correlate radio astronomy signalsabstractA recent development in radio astronomy is to replace traditional dishes with many small antennas. The signals are combined to form one large, virtual telescope. The enormous data streams are cross-correlated to filter out noise. This is especially challenging, since the computational demands grow quadratically with the number of data streams. Moreover, the correlator is not only computationally intensive, but also very I/O intensive. The LOFAR telescope, for instance, will produce over 100 terabytes per day. The future SKA telescope will even require in the order of exaflops, and petabits/s of I/O. A recent trend is to correlate in software instead of dedicated hardware. This is done to increase flexibility and to reduce development efforts. Examples include e-VLBI and LOFAR. Rob van Nieuwpoort, John W. Romein |
ICS | 1 |
| 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 | 5 |
| 2008 | Radioastronomy Image Synthesis on the Cell/B.E
Ana Lucia Varbanescu, Alexander S. van Amesfoort, Tim Cornwell, Andrew Mattingly, Bruce G. Elmegreen, Rob van Nieuwpoort, Ger van Diepen, Henk J. Sips |
Euro-Par | 6 |
| 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 | 2 |
| 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 | 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 | 1 |
| 2006 | Simple Locality-Aware Co-allocation in Peer-to-Peer Supercomputing
Niels Drost, Rob van Nieuwpoort, Henri E. Bal |
CCGRID | 2 |
| 2006 | Middleware adaptation with the Delphoi serviceabstractAbstract Grid middleware needs to adapt to changing resources for a large variety of operations. Currently, however, there is only low‐level information available about Grid resources, coming from various but functionally isolated monitoring and information systems. In this paper, we present the Delphoi service. It provides a unified interface to the necessary information, and also matches the abstraction level needed by middleware services to adapt their behavior. We describe Delphoi's architecture, the information it provides, and we evaluate the quality of its performance information. Delphoi has been developed as part of the EC‐funded GridLab project and is currently being deployed on the project's testbed for adding adaptivity to GridLab's middleware services. Copyright © 2006 John Wiley & Sons, Ltd. Jason Maassen, Rob van Nieuwpoort, Thilo Kielmann, Kees Verstoep, Mathijs den Burger |
Concurr. Comput. Pract. Exp. | 2 |
| 2005 | Developing Java Grid Applications with Ibis
Kees van Reeuwijk, Rob van Nieuwpoort, Henri E. Bal |
Euro-Par | 2 |
| 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. | 1 |
| 2005 | The Grid Application Toolkit: Toward Generic and Easy Application Programming Interfaces for the GridabstractCore Grid technologies are rapidly maturing, but there remains a shortage of real Grid applications. One important reason is the lack of a simple and high-level application programming toolkit, bridging the gap between existing Grid middleware and application-level needs. The Grid Application Toolkit (GAT), as currently developed by the EC-funded project GridLab, provides this missing functionality. As seen from the application, the GAT provides a unified simple programming interface to the Grid infrastructure, tailored to the needs of Grid application programmers and users. A uniform programming interface will be needed for application developers to create a new generation of "Grid-aware" applications. The GAT implementation handles both the complexity and the variety of existing Grid middleware services via so-called adaptors. Complementing existing Grid middleware, GridLab also provides high-level services to implement the GAT functionality. We present the GridLab software architecture, consisting of the GAT, environment-specific adaptors, and GridLab services. We elaborate the concepts underlying the GAT and outline the corresponding application programming interface. We present the functionality of GridLab's high-level services and demonstrate how a dynamic Grid application can easily benefit from the GAT. All GridLab software is open source and can be downloaded from the project Web site. Gabrielle Allen, Kelly Davis, Tom Goodale, Andrei Hutanu, Hartmut Kaiser, Thilo Kielmann, André Merzky, Rob van Nieuwpoort, Alexander Reinefeld, Florian Schintke, Thorsten Schütt, Edward Seidel, Brygg Ullmer |
Proc. IEEE | 8 |
| 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 | 2 |
| 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. | 4 |
| 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 | 1 |
| 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. | 2 |
| 2000 | Satin: Efficient Parallel Divide-and-Conquer in Java
Rob van Nieuwpoort, Thilo Kielmann, Henri E. Bal |
Euro-Par | 1 |
| 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. | 1 |
| 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 | 2 |