VLDB 2026 Research / reviewers in the wild / expert
Arnab Kumar Paul
dblp:191/3093 · also Arnab K. Paul
· DBLP profile ↗
24ranked-venue papers
8as first author
17since 2021 · last 2026
0000-0002-3694-5511ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 21 · 7 first-author · 16 since 2021Databases, data management, data science and information retrieval · 2 · 1 first-author · 1 since 2021Applied, interdisciplinary, general and emerging computing · 2 · 1 first-authorArtificial intelligence and machine learning · 1 · 1 first-authorSoftware engineering, systems software and programming languages · 1 · 1 since 2021
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | CoDL: A Framework for Studying Cross-Component Interference in Deep Learning Training PipelinesabstractDeep learning training on HPC systems exhibits significant step-time variability—data stalls can dominate 50–85% of training time [12, 31]—that limits efficiency. Existing benchmarks and profilers isolate individual components: I/O benchmarks replace compute with sleep, compute profilers ignore I/O, and GPU tracers operate on separate timescales. This isolation prevents understanding real variability, which stems from interactions among data loading, compute, and communication competing for shared CPU resources. We present CoDL, a framework that combines component-isolation benchmarking with multi-level tracing on a unified microsecond timeline, enabling diagnosis that reframes optimization from isolated component throughput to CPU–GPU scheduling coordination. On a single 4-GPU node, evaluation of ResNet-50 shows compute duration more than doubles (+127.4%) under active data loading and NCCL AllReduce degrades 2.8 ×, despite zero direct GPU-kernel overlap with data loading; the dominant mechanism is data-loading workers preempting GPU-launch threads on the CPU, inflating CUDA runtime API latency by 3.3 ×. Druva Dhakshinamoorthy, Ray A. O. Sinurat, Nikoli Dryden, Arnab Kumar Paul, Hariharan Devarajan |
HPDC | 4 |
| 2025 | UnifyFL: Enabling Decentralized Cross-Silo Federated Learning
Sarang S, Druva Dhakshinamoorthy, Aditya Shiva Sharma, Yuvraj Singh Bhadauria, Siddharth Chaitra Vivek, Arihant Bansal, Arnab Kumar Paul |
Middleware | 7 |
| 2024 | FedCaSe: Enhancing Federated Learning with Heterogeneity-aware Caching and SchedulingabstractFederated learning (FL) has emerged as a new paradigm of machine learning (ML) with the goal of collaborative learning on the vast pool of private data available across distributed edge devices. The focus of most existing works in FL systems has been on addressing the challenges of computation and communication heterogeneity inherent in training with edge devices. However, the crucial impact of I/O and the role of limited on-device storage has not been explored fully in FL context. Without policies to exploit the on-device storage for placement of client data samples, and schedule clients based on I/O benefits, FL training can lead to inefficiencies, such as increased training time and impacted accuracy convergence. Redwan Ibne Seraj Khan, Arnab Kumar Paul, Yue Cheng 0001, Xun Jian 0002, Ali Raza Butt |
SoCC | 2 |
| 2024 | When Less is More: Achieving Faster Convergence in Distributed Edge Machine LearningabstractDistributed Machine Learning (DML) on resource-constrained edge devices holds immense potential for real-world applications. However, achieving fast convergence in DML in these heterogeneous environments remains a significant challenge. Traditional frameworks like Bulk Synchronous Parallel (BSP) and Asynchronous Stochastic Parallel (ASP) rely on frequent, small updates that incur substantial communication overhead and hinder convergence speed. Furthermore, these frameworks often employ static dataset sizes, neglecting the heterogeneity of edge devices and potentially leading to straggler nodes that slow down the entire training process. The straggler nodes, i.e., edge devices that take significantly longer to process their assigned data chunk - hinder the overall training speed. To address these limitations, this paper proposes Hermes, a novel probabilistic framework for efficient DML on edge devices. This framework leverages a dynamic threshold based on recent test loss behavior to identify statistically significant improvements in the model's generalization capability, hence transmitting updates only when major improvements are detected, thereby significantly reducing communication overhead. Additionally, Hermes employs dynamic dataset allocation to optimize resource utilization and prevents performance degradation caused by straggler nodes. Our evaluations on a real-world heterogeneous resource-constrained environment demonstrate that Hermes achieves faster convergence compared to state-of-the-art methods, resulting in a remarkable 13.22x reduction in training time and a 62.1% decrease in communication overhead. Advik Raj Basani, Siddharth Chaitra Vivek, Advaith Krishna, Arnab Kumar Paul |
HiPC | 4 |
| 2024 | Tarazu: An Adaptive End-to-end I/O Load-balancing Framework for Large-scale Parallel File SystemsabstractThe imbalanced I/O load on large parallel file systems affects the parallel I/O performance of high-performance computing (HPC) applications. One of the main reasons for I/O imbalances is the lack of a global view of system-wide resource consumption. While approaches to address the problem already exist, the diversity of HPC workloads combined with different file striping patterns prevents widespread adoption of these approaches. In addition, load-balancing techniques should be transparent to client applications. To address these issues, we propose Tarazu , an end-to-end control plane where clients transparently and adaptively write to a set of selected I/O servers to achieve balanced data placement. Our control plane leverages real-time load statistics for global data placement on distributed storage servers, while our design model employs trace-based optimization techniques to minimize latency for I/O load requests between clients and servers and to handle multiple striping patterns in files. We evaluate our proposed system on an experimental cluster for two common use cases: the synthetic I/O benchmark IOR and the scientific application I/O kernel HACC-I/O. We also use a discrete-time simulator with real HPC application traces from emerging workloads running on the Summit supercomputer to validate the effectiveness and scalability of Tarazu in large-scale storage environments. The results show improvements in load balancing and read performance of up to 33% and 43%, respectively, compared to the state-of-the-art. Arnab Kumar Paul, Sarah Neuwirth, Bharti Wadhwa, Feiyi Wang, Sarp Oral, Ali Raza Butt |
ACM Trans. Storage | 1 |
| 2024 | An End-to-end High-performance Deduplication Scheme for Docker Registries and Docker Container Storage SystemsabstractThe wide adoption of Docker containers for supporting agile and elastic enterprise applications has led to a broad proliferation of container images. The associated storage performance and capacity requirements place a high pressure on the infrastructure of container registries that store and distribute images and container storage systems on the Docker client side that manage image layers and store ephemeral data generated at container runtime. The storage demand is worsened by the large amount of duplicate data in images. Moreover, container storage systems that use Copy-on-Write (CoW) file systems as storage drivers exacerbate the redundancy. Exploiting the high file redundancy in real-world images is a promising approach to drastically reduce the growing storage requirements of container registries and improve the space efficiency of container storage systems. However, existing deduplication techniques significantly degrade the performance of both registries and container storage systems because of data reconstruction overhead as well as the deduplication cost. We propose DupHunter, an end-to-end deduplication scheme that deduplicates layers for both Docker registries and container storage systems while maintaining a high image distribution speed and container I/O performance. DupHunter is divided into three tiers: registry tier, middle tier, and client tier. Specifically, we first build a high-performance deduplication engine at the registry tier that not only natively deduplicates layers for space savings but also reduces layer restore overhead. Then, we use deduplication offloading at the middle tier to eliminate the redundant files from the client tier and avoid bringing deduplication overhead to the clients. To further reduce the data duplicates caused by CoWs and improve the container I/O performance, we utilize a container-aware storage system at the client tier that reserves space for each container and arranges the placement of files and their modifications on the disk to preserve locality. Under real workloads, DupHunter reduces storage space by up to 6.9× and reduces the GET layer latency up to 2.8× compared to the state-of-the-art. Moreover, DupHunter can improve the container I/O performance by up to 93% for reads and 64% for writes. Muhui Lin, Hadeel Albahar, Arnab Kumar Paul, Zhijie Huan, Subil Abraham, Vasily Tarasov, Dimitrios Skourtis, Ali Anwar 0001, Ali Raza Butt |
ACM Trans. Storage | 4 |
| 2023 | SHADE: Enable Fundamental Cacheability for Distributed Deep Learning Training
Redwan Ibne Seraj Khan, Ahmad Hossein Yazdani, Yuqi Fu, Arnab Kumar Paul, Bo Ji 0001, Xun Jian 0002, Yue Cheng 0001, Ali Raza Butt |
FAST | 4 |
| 2023 | Modeling the Impact of System-Level Parameters on I/O Performance of HPC ApplicationsabstractModern High Performance Computing (HPC) workloads prioritize data, challenging storage systems to meet rising I/O demands. The network's role in inter-node communication and data transfer significantly impacts overall application performance. Our study systematically examines bandwidth and latency variations in CephFS due to diverse I/O patterns. We formalize latency, bandwidth, I/O size, and block-to-file size ratio relationships and train a model for predicting bandwidth and estimating latency in similar I/O workloads, offering valuable insights into optimizing HPC storage. Debasmita Biswas, Arnab Kumar Paul, Sarah Neuwirth, Ali Raza Butt |
MASCOTS | 2 |
| 2023 | Analyzing File Access Patterns on Large-Scale HPC Systems: Opportunities for File PrefetchingabstractThis paper explores the potential opportunities for implementing file prefetching techniques on large-scale high-performance computing (HPC) systems. Specifically, we investigate the file access patterns of various applications across multiple scientific domains using two years' worth of Darshan I/O traces obtained from the Summit supercomputer. We identify recurring trends and patterns which indicate that prefetching can be effectively leveraged to improve data access performance on HPC systems. This study serves as a valuable reference for system architects and developers in the HPC community, providing insights into the opportunities and challenges associated with enabling file prefetching on large-scale HPC systems. Ahmad Maroof Karimi, Arnab Kumar Paul, Jong Choi 0001, Lipeng Wan 0001, Feiyi Wang |
MASCOTS | 2 |
| 2022 | SchedTune: A Heterogeneity-Aware GPU Scheduler for Deep LearningabstractModern cluster management systems, such as Kubernetes, support heterogeneous workloads and resources. However, existing resource schedulers in these systems do not differentiate between heterogeneous G PU resources-which are becoming a norm-and do not support GPU sharing-which is necessary to support emerging collocation of jobs and multi-tenant applications. Thus the systems suffer from low GPU resource utilization, higher queuing delays, and an increase in application makespan, i.e., the duration between the arrival of the first job and the completion of the last job of a workflow. This is especially a problem in supporting crucial deep learning (DL) applications. To this end, in this paper, we profile and analyze DL jobs on heterogeneous GPUs, investigate the interference caused by collocating jobs on GPUs, and use this information to predict the GPU memory demand and job completion times. We propose SCHEDTUNE, a machine-learning-based heterogeneity-aware scheduler that ensures higher GPU memory utilization and reduced out-of-memory (OOM) failures, while supporting improved makespan. Our evaluation shows that SCHEDTUNE GPU memory predictors and scheduler outperform the state-of-the-art predictors by achieving 81% higher GPU memory utilization, 100% detection and avoidance of OOM errors, and 17.5% reduction in makespan compared to the default Kubernetes scheduler. Hadeel Albahar, Shruti Dongare, Yanlin Du, Arnab Kumar Paul, Ali Raza Butt |
CCGRID | 5 |
| 2022 | Hvac: Removing I/O Bottleneck for Large-Scale Deep Learning ApplicationsabstractScientific communities are increasingly adopting deep learning (DL) models in their applications to accelerate scientific discovery processes. However, with rapid growth in the computing capabilities of HPC supercomputers, large-scale DL applications have to spend a significant portion of training time performing I/O to a parallel storage system. Previous research works have investigated optimization techniques such as prefetching and caching. Unfortunately, there exist non-trivial challenges to adopting the existing solutions on HPC supercomputers for large-scale DL training applications, which include non-performance and/or failures at extreme scale, lack of portability and generality in design, complex deployment methodology, and being limited to a specific application or dataset. To address these challenges, we propose High-Velocity AI Cache (HVAC), a distributed read-cache layer that targets and fully exploits the node-local storage or near node-local storage technology. HVAC seamlessly accelerates read I/O by aggregating node-local or near node-local storage, avoiding metadata lookups and file locking while preserving portability in the application code. We deploy and evaluate HVAC on 1,024 nodes (with over 6000 NVIDIA V100 GPUS) of the Summit supercomputer. In particular, we evaluate the scalability, efficiency, accuracy, and load distribution of HVAC compared to GPFS and XFS-on-NVMe. With four different DL applications, we observe an average 25 % performance improvement atop GPFS and 9% drop against XFS-on-NVMe, which scale linearly and are considered the performance upper bound. We envision HVAC as an important caching library for upcoming HPC supercomputers such as Frontier. Awais Khan 0002, Arnab Kumar Paul, Christopher Zimmer 0001, Sarp Oral, Sajal Dash, Scott Atchley, Feiyi Wang |
CLUSTER | 2 |
| 2022 | Access Patterns and Performance Behaviors of Multi-layer Supercomputer I/O Subsystems under Production LoadabstractScientific computing workloads at HPC facilities have been shifting from traditional numerical simulations to AI/ML applications for training and inference while processing and producing ever-increasing amounts of scientific data. To address the growing need for increased storage capacity, lower access latency, and higher bandwidth, emerging technologies such as non-volatile memory are integrated into supercomputer I/O subsystems. With these emerging trends, we need a better understanding of the multilayer supercomputer I/O systems and ways to use these subsystems efficiently. In this work, we study the I/O access patterns and performance characteristics of two representative supercomputer I/O subsystems. Through an extensive analysis of year-long I/O logs on each system, we report new observations in I/O reads and writes, unbalanced use of storage system layers, and new trends in user behaviors at the HPC I/O middleware stack. Jean Luca Bez, Ahmad Maroof Karimi, Arnab Kumar Paul, Surendra Byna, Philip H. Carns, Sarp Oral, Feiyi Wang, Jesse Hanley |
HPDC | 3 |
| 2022 | Machine Learning Assisted HPC Workload Trace Generation for Leadership Scale Storage SystemsabstractMonitoring and analyzing a wide range of I/O activities in an HPC cluster is important in maintaining mission-critical performance in a large-scale, multi-user, parallel storage system. Center-wide I/O traces can provide high-level information and fine-grained activities per application or per user running in the system. Studying such large-scale traces can provide helpful insights into the system. It can be used to develop predictive methods for making predictive decisions, adjusting scheduling policies, or providing decisions for the design of next-generation systems. However, sharing real-world I/O traces to expedite such research efforts leaves a few concerns; i) the cost of sharing the large traces is expensive due to this large size, and ii) privacy concern is an issue. Arnab Kumar Paul, Jong Choi 0001, Ahmad Maroof Karimi, Feiyi Wang |
HPDC | 1 |
| 2022 | I/O performance analysis of machine learning workloads on leadership scale supercomputer
Ahmad Maroof Karimi, Arnab Kumar Paul, Feiyi Wang |
Perform. Evaluation | 2 |
| 2021 | Parallel I/O Evaluation Techniques and Emerging HPC Workloads: A PerspectiveabstractEmerging workloads such as artificial intelligence, big data analytics and complex multi-step workflows alongside future exascale applications are anticipated future HPC workloads, which will result in a more diverse I/O system workload and even less predictable I/O behavior and access patterns. Along with the ever increasing gap between the compute and storage performance capabilities, the in-depth understanding of extreme-scale I/O behavior and the I/O performance modeling and prediction are essential tools of the large-scale I/O evaluation process for addressing the needs of extreme-scale hybrid workloads. In this survey article, we focus on the state-of-the-art of the I/O behavior and performance analysis process for HPC systems in a 5-year time window and identify future research challenges. Sarah Neuwirth, Arnab Kumar Paul |
CLUSTER | 2 |
| 2021 | Characterizing Machine Learning I/O Workloads on Leadership Scale HPC SystemsabstractHigh performance computing (HPC) is no longer solely limited to traditional workloads such as simulation and modeling. With the increase in the popularity of machine learning (ML) and deep learning (DL) technologies, we are observing that an increasing number of HPC users are incorporating ML methods into their workflow and scientific discovery processes, across a wide spectrum of science domains such as biology, earth science, and physics. This gives rise to a diverse set of I/O patterns than the traditional checkpoint/restart-based HPC I/O behavior. The details of the I/O characteristics of such ML I/O workloads have not been studied extensively for large-scale leadership HPC systems. This paper aims to fill that gap by providing an in-depth analysis to gain an understanding of the I/O behavior of ML I/O workloads using darshan - an I/O characterization tool designed for lightweight tracing and profiling. We study the darshan logs of more than 23, 000 HPC ML I/O jobs over a time period of one year running on Summit - the second-fastest supercomputer in the world. This paper provides a systematic I/O characterization of ML I/O jobs running on a leadership scale supercomputer to understand how the I/O behavior differs across science domains and the scale of workloads, and analyze the usage of parallel file system and burst buffer by ML I/O workloads. Arnab Kumar Paul, Ahmad Maroof Karimi, Feiyi Wang |
MASCOTS | 1 |
| 2021 | Large-Scale Analysis of Docker Images and Performance Implications for Container Storage SystemsabstractDocker containers have become a prominent solution for supporting modern enterprise applications due to the highly desirable features of isolation, low overhead, and efficient packaging of the application’s execution environment. Containers are created from images which are shared between users via a registry. The amount of data registries store is massive. For example, Docker Hub, a popular public registry, stores at least half a million public images. In this article, we analyze over 167 TB of uncompressed Docker Hub images, characterize them using multiple metrics and evaluate the potential of file-level deduplication. Our analysis helps to make conscious decisions when designing storage for containers in general and Docker registries in particular. For example, only 3 percent of the files in images are unique while others are redundant file copies, which means file-level deduplication has a great potential to save storage space. Furthermore, we carry out a comprehensive analysis of both small I/O request performance and copy-on-write performance for multiple popular container storage drivers. Our findings can motivate and help improve the design of data reduction and caching methods for images, pulling optimizations for registries, and storage drivers. Vasily Tarasov, Hadeel Albahar, Ali Anwar 0001, Lukas Rupprecht, Dimitrios Skourtis, Arnab Kumar Paul, Ali Raza Butt |
IEEE Trans. Parallel Distributed Syst. | 7 |
| 2020 | On the Use of Containers in High Performance Computing EnvironmentsabstractThe lightweight nature, application portability, and deployment flexibility of containers is driving their widespread adoption in cloud solutions. Data analysis and deep learning (DL)/machine learning (ML) applications have especially benefited from containerization. As such data analysis is adopted in high performance computing (HPC), the need for container support in HPC has become paramount. However, containers face crucial performance and I/O challenges in HPC. One obstacle is that while there have been HPC containers, such solutions have not been thoroughly investigated, especially from the aspect of their impact on the crucial HPC I/O throughput. To this end, this paper provides a first-of-its-kind empirical analysis of state-of-the-art representative container solutions (Docker, Podman, Singularity, and Charliecloud) in HPC environments. We also explore how containers interact with an HPC parallel file system like Lustre. We present the design of an analysis framework that is deployed on all nodes in an HPC environment, and captures CPU, memory, network, and file I/O statistics from the nodes and the storage system. We are able to garner key insights from our analysis, e.g., Charliecloud outperforms other container solutions in terms of container start-up time, while Singularity and Charliecloud are equivalent in I/O throughput. But this comes at a cost, as Charliecloud invokes the most metadata and I/O operations on the underlying Lustre file system. By identifying such trade-offs and optimization opportunities, we can enhance HPC containers performance and the ML/DL applications that increasingly rely on them. Subil Abraham, Arnab Kumar Paul, Redwan Ibne Seraj Khan, Ali Raza Butt |
CLOUD | 2 |
| 2020 | Efficient Metadata Indexing for HPC Storage SystemsabstractThe increase in data generation rate along with the scale of today's high performance computing (HPC) storage systems make finding and managing files extremely difficult. Efficient file system metadata indexing and querying tools are needed to ease file system management. Current metadata indexing techniques either use spatial trees or an external database to index metadata. Both approaches have their drawbacks which reduce the performance of indexing and querying the metadata on large scale file systems. In this paper, we have developed Brindexer, a metadata indexing and search tool specifically designed for large-scale HPC storage systems. Brindexer is mainly designed for system administrators to help them manage the file system effectively. It uses a leveled partitioning approach to partition the file system namespace, and has an in-tree design to reduce resource utilization from an external database. Brindexer uses RDBMS for efficient querying of the metadata index database, also uses a changelog-based approach to effectively handle real-time metadata changes and re-index the metadata at regular intervals. We implement and evaluate Brindexer on a 4.8 TB Lustre store and show that it improves the indexing and querying performance by 69% and 91% when compared to state-of-the-art metadata indexing tools. Arnab Kumar Paul, Nathan Rutman, Cory Spitz, Ali Raza Butt |
CCGRID | 1 |
| 2020 | Understanding HPC Application I/O Behavior Using System Level StatisticsabstractThe processor performance of high performance computing (HPC) systems is increasing at a much higher rate than storage performance. This imbalance leads to I/O performance bottlenecks in massively parallel HPC applications. Therefore, there is a need for improvements in storage and file system designs to meet the ever-growing I/O needs of HPC applications. Storage and file system designers require a deep understanding of how HPC application I/O behavior affects current storage system installations in order to improve them. In this work, we contribute to this understanding using application-agnostic file system statistics gathered on compute nodes as well as metadata and object storage file system servers. We analyze file system statistics of more than 4 million jobs over a period of three years on two systems at Lawrence Livermore National Laboratory that include a 15 PiB Lustre file system for storage. The results of our study add to the state-of-the-art in I/O understanding by providing insight into how general HPC workloads affect the performance of large-scale storage systems. Some key observations in our study show that reads and writes are evenly distributed across the storage system; applications which perform I/O, spread that I/O across ~78% of the minutes of their runtime on average; less than 22% of HPC users who submit write-intensive jobs perform efficient writes to the file system; and I/O contention seriously impacts I/O performance. Arnab Kumar Paul, Olaf Faaland, Adam Moody, Elsa Gonsiorowski, Kathryn Mohror, Ali Raza Butt |
HiPC | 1 |
| 2019 | FSMonitor: Scalable File System Monitoring for Arbitrary Storage SystemsabstractData automation, monitoring, and management tools are reliant on being able to detect, report, and respond to file system events. Various data event reporting tools exist for specific operating systems and storage devices, such as inotify for Linux, kqueue for BSD, and FSEvents for macOS. However, these tools are not designed to monitor distributed file systems. Indeed, many cannot scale to monitor many thousands of directories, or simply cannot be applied to distributed file systems. Moreover, each tool implements a custom API and event representation, making the development of generalized and portable event-based applications challenging. As file systems grow in size and become increasingly diverse, there is a need for scalable monitoring solutions that can be applied to a wide range of both distributed and local systems. We present here a generic and scalable file system monitor and event reporting tool, FSMonitor, that provides a file-system-independent event representation and event capture interface. FSMonitor uses a modular Data Storage Interface (DSI) architecture to enable the selection and application of appropriate event monitoring tools to detect and report events from a target file system, and implements efficient and fault-tolerant mechanisms that can detect and report events even on large file systems. We describe and evaluate DSIs for common UNIX, macOS, and Windows storage systems, and for the Lustre distributed file system. Our experiments on a 897 TB Lustre file system show that FSMonitor can capture and process almost 38 000 events per second. Arnab Kumar Paul, Ryan Chard, Kyle Chard, Steven Tuecke, Ali Raza Butt, Ian T. Foster |
CLUSTER | 1 |
| 2019 | iez: Resource Contention Aware Load Balancing for Large-Scale Parallel File SystemsabstractParallel I/O performance is crucial to sustaining scientific applications on large-scale High-Performance Computing (HPC) systems. However, I/O load imbalance in the underlying distributed and shared storage systems can significantly reduce overall application performance. There are two conflicting challenges to mitigate this load imbalance: (i) optimizing systemwide data placement to maximize the bandwidth advantages of distributed storage servers, i.e., allocating I/O resources efficiently across applications and job runs; and (ii) optimizing client-centric data movement to minimize I/O load request latency between clients and servers, i.e., allocating I/O resources efficiently in service to a single application and job run. Moreover, existing approaches that require application changes limit wide-spread adoption in commercial or proprietary deployments. We propose iez, an “end-to-end control plane” where clients transparently and adaptively write to a set of selected I/O servers to achieve balanced data placement. Our control plane leverages realtime load information for distributed storage server global data placement while our design model leverages trace-based optimization techniques to minimize I/O load request latency between clients and servers. We evaluate our proposed system on an experimental cluster for two common use cases: synthetic I/O benchmark IOR for large sequential writes and a scientific application I/O kernel, HACC-I/O. Results show read and write performance improvements of up to 34% and 32%, respectively, compared to the state of the art. Bharti Wadhwa, Arnab Kumar Paul, Sarah Neuwirth, Feiyi Wang, Sarp Oral, Ali Raza Butt, Jon Bernard, Kirk W. Cameron |
IPDPS | 2 |
| 2017 | I/O load balancing for big data HPC applicationsabstractHigh Performance Computing (HPC) big data problems require efficient distributed storage systems. However, at scale, such storage systems often experience load imbalance and resource contention due to two factors: the bursty nature of scientific application I/O; and the complex I/O path that is without centralized arbitration and control. For example, the extant Lustre parallel file system-that supports many HPC centers-comprises numerous components connected via custom network topologies, and serves varying demands of a large number of users and applications. Consequently, some storage servers can be more loaded than others, which creates bottlenecks and reduces overall application I/O performance. Existing solutions typically focus on per application load balancing, and thus are not as effective given their lack of a global view of the system. In this paper, we propose a data-driven approach to load balance the I/O servers at scale, targeted at Lustre deployments. To this end, we design a global mapper on Lustre Metadata Server, which gathers runtime statistics from key storage components on the I/O path, and applies Markov chain modeling and a minimum-cost maximum-flow algorithm to decide where data should be placed. Evaluation using a realistic system simulator and a real setup shows that our approach yields better load balancing, which in turn can improve end-to-end performance. Arnab Kumar Paul, Arpit Goyal, Feiyi Wang, Sarp Oral, Ali Raza Butt, Michael J. Brim, Sangeetha B. Srinivasa |
IEEE BigData | 1 |
| 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 | 1 |