VLDB 2026 Research / reviewers in the wild / expert
M. Mustafa Rafique
dblp:32/5404 · also Muhammad Mustafa Rafique
· DBLP profile ↗
42ranked-venue papers
6as first author
20since 2021 · last 2026
0000-0002-5034-2880ORCID · corroborated
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 32 · 5 first-author · 16 since 2021Artificial intelligence and machine learning · 3 · 1 since 2021Computer networks · 3 · 1 since 2021Software engineering, systems software and programming languages · 3 · 1 first-author · 2 since 2021Applied, interdisciplinary, general and emerging computing · 2Databases, data management, data science and information retrieval · 1Human-computer interaction and ubiquitous computing · 1
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | Beyond Fixed Budgets: Characterizing the Inelasticity and Limitations of Tree-of-Thought Reasoning StrategiesabstractTree of Thought (ToT) search has become a promising direction for improving the reasoning capabilities of large language models, but deploying these methods in practice raises a question that has received little systematic attention: how do different search strategies behave under varying compute budgets, model sizes, and problem difficulties? In this work, we evaluate two representative ToT methods; DPTS, a Monte Carlo tree search based approach, and SSDP, a semantic deduplication based approach, across two mathematical reasoning benchmarks (Math500 and GSM8K), two model scales (Llama-3B and Llama-8B), and four token budgets (3k–10k). Our analysis reveals that the two methods exhibit limitations that pull in opposite directions. DPTS suffers from a cold-start bottleneck at low budgets: it requires sufficient exploration before its value estimates become reliable, making it a poor fit for resource-constrained settings despite strong scaling behavior at higher budgets. SSDP, on the other hand, reaches candidate solutions efficiently but is prone to frontier depletion; its aggressive node merging permanently discards unexplored paths, leaving it unable to improve regardless of how much budget remains. Together, these findings suggest that neither a fixed exploration strategy nor a fixed pruning strategy is sufficient across compute continuum. We argue that effective search for scientific reasoning agents requires strategies that can adapt their behavior based on search progress and available resources. Atkia Mahila, Avinash Maurya, M. Mustafa Rafique, Bogdan Nicolae |
HPDC | 3 |
| 2026 | Been There, Scanned That: Nostalgia-Driven Point Cloud Compression for Self-Driving CarsabstractAn autonomous vehicle generates several terabytes of sensor data per day. A significant portion of this data consists of 3D point clouds produced by depth sensors such as LiDAR. This data is transferred to cloud storage, where it is utilized for training machine learning models or conducting analyses, e.g., forensic investigations in the event of an accident. To reduce network and storage costs, this paper introduces DejaView that searches for and uses redundancies on larger temporal scales (days and months) for more effective compression. We designed DejaView with the insight that the operating area of autonomous vehicles is limited and that vehicles mostly traverse the same routes daily. Consequently, the daily collected 3D data is likely similar to the data they’ve captured in the past. To capture this, the core of DejaView is a diff operation that compactly represents point clouds as delta w.r.t. 3D data from the past. Using two months of LiDAR data, DejaView can compress point clouds by a factor of 210 at a reconstruction error of only 15 cm. Ali Khalid, Jaiaid Mobin, Sumanth Rao Appala, Avinash Maurya, Julie Stephany Berrio, M. Mustafa Rafique, Fawad Ahmad 0002 |
SenSys | 6 |
| 2026 | DataStates-LLM: Scalable Checkpointing for Transformer Models Using Composable State ProvidersabstractThe rapid growth of Large Transformer-based models, specifically Large Language Models (LLMs), now scaling to trillions of parameters, has necessitated training across thousands of GPUs using complex hybrid parallelism strategies (e.g., data, tensor, and pipeline parallelism). Checkpointing this massive, distributed state is critical for a wide range of use cases, such as resilience, suspend-resume, investigating undesirable training trajectories, and explaining model evolution. However, existing checkpointing solutions typically treat model state as opaque binary blobs, ignoring the “3D heterogeneity” of the underlying data structures–varying by memory location (GPU vs. Host), number of “logical” objects sharded and split across multiple files, data types (tensors vs. Python objects), and their serialization requirements. This results in significant runtime overheads due to blocking device-to-host transfers, data-oblivious serialization, and storage I/O contention. In this paper, we introduce DataStates-LLM, a novel check pointing architecture that leverages State Providers to decouple state abstraction from data movement. DataStates-LLM exploits the immutability of model parameters during the forward and backward passes to perform “lazy”, non-blocking asynchronous snapshots. By introducing State Providers, we efficiently coalesce fragmented, heterogeneous shards and overlap the serialization of metadata with bulk tensor I/O. We evaluate DataStates-LLM on models up to 70B parameters on 256 A100-40GB GPUs. Our results demonstrate that DataStates-LLM achieves up to 4× higher checkpointing throughput and reduces end-to-end training time by up to 2.2× compared to state-of-the-art solutions, effectively mitigating the serialization and heterogeneity bottlenecks in extreme-scale LLM training. Avinash Maurya, M. Mustafa Rafique, Franck Cappello, Bogdan Nicolae |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2025 | On Optimizing Checkpoint Restoration for HPC Applications: Leveraging Merkle Trees and Asynchronous I/OabstractEfficient checkpoint restoration is critical in high-performance computing (HPC) and AI applications, where slow recovery times disrupt workflows, waste resources, and hinder reproducibility. This work introduces a Merkle tree checkpoint restoration method to accelerate failure recovery and improve explainability. Our method integrates asynchronous I/O via the Liburing library to optimize scattered reads in HPC applications. Tested on the Polaris system at Argonne National Laboratory, it exhibits lower restoration time and memory consumption than state-of-the-art checkpoint restoration methods, reaching near-full efficiency with duplicated data. Our work advances scalable and efficient checkpointing solutions for HPC, ensuring reliable and fast failure recovery for large-scale simulations. Zackary Malkmus, Nigel Tan, Ian Lumsden, Kevin Assogba, M. Mustafa Rafique, Bogdan Nicolae, Michela Taufer |
HPDC | 5 |
| 2025 | Bayelemabaga: Creating Resources for Bambara NLPabstractAllahsera Auguste Tapo, Kevin Assogba, Christopher M Homan, M. Mustafa Rafique, Marcos Zampieri. Proceedings of the 2025 Conference of the Nations of the Americas Chapter of the Association for Computational Linguistics: Human Language Technologies (Volume 1: Long Papers). 2025. Allahsera Tapo, Kevin Assogba, Christopher Homan, M. Mustafa Rafique, Marcos Zampieri |
NAACL (Long Papers) | 4 |
| 2025 | MLP-Offload: Multi-Level, Multi-Path Offloading for LLM Pre-training to Break the GPU Memory WallabstractTraining LLMs larger than the aggregated memory of multiple GPUs is increasingly necessary due to the faster growth of LLM sizes compared to GPU memory. To this end, multi-tier host memory or disk offloading techniques are proposed by state of art. Despite advanced asynchronous multi-tier read/write strategies, such offloading strategies result in significant I/O overheads in the critical path of training, resulting in slower iterations. To this end, we propose MLP-Offload, a novel multi-level, multi-path offloading engine specifically designed for optimizing LLM training on resource-constrained setups by mitigating I/O bottlenecks. We make several key observations that drive the design of MLP-Offload, such as I/O overheads during the update dominate the iteration time; I/O bandwidth of the third-level remote storage tier remains unutilized; and, contention due to concurrent offloading amplifies I/O bottlenecks. Driven by these insights, we design and implement MLP-Offload to offload the optimizer states across multiple tiers in a cache-efficient and concurrency-controlled fashion to mitigate I/O bottlenecks during the backward and update phases. Evaluations on models up to 280B parameters shows that MLP-Offload achieves 2.5 × faster iterations compared to the state-of-the-art LLM training runtimes. Avinash Maurya, M. Mustafa Rafique, Franck Cappello, Bogdan Nicolae |
SC | 2 |
| 2024 | DataStates-LLM: Lazy Asynchronous Checkpointing for Large Language ModelsabstractLLMs have seen rapid adoption in all domains. They need to be trained on high-end high-performance computing (HPC) infrastructures and ingest massive amounts of input data. Unsurprisingly, at such a large scale, unexpected events (e.g., failures of components, instability of the software, undesirable learning patterns, etc.), are frequent and typically impact the training in a negative fashion. Thus, LLMs need to be checkpointed frequently so that they can be rolled back to a stable state and subsequently fine-tuned. However, given the large sizes of LLMs, a straightforward checkpointing solution that directly writes the model parameters and optimizer state to persistent storage (e.g., a parallel file system), incurs significant I/O overheads. To address this challenge, in this paper we study how to reduce the I/O overheads for enabling fast and scalable checkpointing for LLMs that can be applied at high frequency (up to the granularity of individual iterations) without significant impact on the training process. Specifically, we introduce a lazy asynchronous multi-level approach that takes advantage of the fact that the tensors making up the model and optimizer state shards remain immutable for extended periods of time, which makes it possible to copy their content in the background with minimal interference during the training process. We evaluate our approach at scales of up to 180 GPUs using different model sizes, parallelism settings, and checkpointing frequencies. The results show up to 48× faster checkpointing and 2.2× faster end-to-end training runtime compared with the state-of-art checkpointing approaches. Avinash Maurya, Robert Underwood, M. Mustafa Rafique, Franck Cappello, Bogdan Nicolae |
HPDC | 3 |
| 2024 | Application-Attuned Memory Management for Containerized HPC WorkflowsabstractHigh-Performance Computing (HPC) jobs consist of data and memory-intensive tasks often executed as workflows or ensembles to facilitate efficient and coordinated execution. These workflows are traditionally executed on HPC systems and have unique memory requirements based on the data size, computational complexity, and I/O activity. Recently containerized execution of these workflows has been extensively explored. Containerized workflow execution of HPC jobs requires several terabytes of memory that exceed node capacity, resulting in excessive data swapping to slower storage, degraded job performance, and failures. Similarly, colocated bandwidth-intensive, latency-sensitive, or short-lived workflows suffer from degraded performance due to contention, memory exhaustion, and higher access latency due to suboptimal memory allocation. Recently, tiered memory systems comprising persistent memory and compute express link (CXL) have been explored to provide additional memory capacity and bandwidth to memory-constrained systems and applications. However, current memory allocation and management techniques for tiered memory subsystems are inadequate to meet the diverse needs of colocated containerized jobs in HPC systems that concurrently run workflows and ensembles at scale. This paper leverages tiered memory systems for containerized HPC workflows and proposes efficient memory management policies including intelligent page placement and eviction policies to improve memory access performance. Our page allocation and replacement policies incorporate task characteristics and enable efficient memory sharing between workflows. We integrate our policies with the popular HPC scheduler, SLURM, and container runtime, Singularity, to show that our approach improves tiered memory utilization and application performance and reduces workflow execution times by up to 51%, 87%, and 35% as compared to the ideal, realistic, and optimized tiered execution environments, respectively. Moiz Arif, Avinash Maurya, M. Mustafa Rafique, Dimitrios S. Nikolopoulos, Ali Raza Butt |
IPDPS | 3 |
| 2024 | Deep Optimizer States: Towards Scalable Training of Transformer Models using Interleaved OffloadingabstractTransformers and large language models (LLMs) have seen rapid adoption in all domains. Their sizes have exploded to hundreds of billions of parameters and keep increasing. Under these circumstances, the training of transformers is very expensive and often hits a "memory wall", i.e., even when using 3D parallelism (pipeline, tensor, data) and aggregating the memory of many GPUs, it is still not enough to hold the necessary data structures (model parameters, optimizer state, gradients, activations) in GPU memory. To compensate, state-of-the-art approaches offload the optimizer state, at least partially, to the host memory and perform hybrid CPU-GPU computations. However, the management of the combined host-GPU memory is often suboptimal and results in poor overlapping between data movements and computations. This leads to missed opportunities to simultaneously leverage the interconnect bandwidth and computational capabilities of CPUs and GPUs. In this paper, we leverage a key observation that the interleaving of the forward, backward and update phases generate fluctuations in the GPU memory utilization, which can be exploited to dynamically move a part of the optimizer state between the host and the GPU memory at each iteration. To this end, we design and implement Deep Optimizer States, a novel technique to split the LLM into subgroups, whose update phase is scheduled on either the CPU or the GPU based on our proposed performance model that addresses the trade-off between data movement cost, acceleration on the GPUs vs the CPUs, and competition for shared resources. We integrate our approach with DeepSpeed and demonstrate 2.5× faster iterations over state-of-the-art approaches using extensive experiments. Avinash Maurya, M. Mustafa Rafique, Franck Cappello, Bogdan Nicolae |
Middleware | 3 |
| 2024 | Towards Affordable Reproducibility Using Scalable Capture and Comparison of Intermediate Multi-Run ResultsabstractEnsuring reproducibility in high-performance computing (HPC) applications is a significant challenge, particularly when nondeterministic execution can lead to untrustworthy results. Traditional methods that compare final results from multiple runs often fail because they provide sources of discrepancies only a posteriori and require substantial resources, making them impractical and unfeasible. This paper introduces an innovative method to address this issue by using scalable capture and comparing intermediate multi-run results. By capitalizing on intermediate checkpoints and hash-based techniques with user-defined error bounds, our method identifies divergences early in the execution paths. We employ Merkle trees for checkpoint data to reduce the I/O overhead associated with loading historical data. Our evaluations on the nondeterministic HACC cosmology simulation show that our method effectively captures differences above a predefined error bound and significantly reduces I/O overhead. Our solution provides a robust and scalable method for improving reproducibility, ensuring that scientific applications on HPC systems yield trustworthy and reliable results. Nigel Tan, Kevin Assogba, Walter J. Ashworth, Befikir Bogale, Franck Cappello, M. Mustafa Rafique, Michela Taufer, Bogdan Nicolae |
Middleware | 6 |
| 2023 | PredictDDL: Reusable Workload Performance Prediction for Distributed Deep LearningabstractAccurately predicting the training time of deep learning (DL) workloads is critical for optimizing the utilization of data centers and allocating the required cluster resources for completing critical model training tasks before a deadline. The state-of-the-art prediction models, e.g., Ernest and Cherrypick, treat DL workloads as black boxes, and require running the given DL job on a fraction of the dataset. Moreover, they require retraining their prediction models every time a change occurs in the given DL workload. This significantly limits the reusability of prediction models across DL workloads with different deep neural network (DNN) architectures. In this paper, we address this challenge and propose a novel approach where the prediction model is trained only once for a particular dataset type, e.g., ImageNet, thus completely avoiding tedious and costly retraining tasks for predicting the training time of new DL workloads. Our proposed approach, called PredictDDL, provides an end-to-end system for predicting the training time of DL models in distributed settings. PredictDDL leverages Graph HyperNetworks, a class of neural networks that takes computational graphs as input and produces vector representations of their DNNs. PredictDDL is the first prediction system that eliminates the need of retraining a performance prediction model for each new DL workload and maximizes the reuse of the prediction model by requiring running a DL workload only once for training the prediction model. Our extensive evaluation using representative workloads shows that PredictDDL achieves up to 9.8× lower average prediction error and 10.3× lower inference time compared to the state-of-the-art system, i.e., Ernest, on multiple DNN architectures. Kevin Assogba, Eduardo Lima, M. Mustafa Rafique, Minseok Kwon |
CLUSTER | 3 |
| 2023 | Optimizing the Training of Co-Located Deep Learning Models Using Cache-Aware StaggeringabstractDespite significant advances, training deep learning models remains a time-consuming and resource-intensive task. One of the key challenges in this context is the ingestion of the training data, which involves non-trivial overheads: read the training data from a remote repository, apply augmentations and transformations, shuffle the training samples, and assemble them into mini-batches. Despite the introduction of abstractions such as data pipelines that aim to hide such overheads asynchronously, it is often the case that the data ingestion is slower than the training, causing a delay at each training iteration. This problem is further augmented when training multiple deep learning models simultaneously on powerful compute nodes that feature multiple GPUs. In this case, the training data is often reused across different training instances (e.g., in the case of multi-model or ensemble training) or even within the same training instance (e.g., data-parallel training). However, transparent caching solutions (e.g., OS-level POSIX caching) are not suitable to directly mitigate the competition between training instances that reuse the same training data. In this paper, we study the problem of how to minimize the makespan of running two training instances that reuse the same training data. The makespan is subject to a trade-off: if the training instances start at the same time, competition for I/O bandwidth slows down the data pipelines and increases the makespan. If one training instance is staggered, competition is reduced but the makespan increases. We aim to optimize this trade-off by proposing a performance model capable of predicting the makespan based on the staggering between the training instances, which can be used to find the optimal staggering that triggers just enough competition to make optimal use of transparent caching in order to minimize the makespan. Experiments with different combinations of learning models using the same training data demonstrate that (1) staggering is important to minimize the makespan; (2) our performance model is accurate and can predict the optimal staggering in advance based on calibration overhead. Kevin Assogba, Bogdan Nicolae, M. Mustafa Rafique |
HiPC | 3 |
| 2023 | Towards Efficient I/O Pipelines Using Accumulated CompressionabstractHigh-Performance Computing (HPC) workloads generate large volumes of data at high-frequency during their execution, which needs to be captured concurrently at scale. These workloads exploit accelerators such as GPU for faster performance. However, the limited onboard high-bandwidth memory (HBM) on the GPU, and slow device-to-host memory PCIe interconnects lead to I/O overheads during application execution, thereby exacerbating their overall runtime. To overcome the aforementioned limitations, techniques such as compression and asynchronous transfers have been used by data management runtimes. However, compressing small blocks of data leads to a significant runtime penalty on the application. In this paper, we design and develop strategies to optimize the trade-off between compressing checkpoints instantly and enqueuing transfers immediately versus accumulating snapshots and delaying compression to achieve faster compression throughput. Our evaluations on synthetic and real-life workloads for different systems and workload configurations demonstrate 1.3 × to 8.3 × speedup compared to the existing checkpoint approaches. Avinash Maurya, Bogdan Nicolae, M. Mustafa Rafique, Franck Cappello |
HiPC | 3 |
| 2023 | GPU-Enabled Asynchronous Multi-level Checkpoint Caching and PrefetchingabstractCheckpointing is an I/O intensive operation increasingly used by High-Performance Computing (HPC) applications to revisit previous intermediate datasets at scale. Unlike the case of resilience, where only the last checkpoint is needed for application restart and rarely accessed to recover from failures, in this scenario, it is important to optimize frequent reads and writes of an entire history of checkpoints. State-of-the-art checkpointing approaches often rely on asynchronous multi-level techniques to hide I/O overheads by writing to fast local tiers (e.g. an SSD) and asynchronously flushing to slower, potentially remote tiers (e.g. a parallel file system) in the background, while the application keeps running. However, such approaches have two limitations. First, despite the fact that HPC infrastructures routinely rely on accelerators (e.g. GPUs), and therefore a majority of the checkpoints involve GPU memory, efficient asynchronous data movement between the GPU memory and host memory is lagging behind. Second, revisiting previous data often involves predictable access patterns, which are not exploited to accelerate read operations. In this paper, we address these limitations by proposing a scalable and asynchronous multi-level checkpointing approach optimized for both reading and writing of an arbitrarily long history of checkpoints. Our approach exploits GPU memory as a first-class citizen in the multi-level storage hierarchy to enable informed caching and prefetching of checkpoints by leveraging foreknowledge about the access order passed by the application as hints. Our evaluation using a variety of scenarios under I/O concurrency shows up to 74× faster checkpoint and restore throughput as compared to the state-of-art runtime and optimized unified virtual memory (UVM) based prefetching strategies and at least 2× shorter I/O wait time for the application across various workloads and configurations. Avinash Maurya, M. Mustafa Rafique, Thierry-Laurent D. Tonellot, Hussain J. AlSalem, Franck Cappello, Bogdan Nicolae |
HPDC | 2 |
| 2023 | COLTI: Towards Concurrent and Co-located DNN Training and InferenceabstractDeep learning models are extensively used in a wide range of domains, e.g., scientific simulations, predictions, and modeling. However, training these dense networks is both compute and memory intensive, and typically requires accelerators such as Graphics Processing Units (GPUs). While such DNN workloads consume a major proportion of the limited onboard high-bandwidth memory (HBM), they typically underutilize the GPU compute resources. In such scenarios, the idle compute resources on the GPU can be leveraged to run pending jobs that can either be (1) accommodated on the remainder HBM, or (2) can share memory resources with other concurrent workloads. However, state-of-the-art workload schedulers and DNN runtimes are not designed to leverage HBM co-location to improve resource utilization and throughput. In this work, we propose COLTI, which introduces a set of novel techniques to solve the aforementioned challenges by co-locating DNN training and inference on memory-constrained GPU devices. Our preliminary evaluations of three different DNN models implemented in the PyTorch framework demonstrate up to 37% and 40% improvement in makespan and memory utilization, respectively. Jaiaid Mobin, Avinash Maurya, M. Mustafa Rafique |
HPDC | 3 |
| 2022 | On Realizing Efficient Deep Learning Using Serverless ComputingabstractServerless computing is gaining rapid popularity as it enables quick application deployment and seamless application scaling without managing complex computing resources. Re-cently, it has been explored for running data-intensive, e.g., deep learning (DL), workloads for improving application performance and reducing execution cost. However, serverless computing imposes resource-level constraints, specifically fixed memory allocation and short task timeouts, that lead to job failures. In this paper, we address these constraints and develop an effective runtime framework, DiSDeL, that improves the performance of DL jobs by leveraging data splitting techniques, and ensuring that an appropriate amount of memory is allocated to containers for storing application data and a suitable timeout is selected for each job based on its complexity in serverless deployments. We implement our approach using Apache OpenWhisk and TensorFlow platforms and evaluate it using representative DL workloads to show that it eliminates DL job failures and reduces action memory consumption and total training time by up to 44% and 46%, respectively as compared to a default serverless computing framework. Our evaluation also shows that DiSDeL achieves a performance improvement of up to 29% as compared to bare-metal TensorFlow environment in a multi-tenant setting. Kevin Assogba, Moiz Arif, M. Mustafa Rafique, Dimitrios S. Nikolopoulos |
CCGRID | 3 |
| 2022 | Towards Efficient Cache Allocation for High-Frequency CheckpointingabstractWhile many HPC applications are known to have long runtimes, this is not always because of single large runs: in many cases, this is due to ensembles composed of many short runs (runtime in the order of minutes). When each such run needs to checkpoint frequently (e.g. adjoint computations using a checkpoint interval in the order of milliseconds), it is important to minimize both checkpointing overheads at each iteration, as well as initialization overheads. With the rising popularity of GPUs, minimizing both overheads simultaneously is challenging: while it is possible to take advantage of efficient asynchronous data transfers between GPU and host memory, this comes at the cost of high initialization overhead needed to allocate and pin host memory. In this paper, we contribute with an efficient technique to address this challenge. The key idea is to use an adaptive approach that delays the pinning of the host memory buffer holding the checkpoints until all memory pages are touched, which greatly reduces the overhead of registering the host memory with the CUDA driver. To this end, we use a combination of asynchronous touching of memory pages and direct writes of checkpoints to untouched and touched memory pages in order to minimize end-to-end checkpointing overheads based on performance modeling. Our evaluations show a significant improvement over a variety of alternative static allocation strategies and state-of-art approaches. Avinash Maurya, Bogdan Nicolae, M. Mustafa Rafique, Amr M. Elsayed, Thierry-Laurent D. Tonellot, Franck Cappello |
HIPC | 3 |
| 2022 | Exploiting CXL-based Memory for Distributed Deep LearningabstractDeep learning (DL) is being widely used to solve complex problems in scientific applications from diverse domains, such as weather forecasting, medical diagnostics, and fluid dynamics simulation. DL applications consume a large amount of data using large-scale high-performance computing (HPC) systems to train a given model. These workloads have large memory and storage requirements that typically go beyond the limited amount of main memory available on an HPC server. This significantly increases the overall training time as the input training data and model parameters are frequently swapped to slower storage tiers during the training process. In this paper, we use the latest advancements in the memory subsystem, specifically Compute Express Link (CXL), to provide additional memory and fast scratch space for DL workloads to reduce the overall training time while enabling DL jobs to efficiently train models using data that is much larger than the installed system memory. We propose a framework, called DeepMemoryDL, that manages the allocation of additional CXL-based memory, introduces a fast intermediate storage tier, and provides intelligent prefetching and caching mechanisms for DL workloads. We implement and integrate DeepMemoryDL with a popular DL platform, TensorFlow, to show that our approach reduces read and write latencies, improves the overall I/O throughput, and reduces the training time. Our evaluation shows a performance improvement of up to 34% and 27% compared to the default TensorFlow platform and CXL-based memory expansion approaches, respectively. Moiz Arif, Kevin Assogba, M. Mustafa Rafique, Sudharshan S. Vazhkudai |
ICPP | 3 |
| 2022 | Canary: Fault-Tolerant FaaS for Stateful Time-Sensitive ApplicationsabstractFunction-as-a-Service (FaaS) platforms have recently gained rapid popularity. Many stateful applications have been migrated to FaaS platforms due to their ease of deployment, scalability, and minimal management overhead. However, failures in FaaS have not been thoroughly investigated, thus making these desirable platforms unreliable for guaranteeing function execution and ensuring performance requirements. In this paper, we propose Canary, a highly resilient and fault-tolerant framework for FaaS that mitigates the impact of failures and reduces the overhead of function restart. Canary utilizes replicated container runtimes and application-level checkpoints to reduce application recovery time over FaaS platforms. Our evaluations using representative stateful FaaS applications show that Canary reduces the application recovery time and dollar cost by up to 83% and 12%, respectively over the default retry-based strategy. Moreover, it improves application availability with an additional average execution time and cost overhead of 14% and 8%, respectively, as compared to the ideal failure-free execution. Moiz Arif, Kevin Assogba, M. Mustafa Rafique |
SC | 3 |
| 2021 | Towards Efficient I/O Scheduling for Collaborative Multi-Level CheckpointingabstractEfficient checkpointing of distributed data structures periodically at key moments during runtime is a recurring fundamental pattern in a large number of uses cases: fault tolerance based on checkpoint-restart, in-situ or post-analytics, reproducibility, adjoint computations, etc. In this context, multilevel checkpointing is a popular technique: distributed processes can write their shard of the data independently to fast local storage tiers, then flush asynchronously to a shared, slower tier of higher capacity. However, given the limited capacity of fast tiers (e.g. GPU memory) and the increasing checkpoint frequency, the processes often run out of space and need to fall back to blocking writes to the slow tiers. To mitigate this problem, compression is often applied in order to reduce the checkpoint sizes. Unfortunately, this reduction is not uniform: some processes will have spare capacity left on the fast tiers, while others still run out of space. In this paper, we study the problem of how to leverage this imbalance in order to reduce I/O overheads for multi-level checkpointing. To this end, we solve an optimization problem of how much data to send from each process that runs out of space to the processes that have spare capacity in order to minimize the amount of time spent blocking in I/O. We propose two algorithms: one based on a greedy approach and the other based on modified minimum cost flows. We evaluate our proposal using synthetic and real-life application traces. Our evaluation shows that both algorithms achieve significant improvements in checkpoint performance over traditional multilevel checkpointing. Avinash Maurya, Bogdan Nicolae, M. Mustafa Rafique, Thierry-Laurent D. Tonellot, Franck Cappello |
MASCOTS | 3 |
| 2020 | MARBLE: A Multi-GPU Aware Job Scheduler for Deep Learning on HPC SystemsabstractDeep learning (DL) has become a key tool for solving complex scientific problems. However, managing the multi-dimensional large-scale data associated with DL, especially atop extant multiple graphics processing units (GPUs) in modern supercomputers poses significant challenges. Moreover, the latest high-performance computing (HPC) architectures bring different performance trends in training throughput compared to the existing studies. Existing DL optimizations such as larger batch size and GPU locality-aware scheduling have little effect on improving DL training throughput performance due to fast CPU-to-GPU connections. Additionally, DL training on multiple GPUs scales sublinearly. Thus, simply adding more GPUs to a system is ineffective. To this end, we design MARBLE, a first-of-its-kind job scheduler, which considers the non-linear scalability of GPUs at the intra-node level to schedule an appropriate number of GPUs per node for a job. By sharing the GPU resources on a node with multiple DL jobs, MARBLE avoids low GPU utilization in current multi-GPU DL training on HPC systems. Our comprehensive evaluation in the Summit supercomputer shows that MARBLE is able to improve DL training performance by up to 48.3% compared to the popular Platform Load Sharing Facility (LSF) scheduler. Compared to the state-of-the-art of DL scheduler, Optimus, MARBLE reduces the job completion time by up to 47%. Jingoo Han, M. Mustafa Rafique, Luna Xu, Ali Raza Butt, Seung-Hwan Lim, Sudharshan S. Vazhkudai |
CCGRID | 2 |
| 2020 | CuVPP: Filter-based Longest Prefix Matching in Software Data PlanesabstractProgrammability in the data plane has become increasingly important as virtualization is introduced into networking and software-defined networking becomes more prevalent. Yet, the performance of programmable data planes on commodity hardware is a major concern, in light of ever-increasing network speed and routing table size. This paper focuses on IP lookup, specifically the longest prefix matching for IPv6 addresses, which is a major performance bottleneck in programmable switches. As a solution, the paper presents CuVPP, a programmable switch that uses packet batch processing and cache locality for both instructions and data by leveraging Vector Packet Processing (VPP). We thoroughly evaluate CuVPP with both real network traffic and file-based lookup on a commodity hardware server connected via 80 Gbps network links and compare its performance with the other popular approaches. Our evaluation shows that CuVPP can achieve up to 4.5 million lookups per second with real traffic, higher than the other trie- or filter-based lookup approaches, and scales well even when the routing table size grows to 2 million prefixes. Minseok Kwon, Krishna Prasad Neupane, John Marshall, M. Mustafa Rafique |
CLUSTER | 4 |
| 2020 | CoSim: A Simulator for Co-Scheduling of Batch and On-Demand Jobs in HPC DatacentersabstractThe increasing scale and complexity of scientific applications are rapidly transforming the ecosystem of tools, methods, and workflows adopted by the high-performance computing (HPC) community. Big data analytics and deep learning are gaining traction as essential components in this ecosystem in a variety of scenarios, such as, steering of experimental instruments, acceleration of high-fidelity simulations through surrogate computations, and guided ensemble searches. In this context, the batch job model traditionally adopted by the supercomputing infrastructures needs to be complemented with support to schedule opportunistic on-demand analytics jobs, leading to the problem of efficient preemption of batch jobs with minimum loss of progress. In this paper, we design and implement a simulator, CoSim, that enables on-the-fly analysis of the trade-offs arising between delaying the start of opportunistic on-demand jobs, which leads to longer analytics latency, and loss of progress due to preemption of batch jobs, which is necessary to make room for on-demand jobs. To this end, we propose an algorithm based on dynamic programming with predictable performance and scalability that enables supercomputing infrastructure schedulers to analyze the aforementioned trade-off and take decisions in near real-time. Compared with other state-of-art approaches using traces of the Theta pre-Exascale machine, our approach is capable of finding the optimal solution, while achieving high performance and scalability. Avinash Maurya, Bogdan Nicolae, Ishan Guliani, M. Mustafa Rafique |
DS-RT | 4 |
| 2020 | Infrastructure-Aware TensorFlow for Heterogeneous DatacentersabstractHeterogeneous datacenters, with a variety of compute, memory, and network resources, are becoming increasingly popular to address the resource requirements of time-sensitive applications. One such application framework is the TensorFlow platform, which has become a platform of choice for running machine learning workloads. The state-of-the-art TensorFlow platform is oblivious to the availability and performance profiles of the underlying datacenter resources and does not incorporate resource requirements of the given workloads for distributed training. This leads to executing the training tasks on busy and resource-constrained worker nodes, which results in a significant increase in the overall training time. In this paper, we address this challenge and propose architectural improvements and new software modules in the default TensorFlow platform to make it aware of the availability and capabilities of the underlying datacenter resources. The proposed Infrastructure-Aware Tensor-Flow efficiently schedules the training tasks on the best possible resources for execution and reduces the overall training time. Our evaluation using the worker nodes with varying availability and performance profiles shows that the proposed enhancements yield up to 54 % reduced training time as compared to the default TensorFlow platform. Moiz Arif, M. Mustafa Rafique, Seung-Hwan Lim, Zaki Malik |
MASCOTS | 2 |
| 2019 | A Quantitative Study of Deep Learning Training on Heterogeneous SupercomputersabstractThe following topics are dealt with: parallel processing; storage management; learning (artificial intelligence); resource allocation; scheduling; computer centres; multiprocessing systems; graphics processing units; neural nets; message passing. Jingoo Han, Luna Xu, M. Mustafa Rafique, Ali Raza Butt, Seung-Hwan Lim |
CLUSTER | 3 |
| 2019 | Container Orchestration by Kubernetes for RDMA NetworkingabstractWith the widespread usage of containerized virtualization in data centers and clouds, it is important to enabling high-throughput and zero-copy data transfer between those containers. Remote Direct Memory Access (RDMA) allows bypassing the kernel for packet processing by offloading it to specific RDMA-enabled NICs. The existing solutions enabling RDMA with containers are either based on custom container orchestrators (e.g., FreeFlow) or lack the ability for the control plane to manage the underlying RDMA traffic (e.g., Kubernetes RDMA plug-in via SR-IOV). The work in this paper builds off of previous work in Kubernetes to make an architecture that allows control over bandwidth requirements of RDMA within a Kubernetes cluster. Coleman Link, Jesse Sarran, Garegin Grigoryan, Minseok Kwon, M. Mustafa Rafique, Warren R. Carithers |
ICNP | 5 |
| 2017 | Evaluation of Data Locality Strategies for Hybrid Cloud Bursting of Iterative MapReduceabstractHybrid cloud bursting (i.e., leasing temporary off-premise cloud resources to boost the overall capacity during peak utilization) is a popular and cost-effective way to deal with the increasing complexity of big data analytics. It is particularly promising for iterative MapReduce applications that reuse massive amounts of input data at each iteration, which compensates for the high overhead and cost of concurrent data transfers from the on-premise to the off-premise VMs over a weak inter-site link that is of limited capacity. In this paper we study how to combine various MapReduce data locality techniques designed for hybrid cloud bursting in order to achieve scalability for iterative MapReduce applications in a cost-effective fashion. This is a non trivial problem due to the complex interaction between the data movements over the weak link and the scheduling of computational tasks that have to adapt to the shifting data distribution. We show that using the right combination of techniques, iterative MapReduce applications can scale well in a hybrid cloud bursting scenario and come even close to the scalability observed in single sites. Francisco J. Clemente-Castelló, Bogdan Nicolae, M. Mustafa Rafique, Rafael Mayo 0002, Juan Carlos Fernández 0002 |
CCGrid | 3 |
| 2017 | Optimization of data-intensive workflows in stream-based data processing models
Saima Gulzar Ahmad, Chee Sun Liew, M. Mustafa Rafique, Ehsan Ullah Munir |
J. Supercomput. | 3 |
| 2016 | On exploiting data locality for iterative mapreduce applications in hybrid cloudsabstractHybrid cloud bursting (i.e., leasing temporary off-premise cloud resources to boost the capacity during peak utilization), has made significant impact especially for big data analytics, where the explosion of data sizes and increasingly complex computations frequently leads to insufficient local data center capacity. Cloud bursting however introduces a major challenge to runtime systems due to the limited throughput and high latency of data transfers between on-premise and off-premise resources (weak link). This issue and how to address it is not well understood. We contribute with a comprehensive study on what challenges arise in this context, what potential strategies can be applied to address them and what best practices can be leveraged in real-life. Specifically, we focus our study on iterative MapReduce applications, which are a class of large-scale data intensive applications particularly popular on hybrid clouds. In this context, we study how data locality can be leveraged over the weak link both from the storage layer perspective (when and how to move it off-premise) and from the scheduling perspective (when to compute off-premise). We conclude with a brief discussion on how to set up an experimental framework suitable to study the effectiveness of our proposal in future work. Francisco J. Clemente-Castelló, Bogdan Nicolae, Rafael Mayo 0002, Juan Carlos Fernández 0002, M. Mustafa Rafique |
BDCAT | 5 |
| 2016 | On Efficient Hierarchical Storage for Big Data ProcessingabstractA promising trend in storage management for big data frameworks, such as Hadoop and Spark, is the emergence of heterogeneous and hybrid storage systems that employ different types of storage devices, e.g. SSDs, RAMDisks, etc., alongside traditional HDDs. However, scheduling data accesses or requests to an appropriate storage device is non-trivial and depends on several factors such as data locality, device performance, and application compute and storage resources utilization. To this end, we present DUX, an application-attuned dynamic data management system for data processing frameworks, which aims to improve overall application I/O throughput by efficiently using SSDs only for workloads that are expected to benefit from them rather than the extant approach of storing a fraction of the overall workloads in SSDs. The novelty of DUX lies in profiling application performance on SSDs and HDDs, analyzing the resulting I/O behavior, and considering the available SSDs at runtime to dynamically place data in an appropriate storage tier. Evaluation of DUX with trace-driven simulations using synthetic Facebook workloads shows that even when using 5.5× fewer SSDs compared to a SSD-only solution, DUX incurs only a small (5%) performance overhead, and thus offers an affordable and efficient storage tier management. Krish K. R., Bharti Wadhwa, M. Safdar Iqbal, M. Mustafa Rafique, Ali Raza Butt |
CCGrid | 4 |
| 2016 | CHOPPER: Optimizing Data Partitioning for In-memory Data Analytics FrameworksabstractThe performance of in-memory based data analytic frameworks such as Spark is significantly affected by how data is partitioned. This is because the partitioning effectively determines task granularity and parallelism. Moreover, different phases of a workload execution can have different optimal partitions. However, in the current implementations, the tuning knobs controlling the partitioning are either configured statically or involve a cumbersome programmatic process for affecting changes at runtime. In this paper, we propose CHOPPER, a system for automatically determining the optimal number of partitions for each phase of a workload and dynamically changing the partition scheme during workload execution. CHOPPER monitors the task execution and DAG scheduling information to determine the optimal level of parallelism. CHOPPER repartitions data as needed to ensure efficient task granularity, avoids data skew, and reduces shuffle traffic. Thus, CHOPPER allows users to write applications without having to hand-tune for optimal parallelism. Experimental results show that CHOPPER effectively improves workload performance by up to 35.2% compared to standard Spark setup. Arnab Kumar Paul, Wenjie Zhuang, Luna Xu, M. Mustafa Rafique, Ali Raza Butt |
CLUSTER | 5 |
| 2015 | Heterogeneous cloud systems monitoring using semantic and linked data technologiesabstractCloud businesses need comprehensive visibility on hardware and software components, their utilization and their configuration. In addition, they need to integrate such information with their asset management systems and publicly available information such as hardware specifications. In this paper, we present an approach for cloud management and monitoring based on a semantic layer that unifies different interfaces and representations, and makes all relevant information accessible from a single point. We show a proof-of-concept based on OpenStack and Linked Data technologies and evaluate it in terms of overhead, query execution times, and effectiveness in the data gathering phase. Our findings indicate that a semantics-based approach is indeed feasible and advantageous for providing uniform access across different cloud environments and levels. Alessandro Portosa, M. Mustafa Rafique, Spyros Kotoulas, Luca Foschini 0001, Antonio Corradi |
IM | 2 |
| 2014 | A Capacity Allocation Approach for Volunteer Cloud Federations Using Poisson-Gamma Gibbs SamplingabstractIn volunteer cloud federations (VCFs), volunteers join and leave without restrictions and may collectively contribute a large number of heterogeneous virtual machine instances. A challenge is to efficiently allocate this dynamic, heterogeneous capacity to a flow of incoming virtual machine (VM) instantiation requests, i.e., maximize the number of virtual machines that may be placed on the VCF. Cloud federations may allocate VMs far more efficiently if they can accurately predict the demand in terms of VM instantiation requests. In this paper, we present a stochastic technique that forecasts future demand to efficiently allocate VMs to VM instantiation requests. Our approach uses a Markov Chain Monte Carlo (MCMC) simulation known as the Poisson-Gamma Gibbs (PGG) sampler. The PGG sampler is used to determine the arrival rate of each type of VM instantiation requests. This arrival rate is then used to determine an optimal VM placement for the incoming VM instantiation requests. We compared our approach to a solution that adopts a static smallest-fit approach. The experimental results showed that our solution reacts quickly to abrupt changes in the frequency of VM instantiation requests and performs 10% better than the static smallest-fit approach in terms of the total number of satisfied requests. Abdelmounaam Rezgui, Gary Quezada, M. Mustafa Rafique, Zaki Malik |
IEEE CLOUD | 3 |
| 2014 | Distributed Detection of Cancer Cells in High-Throughput Cellular Spike StreamsabstractDetection and identification of important biological targets such as, DNA, proteins, and diseased human cells is crucial towards early disease diagnosis and prognosis. The key to differentiate healthy cells from the diseased cells is the biophysical properties that differ significantly. Micro and nanosystems, such as solid-state micropores and nanopores, can measure and translate these properties of human cells and DNA into electrical spikes to decode useful biological insights. Nonetheless, such approaches result in large data streams that are often plagued with inherit noise and baseline wanders. Moreover, the extant detection approaches are tedious, time-consuming, and error-prone, and there is no error-resilient software that can analyze large datasets instantly. The ability to effectively process and detect biological targets in larger datasets lies in the automated and accelerated data processing strategies using state-of-the-art distributed computing systems. To this end, we propose a distributed detection framework, which collects the raw data stream on a server node that then splits/distributes the data into segments across the worker nodes. Each node reduces noise in the assigned data segment using moving-average filtering, and detects the electric spikes by comparing them against a statistical threshold (based on the mean and standard deviation of the data), in a Single Program Multiple Data (SPMD) style. Our proposed framework enables the detection of cancer cells with an accuracy of 63% in a mixture of Cancer cells, Red Blood Cells (RBCs), and White Blood Cells (WBCs), and achieves a maximum speedup of 6X over a single-node machine by processing 10 gigabytes of raw data using an 8-node cluster in less than a minute. Abdul Hafeez, M. Mustafa Rafique, Ali Raza Butt |
CCGRID | 2 |
| 2013 | Leveraging Collaborative Content Exchange for On-Demand VM Multi-deployments in IaaS Clouds
Bogdan Nicolae, M. Mustafa Rafique |
Euro-Par | 2 |
| 2012 | On the Use of GPUs in Realizing Cost-Effective Distributed RAIDabstractThe exponential growth in user and application data entails new means for providing fault tolerance and protection against data loss. High Performance Computing (HPC) storage systems, which are at the forefront of handling the data deluge, typically employ hardware RAID at the backend. However, such solutions are costly, do not ensure end-to-end data integrity, and can become a bottleneck during data reconstruction. In this paper, we design an innovative solution to achieve a flexible, fault-tolerant, and high-performance RAID-6 solution for a parallel file system (PFS). Our system utilizes low-cost, strategically placed GPUs - both on the client and server sides - to accelerate parity computation. In contrast to hardware-based approaches, we provide full control over the size, length and location of a RAID array on a per file basis, end-to-end data integrity checking, and parallelization of RAID array reconstruction. We have deployed our system in conjunction with the widely-used Lustre PFS, and show that our approach is feasible and imposes acceptable overhead. Aleksandr Khasymski, M. Mustafa Rafique, Ali Raza Butt, Sudharshan S. Vazhkudai, Dimitrios S. Nikolopoulos |
MASCOTS | 2 |
| 2011 | Symphony: A Scheduler for Client-Server Applications on Coprocessor-Based Heterogeneous ClustersabstractCoprocessors such as GPUs are increasingly being deployed in clusters to process scientific and compute-intensive jobs. In this work, we study if GPU-based heterogeneous clusters can benefit client-server applications. Specifically, we consider the practical situation where multiple client-server applications share a heterogeneous cluster (multi-tenancy), and experience unpredictable variations in incoming client request rates, including steep load spikes. Even for "compute-intensive" client-server applications, it is unclear if a GPU-based cluster can seamlessly deliver acceptable response times in the presence of multi-tenancy and load spikes. We argue that a cluster-level scheduler that is aware of application load, request deadlines and the heterogeneity is necessary in this situation. We propose a novel scheduler called Symphony that enables efficient, dynamic sharing of a GPU-based heterogeneous cluster across multiple concurrently-executing client-server applications, each with arbitrary load spikes. Symphony performs three key tasks: it (i) monitors the load on each application, (ii) collects past performance data and dynamically builds simple performance models of available processing resources and (iii) computes a priority for pending requests based on the above parameters and the requests' slack. Based on this, it reorders client requests across different applications to achieve acceptable response times. We also define how client-server applications should interact with a scheduler such as Symphony, and develop an API to this end. We deploy Symphony as user-space middleware on a high-end heterogeneous cluster with dual quad-core Xeon CPUs and dual NVIDIA Fermi GPUs. An evaluation using representative applications shows that in the presence of load spikes (i) Symphony incurs 2-20× fewer requests that do not meet response time constraints compared with other schedulers, and (ii) in order to achieve the same performance as Symphony, other schedulers need 2× more cluster nodes. M. Mustafa Rafique, Srihari Cadambi, Kunal Rao, Ali Raza Butt, Srimat T. Chakradhar |
CLUSTER | 1 |
| 2011 | A capabilities-aware framework for using computational accelerators in data-intensive computing
M. Mustafa Rafique, Ali Raza Butt, Dimitrios S. Nikolopoulos |
J. Parallel Distributed Comput. | 1 |
| 2011 | Reusable software components for accelerator-based clusters
M. Mustafa Rafique, Ali Raza Butt, Eli Tilevich |
J. Syst. Softw. | 1 |
| 2010 | A Capabilities-Aware Programming Model for Asymmetric High-End SystemsabstractIn this research, we investigate and address the challenges of asymmetry in High-End Computing (HEC) systems comprising heterogeneous architectures with varying I/O and computation capacities. We focus on developing a flexible, scalable and easy-to-use programming model that automatically adapts to the capabilities of the system resources on largescale asymmetric clusters. Furthermore, we aim to develop innovative and efficient workload distribution techniques that bridge the asymmetry between system components. In particular, we intent to design tools and technologies that enable quick and efficient utilization of high-end asymmetric clusters in large-scale settings for modern scientific and enterprise computing. M. Mustafa Rafique |
CCGRID | 1 |
| 2010 | Designing Accelerator-Based Distributed Systems for High PerformanceabstractMulti-core processors with accelerators are becoming commodity components for high-performance computing at scale. While accelerator-based processors have been studied in some detail, the design and management of clusters based on these processors have not received the same focus. In this paper, we present an exploration of four design and resource management alternatives, which can be used on large-scale asymmetric clusters with accelerators. Moreover, we adapt the popular MapReduce programming model to our proposed configurations. We enhance MapReduce with new dynamic data streaming and workload scheduling capabilities, which enable application writers to use asymmetric accelerator-based clusters without being concerned with the capabilities of individual components. We present an evaluation of the presented designs in a physical setting and show that our designs can provide significant performance advantages. Compared to a standard static MapReduce design, we achieve 62.5%, 73.1%, and 82.2% performance improvement using accelerators with limited general-purpose resources, well-provisioned shared general-purpose resources, and well-provisioned dedicated general-purpose resources, respectively. M. Mustafa Rafique, Ali Raza Butt, Dimitrios S. Nikolopoulos |
CCGRID | 1 |
| 2009 | CellMR: A framework for supporting mapreduce on asymmetric cell-based clustersabstractThe use of asymmetric multi-core processors with on-chip computational accelerators is becoming common in a variety of environments ranging from scientific computing to enterprise applications. The focus of current research has been on making efficient use of individual systems, and porting applications to asymmetric processors. In this paper, we take the next step by investigating the use of multi-core-based systems, especially the popular Cell processor, in a cluster setting. We present CellMR, an efficient and scalable implementation of the MapReduce framework for asymmetric Cell-based clusters. The novelty of CellMR lies in its adoption of a streaming approach to supporting MapReduce, and its adaptive resource scheduling schemes: Instead of allocating workloads to the components once, CellMR slices the input into small work units and streams them to the asymmetric nodes for efficient processing. Moreover, CellMR removes I/O bottlenecks by design, using a number of techniques, such as double-buffering and asynchronous I/O, to maximize cluster performance. Our evaluation of CellMR using typical MapReduce applications shows that it achieves 50.5% better performance compared to the standard nonstreaming approach, introduces a very small overhead on the manager irrespective of application input size, scales almost linearly with increasing number of compute nodes (a speedup of 6.9 on average, when using eight nodes compared to a single node), and adapts effectively the parameters of its resource management policy between applications with varying computation density. M. Mustafa Rafique, Benjamin Rose, Ali Raza Butt, Dimitrios S. Nikolopoulos |
IPDPS | 1 |