VLDB 2026 Research / reviewers in the wild / expert
Randal C. Burns
dblp:b/RandalCBurns
· DBLP profile ↗
76ranked-venue papers
8as first author
12since 2021 · last 2026
0000-0002-2924-1997ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 44 · 5 first-author · 6 since 2021Databases, data management, data science and information retrieval · 25 · 3 first-author · 3 since 2021Artificial intelligence and machine learning · 6 · 2 since 2021Applied, interdisciplinary, general and emerging computing · 5 · 2 since 2021Software engineering, systems software and programming languages · 3 · 1 since 2021Computer networks · 2 · 1 since 2021Security and privacy · 2Graphics, computer vision, multimedia, augmented reality and games · 2Theory of computation · 1 · 1 since 2021
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | PTStore (Prefix Tensor Store): Distributed Prefix Caching and Replication for High Throughput Inference Serving
Meghana Maghyastha, Robert Underwood, Randal C. Burns, Bogdan Nicolae |
Euro-Par (2) | 3 |
| 2025 | Poster: On Harnessing Idle Compute at the Edge for Foundation Model TrainingabstractFoundation model training is increasingly centralized in large cloud data centers because it demands immense compute and memory resources. Training over decentralized edge devices could democratize this ecosystem by harnessing otherwise idle compute, but prior edge-training systems fall short: they scale poorly with model size and device count, exceed per-device memory budgets, incur prohibitive collective communication, and are fragile to heterogeneous and dynamic device availability. We present Cleave, a parameter-server-centric framework that makes tensor-parallel training practical at the edge. Cleave introduces selective hybrid tensor parallelism, which finely shards GEMM-dominated training operations into memory-feasible sub-tasks while avoiding peer-to-peer collectives that become bottlenecks on asymmetric edge links. A cost model guides device selection and shard placement to mitigate stragglers and rapidly adapt to churn. Across OPT and Llama2 models, Cleave matches cloud GPU training efficiency while scaling to thousands of devices. It supports up to 8× more devices than prior edge approaches, reduces per-batch training time by up to 10×, and achieves 100× faster recovery from device failures. Leyang Xue, Meghana Madhyastha, Myungjin Lee, Amos J. Storkey, Randal C. Burns, Mahesh K. Marina |
MobiCom | 5 |
| 2024 | EvoStore: Towards Scalable Storage of Evolving Learning ModelsabstractDeep Learning (DL) has seen rapid adoption in all domains. Since training DL models is expensive, both in terms of time and resources, application workflows that make use of DL increasingly need to operate with a large number of derived learning models, which are obtained through transfer learning and fine-tuning. At scale, thousands of such derived DL models are accessed concurrently by a large number of processes. In this context, an important question is how to design and develop specialized DL model repositories that remain scalable under concurrent access, while addressing key challenges: how to query the DL model architectures for specific patterns? How to load/store a subset of layers/tensors from a DL model? How to efficiently share unmodified layers/tensors between DL models derived from each other through transfer learning? How to maintain provenance and answer ancestry queries? State of art leaves a gap regarding these challenges. To fill this gap, we introduce EvoStore, a distributed DL model repository with scalable data and metadata support to store and access derived DL models efficiently. Large-scale experiments on hundreds of GPUs show significant benefits over state-of-art with respect to I/O and metadata performance, as well as storage space utilization. Robert Underwood, Meghana Madhyastha, Randal C. Burns, Bogdan Nicolae |
HPDC | 3 |
| 2024 | T-Rex (Tree-Rectangles): Reformulating Decision Tree Traversal as Hyperrectangle EnclosureabstractTree ensembles, random forests and gradient boosted trees, are useful in resource-limited machine learning deployments. However, traversing tree data structures is not cache friendly, which results in high latency during inference or regression. Tree traversal incurs random I/Os making inference memory bound. We present a system that trades many random I/Os for few sequential I/O by remapping a forest of trees into a single spatial index. It builds on the observation that each leaf in the forest encodes a hyperrectangle in the feature space. We make queries I/O efficient through pruning and space-filling curves. We then optimize computation through quantization of hyperrectangle boundaries and vectorization of enclosure queries. Our evaluation on a diverse set of benchmark datasets shows that the system reduces inference latency by 2 times in memory and 10 times for external memory with no detectable loss of accuracy. Meghana Madhyastha, Tamás Budavári, Vladimir Braverman, Joshua T. Vogelstein, Randal C. Burns |
ICDE | 5 |
| 2024 | CPMA: An Efficient Batch-Parallel Compressed Set Without PointersabstractThis paper introduces the batch-parallel Compressed Packed Memory Array (CPMA), a compressed, dynamic, ordered set data structure based on the Packed Memory Array (PMA). Traditionally, batch-parallel sets are built on pointer-based data structures such as trees because pointer-based structures enable fast parallel unions via pointer manipulation. When compared with cache-optimized trees, PMAs were slower to update but faster to scan. Brian Wheatman, Randal C. Burns, Aydin Buluç, Helen Xu 0001 |
PPoPP | 2 |
| 2023 | Optimizing Search Layouts in Packed Memory ArraysabstractThis paper introduces Search-optimized Packed Memory Arrays (SPMAs), a collection of data structures based on Packed Memory Arrays (PMAs) that address suboptimal search via cache-optimized search layouts. Traditionally, PMAs and B-trees have tradeoffs between searches/inserts and scans: B-trees were faster for searches and inserts, while PMAs were faster for scans. Our empirical evaluation shows that SPMAs overcome this tradeoff for unsorted input distributions: on average, SPMAs are faster than B+-trees (a variant of B-trees optimized for scans) on all major operations. We generated datasets and search/insert workloads from the Yahoo! Cloud Serving Benchmark (YCSB) and found that SPMAs are about 2× faster than B+-trees regardless of the ratio of searches to inserts. On uniform random inputs, SPMAs are on average between 1.3× −2.3× faster than B+-trees on all operations. Finally, we vary the amount of sortedness in the inputs to stress the worst-case insert distribution in the PMA. We find that the worst-case B+-tree insertion throughput is about 1.5× faster than the worst-case PMA insertion throughput. However, the worst-case input for the PMA is sorted and highly unlikely to appear naturally in practice. The SPMAs maintain higher insertion throughput than the B+-tree when the input is up to 25% sorted. Brian Wheatman, Randal C. Burns, Aydin Buluç, Helen Xu 0001 |
ALENEX | 2 |
| 2023 | Understanding Patterns of Deep Learning Model Evolution in Network Architecture SearchabstractNetwork Architecture Search and specifically Regularized Evolution is a common way to refine the structure of a deep learning model. However, little is known about how models empirically evolve over time which has design implications for designing caching policies, refining the search algorithm for particular applications, and other important use cases. In this work, we algorithmically analyze and quantitatively characterize the patterns of model evolution for a set of models from the Candle project and the Nasbench-201 search space. We show how the evolution of the model structure is influenced by the regularized evolution algorithm. We describe how evolutionary patterns appear in distributed settings and opportunities for caching and improved scheduling. Lastly, we describe the conditions that affect when particular model architectures rise and fall in popularity based on their frequency of acting as a donor in a sliding window. Robert Underwood, Meghana Madhyastha, Randal C. Burns, Bogdan Nicolae |
HiPC | 3 |
| 2023 | DStore: A Lightweight Scalable Learning Model Repository with Fine-Grain Tensor-Level AccessabstractThe ability to share and reuse deep learning (DL) models is a key driver that facilitates the rapid adoption of artificial intelligence (AI) in both industrial and scientific applications. However, state-of-the-art approaches to store and access DL models efficiently at scale lag behind. Most often, DL models are serialized by using various formats (e.g., HDF5, SavedModel) and stored as files on POSIX file systems. While simple and portable, such an approach exhibits high serialization and I/O overheads, especially under concurrency. Additionally, the emergence of advanced AI techniques (transfer learning, sensitivity analysis, explainability, etc.) introduces the need for fine-grained access to tensors to facilitate the extraction and reuse of individual or subsets of tensors. Such patterns are underserved by state-of-the-art approaches. Requiring tensors to be read in bulk incurs suboptimal performance, scales poorly, and/or overutilizes network bandwidth. In this paper we propose a lightweight, distributed, RDMA-enabled learning model repository that addresses these challenges. Specifically we introduce several ideas: compact architecture graph representation with stable hashing and client-side metadata caching, scalable load balancing on multiple providers, RDMA-optimized data staging, and direct access to raw tensor data. We evaluate our proposal in extensive experiments that involve different access patterns using learning models of diverse shapes and sizes. Our evaluations show a significant improvement (between 2 and 30× over a variety of state-of-the-art model storage approaches while scaling to half the Cooley cluster at the Argonne Leadership Computing Facility. Meghana Madhyastha, Robert Underwood, Randal C. Burns, Bogdan Nicolae |
ICS | 3 |
| 2022 | Towards Optimal Line of Sight CoverageabstractMaintaining the line of sight to a moving object or person over long distances is critical in many applications, e.g., mobile communications, security, surveillance. Determining the best places to position (or build) technologies is difficult because even small changes in the location can greatly affect the so-called viewshed, which is the collection of land areas within line of sight of a given observer. The need for multiple sensors or towers further complicates this problem, as they often need to work cooperatively to achieve the best possible coverage. This study proposes a novel approach that consists of two separate inventions: 1) Introduction of a meaningful measure of quality for coverage to compare competing configurations; and 2) Optimization of that well-defined objective function to find the best suitable sensor parameters for practical applications. Preliminary results suggest unprecedented performance on a wide range of real terrains. Peter Gu, Tamás Budavári, Amanda Galante, Randal C. Burns |
e-Science | 4 |
| 2021 | Streaming Sparse Graphs using Efficient Dynamic SetsabstractWe present the SSTGraph framework for the storage and analysis of dynamic graphs. Its performance matches or exceeds state-of-the-art static graph engines and supports streaming updates. SSTGraph builds on top of the tinyset parallel, dynamic set data structure. Tinyset implements set membership in a shallow hierarchy of sorted packed memory arrays to achieve logarithmic time access and updates, and it scans in optimal linear time. Tinyset uses space comparable to that of systems that use data compression while avoiding compression’s computation and serialization overhead.SSTGraph outperforms other streaming, dynamic graph engines on a suite of four graph algorithms. Our evaluation includes a comparison with the Aspen streaming graph system. SSTGraph reduces runtime by 40% on average, updates are 2x-5x faster on batch sizes up to 10 million, and graphs are smaller. The partitioned data structure scales well and runs on billion edge graphs in just 15 GB of memory. Brian Wheatman, Randal C. Burns |
IEEE BigData | 2 |
| 2021 | Understanding and dealing with hard faults in persistent memory systemsabstractThe advent of Persistent Memory (PM) devices enables systems to actively persist information at low costs, including program state traditionally in volatile memory. However, this trend poses a reliability challenge in which multiple classes of soft faults that go away after restart in traditional systems turn into hard (recurring) faults in PM systems. In this paper, we first characterize this rising problem with an empirical study of 28 real-world bugs. We analyze how they cause hard faults in PM systems. We then propose Arthas, a tool to effectively recover PM systems from hard faults. Arthas checkpoints PM states via fine-grained versioning and uses program slicing of fault instructions to revert problematic PM states to good versions. We evaluate Arthas on 12 real-world hard faults from five large PM systems. Arthas successfully recovers the systems for all cases while discarding 10× less data on average compared to state-of-the-art checkpoint-rollback solutions. Randal C. Burns |
EuroSys | 2 |
| 2021 | BLOCKSET (Block-Aligned Serialized Trees): Reducing Inference Latency for Tree ensemble DeploymentabstractWe present methods to serialize and deserialize gradient-boosted trees and random forests that optimize inference latency when models are not loaded into memory. This arises when models are larger than memory, but also systematically when models are deployed on low-resource devices in the Internet of Things or run as cloud microservices where resources are allocated on demand. Block-Aligned Serialized Trees (BLOCKSET) introduce the concept of selective access for random forests and gradient boosted trees in which only the parts of the model needed for inference are deserialized and loaded into memory. %BLOCKSET combines concepts from external memory algorithms and data-parallel %layouts of random forests that maximize I/O-density for in-memory models. Using principles from external memory algorithms, we block-align the serialization format in order to minimize the number of I/Os. For gradient boosted trees, this results in a more than five time reduction in inference latency over layouts that do not perform selective access and a 2 times latency reduction over techniques that are selective, but do not encode I/O block boundaries in the layout. Meghana Madhyastha, Kunal Lillaney, James Browne, Joshua T. Vogelstein, Randal C. Burns |
KDD | 5 |
| 2020 | Geodesic ForestsabstractTogether with the curse of dimensionality, nonlinear dependencies in large data sets persist as major challenges in data mining tasks. A reliable way to accurately preserve nonlinear structure is to compute geodesic distances between data points. Manifold learning methods, such as Isomap, aim to preserve geodesic distances in a Riemannian manifold. However, as manifold learning algorithms operate on the ambient dimensionality of the data, the essential step of geodesic distance computation is sensitive to high-dimensional noise. Therefore, a direct application of these algorithms to high-dimensional, noisy data often yields unsatisfactory results and does not accurately capture nonlinear structure. Meghana Madhyastha, Gongkai Li, Veronika Strnadová-Neeley, James Browne, Joshua T. Vogelstein, Randal C. Burns, Carey E. Priebe |
KDD | 6 |
| 2020 | Sparse Projection Oblique Randomer ForestsabstractDecision forests, including Random Forests and Gradient Boosting Trees, have recently demonstrated state-of-the-art performance in a variety of machine learning settings. Decision forests are typically ensembles of axis-aligned decision trees; that is, trees that split only along feature dimensions. In contrast, many recent extensions to decision forests are based on axis-oblique splits. Unfortunately, these extensions forfeit one or more of the favorable properties of decision forests based on axis-aligned splits, such as robustness to many noise dimensions, interpretability, or computational efficiency. We introduce yet another decision forest, called “Sparse Projection Oblique Randomer Forests” (SPORF). SPORF trees recursively split along very sparse random projections. Our method significantly improves accuracy over existing state-of-the-art algorithms on a standard benchmark suite for classification with $>100$ problems of varying dimension, sample size, and number of classes. To illustrate how SPORF addresses the limitations of both axis-aligned and existing oblique decision forest methods, we conduct extensive simulated experiments. SPORF typically yields improved performance over existing decision forest methods, while mitigating computational efficiency and scalability and maintaining interpretability. Very sparse random projections can be incorporated into gradient boosted trees to obtain potentially similar gains. Tyler M. Tomita, James Browne, Cencheng Shen, Jaewon Chung, Jesse Patsolic, Benjamin Falk, Carey E. Priebe, Jason Yim, Randal C. Burns, Mauro Maggioni, Joshua T. Vogelstein |
J. Mach. Learn. Res. | 9 |
| 2019 | Agni: An Efficient Dual-access File System over Object StorageabstractObject storage is a low-cost, scalable component of cloud ecosystems. However, interface incompatibilities and performance limitations inhibit its adoption for emerging cloud-based workloads. Users are compelled to either run their applications over expensive block storage-based file systems or use inefficient file connectors over object stores. Dual access, the ability to read and write the same data through file systems interfaces and object storage APIs, has promise to improve performance and eliminate storage sprawl. Kunal Lillaney, Vasily Tarasov, David Pease, Randal C. Burns |
SoCC | 4 |
| 2019 | The Case for Dual-access File Systems over Object Storage
Kunal Lillaney, Vasily Tarasov, David Pease, Randal C. Burns |
HotStorage | 4 |
| 2019 | Forest Packing: Fast Parallel, Decision ForestsabstractDecision Forests are popular machine learning techniques that assist scientists to extract knowledge from massive data sets. This class of tool remains popular because of their interpretability and ease of use, unlike other modern machine learning methods, such as kernel machines and deep learning. Decision forests also scale well for use with large data because training and run time operations are trivially parallelizable allowing for high inference throughputs. A negative aspect of these forests, and an untenable property for many real time applications, is their high inference latency caused by the combination of large model sizes with random memory access patterns. We present memory packing techniques and a novel tree traversal method to overcome this deficiency. The result of our system is a grouping of trees into a hierarchical structure. At low levels, we pack the nodes of multiple trees into contiguous memory blocks so that each memory access fetches data for multiple trees. At higher levels, we use leaf cardinality to identify the most popular paths through a tree and collocate those paths in contiguous cache lines. We extend this layout with a re-ordering of the tree traversal algorithm to take advantage of the increased memory throughput provided by out-of-order execution and cache-line prefetching. Together, these optimizations increase the performance and parallel scalability of classification in ensembles by a factor of ten over an optimized C++ implementation and a popular R-language implementation. James Browne, Disa Mhembere, Tyler M. Tomita, Joshua T. Vogelstein, Randal C. Burns |
SDM | 5 |
| 2018 | Building NDStore Through Hierarchical Storage Management and Microservice ProcessingabstractWe describe NDStore, a scalable multi-hierarchical data storage deployment for spatial analysis of neuroscience data on the AWS cloud. The system design is inspired by the requirement to maintain high I/O throughput for workloads that build neural connectivity maps of the brain from peta-scale imaging data using computer vision algorithms. We store all our data on the AWS object store S3 to limit our deployment costs. S3 serves as our base-tier of storage. Redis, an in-memory key-value engine, is used as our caching tier. The data is dynamically moved between the different storage tiers based on user access. All programming interfaces to this system are RESTful web-services. We include a performance evaluation that shows that our production system provides good performance for a variety of workloads by combining the assets of multiple cloud services. Kunal Lillaney, Dean Kleissas, Alexander Eusman, Eric A. Perlman, William R. Gray Roncal, Joshua T. Vogelstein, Randal C. Burns |
eScience | 7 |
| 2018 | FlashR: parallelize and scale R for machine learning using SSDsabstractR is one of the most popular programming languages for statistics and machine learning, but it is slow and unable to scale to large datasets. The general approach for having an efficient algorithm in R is to implement it in C or FORTRAN and provide an R wrapper. FlashR accelerates and scales existing R code by parallelizing a large number of matrix functions in the R base package and scaling them beyond memory capacity with solid-state drives (SSDs). FlashR performs memory hierarchy aware execution to speed up parallelized R code by (i) evaluating matrix operations lazily, (ii) performing all operations in a DAG in a single execution and with only one pass over data to increase the ratio of computation to I/O, (iii) performing two levels of matrix partitioning and reordering computation on matrix partitions to reduce data movement in the memory hierarchy. We evaluate FlashR on various machine learning and statistics algorithms on inputs of up to four billion data points. Despite the huge performance gap between SSDs and RAM, FlashR on SSDs closely tracks the performance of FlashR in memory for many algorithms. The R implementations in FlashR outperforms H2O and Spark MLlib by a factor of 3 -- 20. Da Zheng 0004, Disa Mhembere, Joshua T. Vogelstein, Carey E. Priebe, Randal C. Burns |
PPoPP | 5 |
| 2018 | Remote visual analysis of large turbulence databases at multiple scales
Jesus Pulido, Daniel Livescu, Kalin Kanov, Randal C. Burns, Curtis Canada, James P. Ahrens, Bernd Hamann |
J. Parallel Distributed Comput. | 4 |
| 2017 | knor: A NUMA-Optimized In-Memory, Distributed and Semi-External-Memory k-means Libraryabstractk-means is one of the most influential and utilized machine learning algorithms. Its computation limits the performance and scalability of many statistical analysis and machine learning tasks. We rethink and optimize k-means in terms of modern NUMA architectures to develop a novel parallelization scheme that delays and minimizes synchronization barriers. The k-means NUMA Optimized Routine knor) library has (i) in-memory knori), (ii) distributed memory (knord), and (ii) semi-external memory (\textsf{knors}) modules that radically improve the performance of k-means for varying memory and hardware budgets. knori boosts performance for single machine datasets by an order of magnitude or more. \textsf{knors} improves the scalability of k-means on a memory budget using SSDs. knors scales to billions of points on a single machine, using a fraction of the resources that distributed in-memory systems require. knord retains knori's performance characteristics, while scaling in-memory through distributed computation in the cloud. knor modifies Elkan's triangle inequality pruning algorithm such that we utilize it on billion-point datasets without the significant memory overhead of the original algorithm. We demonstrate knor outperforms distributed commercial products like H2O, Turi (formerly Dato, GraphLab) and Spark's MLlib by more than an order of magnitude for datasets of 107 to 109 points. Disa Mhembere, Da Zheng 0004, Carey E. Priebe, Joshua T. Vogelstein, Randal C. Burns |
HPDC | 5 |
| 2017 | Semi-External Memory Sparse Matrix Multiplication for Billion-Node GraphsabstractSparse matrix multiplication is traditionally performed in memory and scales to large matrices using the distributed memory of multiple nodes. In contrast, we scale sparse matrix multiplication beyond memory capacity by implementing sparse matrix dense matrix multiplication (SpMM) in a semi-external memory (SEM) fashion; i.e., we keep the sparse matrix on commodity SSDs and dense matrices in memory. Our SEM-SpMM incorporates many in-memory optimizations for large power-law graphs. It outperforms the in-memory implementations of Trilinos and Intel MKL and scales to billion-node graphs, far beyond the limitations of memory. Furthermore, on a single large parallel machine, our SEM-SpMM operates as fast as the distributed implementations of Trilinos using five times as much processing power. We also run our implementation in memory (IM-SpMM) to quantify the overhead of keeping data on SSDs. SEM-SpMM achieves almost 100 percent performance of IM-SpMM on graphs when the dense matrix has more than four columns; it achieves at least 65 percent performance of IM-SpMM on all inputs. We apply our SpMM to three important data analysis tasks-PageRank, eigensolving, and non-negative matrix factorization-and show that our SEM implementations significantly advance the state of the art. Da Zheng 0004, Disa Mhembere, Vince Lyzinski, Joshua T. Vogelstein, Carey E. Priebe, Randal C. Burns |
IEEE Trans. Parallel Distributed Syst. | 6 |
| 2015 | VESICLE: Volumetric Evaluation of Synaptic Inferfaces using Computer Vision at Large Scale
William R. Gray Roncal, Michael J. Pekala, Verena Kaynig, Dean Kleissas, Joshua T. Vogelstein, Hanspeter Pfister, Randal C. Burns, R. Jacob Vogelstein, Mark A. Chevillet, Gregory D. Hager |
BMVC | 7 |
| 2015 | Streaming Algorithms for Halo FindersabstractCosmological N-body simulations are essential for studies of the large-scale distribution of matter and galaxies in the Universe. This analysis often involves finding clusters of particles and retrieving their properties. Detecting such "halos" among a very large set of particles is a computationally intensive problem, usually executed on the same super-computers that produced the simulations, requiring huge amounts of memory. Recently, a new area of computer science emerged. This area, called streaming algorithms, provides new theoretical methods to compute data analytics in a scalable way using only a single pass over a data sets and logarithmic memory. The main contribution of this paper is a novel connection between the N-body simulations and the streaming algorithms. In particular, we investigate a link between halo finders and the problem of finding frequent items (heavy hitters) in a data stream, that should greatly reduce the computational resource requirements, especially the memory needs. Based on this connection, we can build a new halo finder by running efficient heavy hitter algorithms as a black-box. We implement two representatives of the family of heavy hitter algorithms, the Count-Sketch algorithm (CS) and the Pick-and-Drop sampling (PD), and evaluate their accuracy and memory usage. Comparison with other halo-finding algorithms from [1] shows that our halo finder can locate the largest haloes using significantly smaller memory space and with comparable running time. This streaming approach makes it possible to run and analyze extremely large data sets from N-body simulations on a smaller machine, rather than on supercomputers. Our findings demonstrate the connection between the halo search problem and streaming algorithms as a promising initial direction of further research. Zaoxing Liu, Nikita Ivkin, Lin Yang 0011, Mark Neyrinck, Gerard Lemson, Alex Szalay, Vladimir Braverman, Tamás Budavári, Randal C. Burns |
e-Science | 9 |
| 2015 | Efficient evaluation of threshold queries of derived fields in a numerical simulation databaseabstractIn this paper, we present a method for the ecient evaluation of threshold queries of derived fields for large numerical simulation datasets stored in a cluster of relational databases. The datasets produced by these simulations are in the TB and even PB ranges. Data-intensive computations that examine entire time-steps of the simulation data are impractical to perform locally by the user, taking days or months to iterate over the entire dataset. The integrated method for the evaluation of threshold queries that we have developed achieves scalability through data-parallel execution of the computations on the nodes of an analysis database cluster. We extend the scientific analysis environment with the introduction of an application-aware cache for query results, building on the concept of semantic caching. The cache has little overhead and improves query performance by over an order of magnitude for queries that hit the cache. Caching the results of threshold queries preserves both the I/O and computation e↵ort used to obtain them. In the case of computational turbulence, this allows scientists to quickly focus on the most intense events and interesting regions in any time-step or the dataset as a whole, which greatly speeds up the rate of scientific exploration and discovery. Kalin Kanov, Randal C. Burns, Cristian Constantin Lalescu |
EDBT | 2 |
| 2015 | FlashGraph: Processing Billion-Node Graphs on an Array of Commodity SSDs
Da Zheng 0004, Disa Mhembere, Randal C. Burns, Joshua T. Vogelstein, Carey E. Priebe, Alex Szalay |
FAST | 3 |
| 2015 | Particle tracking in open simulation laboratoriesabstractParticle tracking along streamlines and pathlines is a common scientific analysis technique, which has demanding data, computation and communication requirements. It has been studied in the context of high-performance computing due to the difficulty in its efficient parallelization and its high demands on communication and computational load. In this paper, we study efficient evaluation methods for particle tracking in open simulation laboratories. Simulation laboratories have a fundamentally different architecture from today's supercomputers and provide publicly-available analysis functionality. We focus on the I/O demands of particle tracking for numerical simulation datasets 100s of TBs in size. We compare data-parallel and task-parallel approaches for the advection of particles and show scalability results on data-intensive workloads from a live production environment. We have developed particle tracking capabilities for the Johns Hopkins Turbulence Databases, which store computational fluid dynamics simulation data, including forced isotropic turbulence, magnetohydrodynamics, channel flow turbulence and homogeneous buoyancy-driven turbulence. Kalin Kanov, Randal C. Burns |
SC | 2 |
| 2013 | Toward millions of file system IOPS on low-cost, commodity hardwareabstractWe describe a storage system that removes I/O bottlenecks to achieve more than one million IOPS based on a user-space file abstraction for arrays of commodity SSDs. The file abstraction refactors I/O scheduling and placement for extreme parallelism and non-uniform memory and I/O. The system includes a set-associative, parallel page cache in the user space. We redesign page caching to eliminate CPU overhead and lock-contention in non-uniform memory architecture machines. We evaluate our design on a 32 core NUMA machine with four, eight-core processors. Experiments show that our design delivers 1.23 million 512-byte read IOPS. The page cache realizes the scalable IOPS of Linux asynchronous I/O (AIO) and increases user-perceived I/O performance linearly with cache hit rates. The parallel, set-associative cache matches the cache hit rates of the global Linux page cache under real workloads. Da Zheng 0004, Randal C. Burns, Alex Szalay |
SC | 2 |
| 2013 | The open connectome project data cluster: scalable analysis and vision for high-throughput neuroscienceabstract- neural connectivity maps of the brain-using the parallel execution of computer vision algorithms on high-performance compute clusters. These services and open-science data sets are publicly available at openconnecto.me. The system design inherits much from NoSQL scale-out and data-intensive computing architectures. We distribute data to cluster nodes by partitioning a spatial index. We direct I/O to different systems-reads to parallel disk arrays and writes to solid-state storage-to avoid I/O interference and maximize throughput. All programming interfaces are RESTful Web services, which are simple and stateless, improving scalability and usability. We include a performance evaluation of the production system, highlighting the effec-tiveness of spatial data organization. Randal C. Burns, Kunal Lillaney, Daniel R. Berger, Logan Grosenick, Karl Deisseroth, R. Clay Reid, William R. Gray Roncal, Priya Manavalan, Davi Bock, Narayanan Kasthuri, Michael M. Kazhdan, Stephen J. Smith, Dean Kleissas, Eric A. Perlman, Kwanghun Chung, Nicholas C. Weiler, Jeff Lichtman, Alex Szalay, Joshua T. Vogelstein, R. Jacob Vogelstein |
SSDBM | 1 |
| 2013 | Inverted indices for particle tracking in petascale cosmological simulationsabstractWe describe the challenges arising from tracking dark matter particles in state of the art cosmological simulations. We are in the process of running the Indra suite of simulations, with an aggregate count of more than 35 trillion particles and 1.1PB of total raw data volume. However, it is not enough just to store the particle positions and velocities in an efficient manner -- analyses also need to be able to track individual particles efficiently through the temporal history of the simulation. The required inverted indices can easily have raw sizes comparable to the original simulation. Daniel Crankshaw, Randal C. Burns, Bridget Falck, Tamás Budavári, Alex Szalay, Jie Wang 0075 |
SSDBM | 2 |
| 2012 | Rethinking erasure codes for cloud file systems: minimizing I/O for recovery and degraded reads
Randal C. Burns, James S. Plank, William Pierce |
FAST | 2 |
| 2012 | A Parallel Page Cache: IOPS and Caching for Multicore Systems
Da Zheng 0004, Randal C. Burns, Alex Szalay |
HotStorage | 2 |
| 2012 | Data-intensive spatial filtering in large numerical simulation datasetsabstractWe present a query processing framework for the efficient evaluation of spatial filters on large numerical simulation datasets stored in a data-intensive cluster. Previously, filtering of large numerical simulations stored in scientific databases has been impractical owing to the immense data requirements. Rather, filtering is done during simulation or by loading snapshots into the aggregate memory of an HPC cluster. Our system performs filtering within the database and supports large filter widths. We present two complementary methods of execution: I/O streaming computes a batch filter query in a single sequential pass using incremental evaluation of decomposable kernels, summed volumes generates an intermediate data set and evaluates each filtered value by accessing only eight points in this dataset. We dynamically choose between these methods depending upon workload characteristics. The system allows us to perform filters against large data sets with little overhead: query performance scales with the cluster's aggregate I/O throughput. Kalin Kanov, Randal C. Burns, Gregory L. Eyink, Charles Meneveau, Alex Szalay |
SC | 2 |
| 2011 | CoScan: cooperative scan sharing in the cloudabstractWe present CoScan, a scheduling framework that eliminates redundant processing in workflows that scan large batches of data in a map-reduce computing environment. CoScan merges Pig programs from multiple users at runtime to reduce I/O contention while adhering to soft deadline requirements in scheduling. This includes support for join workflows that operate on multiple data sources. Our solution maps well to workflows at many Internet companies which reuse data from a common set of inputs. Experiments on the PigMix data analytics benchmark exhibit orders of magnitude reduction in resource contention with minimal impact on latency. Christopher Olston, Anish Das Sarma, Randal C. Burns |
SoCC | 4 |
| 2011 | In Search of I/O-Optimal Recovery from Disk Failures
Randal C. Burns, James S. Plank |
HotStorage | 2 |
| 2011 | MPI-DB, A Parallel Database Services Software Library for Scientific Computing
Edward Givelberg, Alex Szalay, Kalin Kanov, Randal C. Burns |
EuroMPI | 4 |
| 2011 | I/O streaming evaluation of batch queries for data-intensive computational turbulenceabstractWe describe a method for evaluating computational turbulence queries, including Lagrange Polynomial interpolation, based on partial sums that allows the underlying data to be accessed in any order and in parts. We exploit these properties to stream data from disk in a single pass and concurrently evaluate batch queries. The combination of sequential I/O and data sharing improves performance by an order of magnitude when compared with direct evaluation of each query. The technique also supports distributed evaluation of queries in a database cluster, assembling the partial sums from each node at the query mediator. Interpolation is fundamental to computational turbulence, over 95% of queries use these routines, and the partial sums method allows the JHU Turbulence Database Cluster to realize scale and throughput for our scientists' data-intensive workloads. Kalin Kanov, Eric A. Perlman, Randal C. Burns, Yanif Ahmad, Alex Szalay |
SC | 3 |
| 2011 | Remote data checking using provable data possessionabstractWe introduce a model for provable data possession (PDP) that can be used for remote data checking: A client that has stored data at an untrusted server can verify that the server possesses the original data without retrieving it. The model generates probabilistic proofs of possession by sampling random sets of blocks from the server, which drastically reduces I/O costs. The client maintains a constant amount of metadata to verify the proof. The challenge/response protocol transmits a small, constant amount of data, which minimizes network communication. Thus, the PDP model for remote data checking is lightweight and supports large data sets in distributed storage systems. The model is also robust in that it incorporates mechanisms for mitigating arbitrary amounts of data corruption. We present two provably-secure PDP schemes that are more efficient than previous solutions. In particular, the overhead at the server is low (or even constant), as opposed to linear in the size of the data. We then propose a generic transformation that adds robustness to any remote data checking scheme based on spot checking. Experiments using our implementation verify the practicality of PDP and reveal that the performance of PDP is bounded by disk I/O and not by cryptographic computation. Finally, we conduct an in-depth experimental evaluation to study the tradeoffs in performance, security, and space overheads when adding robustness to a remote data checking scheme. Giuseppe Ateniese, Randal C. Burns, Reza Curtmola, Joseph Herring, Lea Kissner, Zachary N. J. Peterson, Dawn Song |
ACM Trans. Inf. Syst. Secur. | 2 |
| 2010 | JAWS: Job-Aware Workload Scheduling for the Exploration of Turbulence SimulationsabstractWe present JAWS, a job-aware, data-driven batch scheduler that improves query throughput for data-intensive scientific database clusters. As datasets reach petabyte-scale, workloads that scan through vast amounts of data to extract features are gaining importance in the sciences. However, acute performance bottlenecks result when multiple queries execute simultaneously and compete for I/O resources. Our solution, JAWS, divides queries into I/O-friendly sub-queries for scheduling. It then identifies overlapping data requirements within the workload and executes sub-queries in batches to maximize data sharing and reduce redundant I/O. JAWS extends our previous work by supporting workflows in which queries exhibit data dependencies, exploiting workload knowledge to coordinate caching decisions, and combating starvation through adaptive and incremental trade-offs between query throughput and response time. Instrumenting JAWS in the Turbulence Database Cluster yields nearly three-fold improvement in query throughput when contention in the workload is high. Eric A. Perlman, Randal C. Burns, Tanu Malik, Tamás Budavári, Charles Meneveau, Alex Szalay |
SC | 3 |
| 2010 | Organization of Data in Non-convex Spatial Domains
Eric A. Perlman, Randal C. Burns, Michael M. Kazhdan, Rebecca R. Murphy, William P. Ball, Nina Amenta |
SSDBM | 2 |
| 2010 | Guest editorial: FAST'10abstractNo abstract available. Randal C. Burns, Kimberly Keeton |
ACM Trans. Storage | 1 |
| 2009 | LifeRaft: Data-Driven, Batch Processing for the Exploration of Scientific Databases
Randal C. Burns, Tanu Malik |
CIDR | 2 |
| 2009 | CA-NFS: A Congestion-Aware Network File System
Alexandros Batsakis, Randal C. Burns, Arkady Kanevsky, James Lentini, Thomas Talpey |
FAST | 2 |
| 2009 | Adaptive Physical Design for Curated Archives
Tanu Malik, Debabrata Dash, Amitabh Chaudhary, Anastasia Ailamaki, Randal C. Burns |
SSDBM | 6 |
| 2009 | CA-NFS: A congestion-aware network file systemabstractWe develop a holistic framework for adaptively scheduling asynchronous requests in distributed file systems. The system is holistic in that it manages all resources, including network bandwidth, server I/O, server CPU, and client and server memory utilization. It accelerates, defers, or cancels asynchronous requests in order to improve application-perceived performance directly. We employ congestion pricing via online auctions to coordinate the use of system resources by the file system clients so that they can detect shortages and adapt their resource usage. We implement our modifications in the Congestion-Aware Network File System (CA-NFS), an extension to the ubiquitous network file system (NFS). Our experimental result shows that CA-NFS results in a 20% improvement in execution times when compared with NFS for a variety of workloads. Alexandros Batsakis, Randal C. Burns, Arkady Kanevsky, James Lentini, Thomas Talpey |
ACM Trans. Storage | 2 |
| 2008 | Workload-Aware Histograms for Remote Applications
Tanu Malik, Randal C. Burns |
DaWaK | 2 |
| 2008 | AWOL: An Adaptive Write Optimizations Layer
Alexandros Batsakis, Randal C. Burns, Arkady Kanevsky, James Lentini, Thomas Talpey |
FAST | 2 |
| 2008 | MR-PDP: Multiple-Replica Provable Data PossessionabstractMany storage systems rely on replication to increase the availability and durability of data on untrusted storage systems. At present, such storage systems provide no strong evidence that multiple copies of the data are actually stored. Storage servers can collude to make it look like they are storing many copies of the data, whereas in reality they only store a single copy. We address this shortcoming through multiple-replica provable data possession (MR-PDP): A provably-secure scheme that allows a client that stores t replicas of a file in a storage system to verify through a challenge-response protocol that (1) each unique replica can be produced at the time of the challenge and that (2) the storage system uses t times the storage required to store a single replica. MR-PDP extends previous work on data possession proofs for a single copy of a file in a client/server storage system (Ateniese et al., 2007). Using MR-PDP to store t replicas is computationally much more efficient than using a single-replica PDP scheme to store t separate, unrelated files (e.g., by encrypting each file separately prior to storing it). Another advantage of MR-PDP is that it can generate further replicas on demand, at little expense, when some of the existing replicas fail. Reza Curtmola, Randal C. Burns, Giuseppe Ateniese |
ICDCS | 3 |
| 2008 | Scientific Data Management: An Orphan in the Database Community?
Randal C. Burns, Susan B. Davidson, Yannis E. Ioannidis, Miron Livny, Jignesh M. Patel |
ICDE | 1 |
| 2008 | Network-Aware Join Processing in Global-Scale Database FederationsabstractWe introduce join scheduling algorithms that employ a balanced network utilization metric to optimize the use of all network paths in a global-scale database federation. This metric allows algorithms to exploit excess capacity in the network, while avoiding narrow, long-haul paths. We give a two- approximate, polynomial-time algorithm for serial (left-deep) join schedules. We also present extensions to this algorithm that explore parallel schedules, reduce resource usage, and define tradeoffs between computation and network utilization. We evaluate these techniques within the SkyQuery federation of Astronomy databases using spatial-join queries submitted by SkyQuery's users. Experiments show that our algorithms realize near-optimal network utilization with minor computational overhead. Randal C. Burns, Andreas Terzis, Amol Deshpande |
ICDE | 2 |
| 2008 | Organizing and indexing non-convex regionsabstractWe demonstrate data indexing and query processing techniques that improve the efficiency of comparing, correlating, and joining data contained in non-convex regions. We use computational geometry techniques to automatically characterize the region of space from which data are drawn, partition the region based on that characterization, and create an index from the partitions. Our motivating application performs distributed data analysis queries among federated database sites that store scientific data sets from the Chesapeake Bay. Our preliminary findings indicate that these techniques often reduce the number of I/Os needed to serve a query by a factor of five---depending on the geometry of the query region. Eric A. Perlman, Randal C. Burns, Michael M. Kazhdan |
Proc. VLDB Endow. | 2 |
| 2008 | NFS-CD: Write-Enabled Cooperative Caching in NFSabstractWe present the network file system with cluster delegation (NFS-CD), an enhancement to the NFSv4 that reduces server load and increases the scalability of distributed file systems in computing clusters. The cluster delegation feature of NFS-CD allows data sharing among clients by extending the NFSv4 delegation model so that multiple clients manage a single file without interacting with the server. Based on cluster delegation, we implement a fast-commit primitive, cooperative caching, and the ability to recover the uncommitted updates of a failed computer. NFS-CD supports both read and write operations in the cooperative cache without degrading the consistency model of NFSv4. We have implemented NFS-CD by modifying the Linux NFSv4 client only. Because the server remains unchanged, NFS-CD preserves the simple administration model of NFSv4 and interoperates with standard NFS clients. NFS-CD offers improved performance when compared to NFSv4 at the expense of slightly weaker reliability guarantees. An experimental evaluation of our implementation, using industry standard benchmarks and application workloads, reveals that NFS-CD reduces server load by more than half. It also demonstrates that under most workloads, file systems must support writes to the cooperative cache to achieve scale. Alexandros Batsakis, Randal C. Burns |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2007 | Provable data possession at untrusted storesabstractWe introduce a model for provable data possession (PDP) that allows a client that has stored data at an untrusted server to verify that the server possesses the original data without retrieving it. The model generates probabilistic proofs of possession by sampling random sets of blocks from the server, which drastically reduces I/O costs. The client maintains a constant amount of metadata to verify the proof. The challenge/response protocol transmits a small, constant amount of data, which minimizes network communication. Thus, the PDP model for remote data checking supports large data sets in widely-distributed storage system. Giuseppe Ateniese, Randal C. Burns, Reza Curtmola, Joseph Herring, Lea Kissner, Zachary N. J. Peterson, Dawn Song |
CCS | 2 |
| 2007 | A Black-Box Approach to Query Cardinality Estimation
Tanu Malik, Randal C. Burns, Nitesh V. Chawla |
CIDR | 2 |
| 2007 | A Workload-Driven Unit of Cache Replacement for Mid-Tier Database Caching
Tanu Malik, Randal C. Burns, Stratos Papadomanolakis, Anastasia Ailamaki |
DASFAA | 3 |
| 2007 | Design and Implementation of Verifiable Audit Trails for a Versioning File System
Zachary N. J. Peterson, Randal C. Burns, Giuseppe Ateniese, Stephen Bono |
FAST | 2 |
| 2007 | Data exploration of turbulence simulations using a database clusterabstractWe describe a new environment for the exploration of turbulent flows that uses a cluster of databases to store complete histories of Direct Numerical Simulation (DNS) results. This allows for spatial and temporal exploration of high-resolution data that were traditionally too large to store and too computationally expensive to produce on demand. We perform analysis of these data directly on the databases nodes, which minimizes the volume of network traffic. The low network demands enable us to provide public access to this experimental platform and its datasets through Web services. This paper details the system design and implementation. Specifically, we focus on hierarchical spatial indexing, cache-sensitive spatial scheduling of batch workloads, localizing computation through data partitioning, and load balancing techniques that minimize data movement. We provide real examples of how scientists use the system to perform high-resolution turbulence research from standard desktop computing environments. Eric A. Perlman, Randal C. Burns, Charles Meneveau |
SC | 2 |
| 2007 | Multilevel streaming for out-of-core surface reconstruction
Matthew Bolitho, Michael M. Kazhdan, Randal C. Burns, Hugues Hoppe |
Symposium on Geometry Processing | 3 |
| 2006 | Improving I/O Performance of Clustered Storage Systems by Adaptive Request DistributionabstractWe develop an adaptive load distribution protocol for logical volume I/O workload in clustered storage systems. It exploits data redundancy among decentralized storage servers to dynamically route I/O workload on a per-request basis, offering short-term load balancing and improved I/O performance. Our protocol builds on tunable hashing techniques and is based purely on client logic. Therefore, it does not limit system scalability and requires no change to the existing infrastructure. It distributes the I/O requests of a client to storage servers selected adoptively by a decentralized tunable hashing scheme, and, applies different policies to read and write requests. It also makes no assumption about inter-server communication latency and thus is robust to different network configurations. It supports both replication and erasure coding data redundancy schemes. Experimental results show that our protocol performs closely to a centralized load-balancing algorithm and verify the robustness of our protocol Changxun Wu, Randal C. Burns |
HPDC | 2 |
| 2006 | Data management and query - Estimating query result sizes for proxy caching in scientific database federationsabstractIn a proxy cache for federations of scientific databases it is important to estimate the size of a query before making a caching decision. With accurate estimates, near-optimal cache performance can be obtained. On the other extreme, inaccurate estimates can render the cache totally ineffective. We present classification and regression over templates (CAROT), a general method for estimating query result sizes, which is suited to the resource-limited environment of proxy caches and the distributed nature of database federations. CAROT estimates query result sizes by learning the distribution of query results, not by examining or sampling data, but from observing workload. We have integrated CAROT into the proxy cache of the National Virtual Observatory (NVO) federation of astronomy databases. Experiments conducted in the NVO show that CAROT dramatically outperforms conventional estimation techniques and provides near-optimal cache performance. Tanu Malik, Randal C. Burns, Nitesh V. Chawla, Alex Szalay |
SC | 2 |
| 2006 | Poster reception - Engineering the 100 terabyte turbulence database (or how to track particles at home)abstractWe describe a new environment for large-scale turbulence simulations that uses a cluster of database nodes to store the complete space-time history of fluid velocities. This allows for rapid access to high resolution data that were traditionally too large to store and too computationally expensive to produce on demand.We perform the actual experimental analysis inside the database nodes, which allows for data-intensive computations to be performed across a large number of nodes with relatively little network traffic.We currently have a limited-scale prototype system running actual turbulence simulations and are in the process of establishing a production cluster with high-resolution data. We will discuss our design choices and initial results with load balancing a data-intensive, migratory workload. Eric A. Perlman, Randal C. Burns |
SC | 2 |
| 2006 | Data analysis tools for sensor-based scienceabstractScience is increasingly driven by data collected automatically from arrays of inexpensive sensors. The collected data volumes require a different approach from the scientists' current Excel spreadsheet storage and analysis model. Spreadsheets work well for small data sets; but scientists want high level summaries of their data for various statistical analyses without sacrificing the ability to drill down to every bit of the raw data. This demonstration describes our prototype data analysis system that is suitable for browsing and visualization - like a spreadsheet - but scalable to much larger data sets. Stuart Ozer, Jim Gray 0001, Alex Szalay, Andreas Terzis, Razvan Musaloiu-Elefteri, Katalin Szlavecz, Randal C. Burns, Joshua Cogan |
SenSys | 7 |
| 2005 | Secure Deletion for a Versioning File System
Zachary N. J. Peterson, Randal C. Burns, Joseph Herring, Adam Stubblefield, Aviel D. Rubin |
FAST | 2 |
| 2005 | Cluster delegation: high-performance, fault-tolerant data sharing in NFSabstractWe present cluster delegation, an enhancement to the NFSv4 file system, that improves both performance and recoverability in computing clusters. Cluster delegation allows data sharing among clients by extending the NFSv4 delegation model so that multiple clients manage a single file without interacting with the server. Based on cluster delegation, we implement a fast commit primitive, cooperative caching, and the ability to recover the uncommitted updates of a failed computer. Cluster delegation supports both read and write operations in the cooperative cache, while preserving the consistency guarantees of NFSv4. We have implemented cluster delegation by modifying the Linux NFSv4 client and show that it improves client performance and reduces server load by more than half. Alexandros Batsakis, Randal C. Burns |
HPDC | 2 |
| 2005 | Bypass Caching: Making Scientific Databases Good Network CitizensabstractScientific database federations are geographically distributed and network bound. Thus, they could benefit from proxy caching. However, existing caching techniques are not suitable for their workloads, which compare and join large data sets. Existing techniques reduce parallelism by conducting distributed queries in a single cache and lose the data reduction benefits of performing selections at each database. We develop the bypass-yield formulation of caching, which reduces network traffic in wide-area database federations, while preserving parallelism and data reduction. Bypass-yield caching is altruistic; caches minimize the overall network traffic generated by the federation, rather than focusing on local performance. We present an adaptive, workload-driven algorithm for managing a bypass-yield cache. We also develop on-line algorithms that make no assumptions about workload: a k-competitive deterministic algorithm and a randomized algorithm with minimal space complexity. We verify the efficacy of bypass-yield caching by running workload traces collected from the Sloan Digital Sky Survey through a prototype implementation. Tanu Malik, Randal C. Burns, Amitabh Chaudhary |
ICDE | 2 |
| 2005 | Ext3cow: a time-shifting file system for regulatory complianceabstractThe ext3cow file system, built on the popular ext3 file system, provides an open-source file versioning and snapshot platform for compliance with the versioning and audtitability requirements of recent electronic record retention legislation. Ext3cow provides a time-shifting interface that permits a real-time and continuous view of data in the past. Time-shifting does not pollute the file system namespace nor require snapshots to be mounted as a separate file system. Further, ext3cow is implemented entirely in the file system space and, therefore, does not modify kernel interfaces or change the operation of other file systems. Ext3cow takes advantage of the fine-grained control of on-disk and in-memory data available only to a file system, resulting in minimal degradation of performance and functionality. Experimental results confirm this hypothesis; ext3cow performs comparably to ext3 on many benchmarks and on trace-driven experiments. Zachary N. J. Peterson, Randal C. Burns |
ACM Trans. Storage | 2 |
| 2005 | Tunable randomization for load management in shared-disk clustersabstractWe develop and evaluate a system for load management in shared-disk file systems built on clusters of heterogeneous computers. It balances workload by moving file sets among cluster server nodes. It responds to changing server resources that arise from failure and recovery, and dynamically adding or removing servers. It also realizes performance consistency---nearly uniform performance across all servers. The system is adaptive and self-tuning. It operates without any a priori knowledge of workload properties, or the capabilities of the servers. Rather, it continuously tunes load placement using a technique called adaptive, nonuniform (ANU) randomization. ANU randomization realizes the scalability and metadata reduction benefits of hash-based, randomized placement techniques, while avoiding hashing's drawbacks: load skew, inability to cope with heterogeneity, and lack of tunability. ANU randomization outperforms virtual-processor approaches to load balancing, while reducing the amount of shared state among servers and the amount of load movement. Changxun Wu, Randal C. Burns |
ACM Trans. Storage | 2 |
| 2004 | Achieving Performance Consistency in Heterogeneous Clusters
Changxun Wu, Randal C. Burns |
HPDC | 2 |
| 2004 | Fastpath Optimizations for Cluster Recovery in Shared-Disk SystemsabstractWe describe the design and implementation of a clustering service for a high-performance, shared-disk file system. The service provides failure detection and recovery, reliable end-to-end messaging, and a centralized and recoverable management interface. We implement novel optimizations in the voting protocol that resolves cluster membership. Optimizations allow clusters to form as quickly as possible without introducing livelock or requiring timeout parameters to be tuned carefully. Our treatment includes performance results that quantify the scalability of the system and measure recovery times. Randal C. Burns |
SC | 1 |
| 2003 | Handling Heterogeneity in Shared-Disk File SystemsabstractWe develop and evaluate a system for load management in shared-disk file systems built on clusters of heterogeneous computers. The system generalizes load balancing and server provisioning. It balances file metadata workload by moving file sets among cluster server nodes. It also responds to changing server resources that arise from failure and recovery and dynamically adding or removing servers. The system is adaptive and self-managing. It operates without any a-priori knowledge of workload properties or the capabilities of the servers. Rather, it continuously tunes load placement using a technique called adaptive, non-uniform (ANU) randomization. ANU randomization realizes the scalability and metadata reduction benefits of hash-based, randomized placement techniques. It also avoids hashing's drawbacks: load skew, inability to cope with heterogeneity, and lack of tunability. Simulation results show that our load-management algorithm performs comparably to a prescient algorithm. Changxun Wu, Randal C. Burns |
SC | 2 |
| 2003 | In-Place Reconstruction of Version DifferencesabstractIn-place reconstruction of differenced data allows information on devices with limited storage capacity to be updated efficiently over low-bandwidth channels. Differencing encodes a version of data compactly as a set of changes from a previous version. Transmitting updates to data as a version difference saves both time and bandwidth. In-place reconstruction rebuilds the new version of the data in the storage or memory the current version occupies-no scratch space is needed for a second version. By combining these technologies, we support highly mobile applications on space-constrained hardware. We present an algorithm that modifies a differentially encoded version to be in-place reconstructible. The algorithm trades a small amount of compression to achieve this property. Our treatment includes experimental results that show our implementation to be efficient in space and time and verify that compression losses are small. Also, we give results on the computational complexity of performing this modification while minimizing lost compression. Randal C. Burns, Larry J. Stockmeyer, Darrell D. E. Long |
IEEE Trans. Knowl. Data Eng. | 1 |
| 2002 | Group-Based Management of Distributed File CachesabstractWe describe a way to manage distributed file system caches based upon groups of files that are accessed together. We use file access patterns to automatically construct dynamic groupings of files and then manage our cache by fetching groups, rather than single files. We present experimental results, based on trace-driven workloads, demonstrating that grouping improves cache performance. At the file system client, grouping can reduce LRU demand fetches by 50 to 60%. At the server cache hit rate improvements are much more pronounced, but vary widely (20 to over 1200%) depending upon the capacity of intervening caches. Our treatment includes information theoretic results that justify our approach to file grouping. Ahmed Amer, Darrell D. E. Long, Randal C. Burns |
ICDCS | 3 |
| 2002 | Compactly encoding unstructured inputs with differential compressionabstractThe subject of this article is differential compression , the algorithmic task of finding common strings between versions of data and using them to encode one version compactly by describing it as a set of changes from its companion. A main goal of this work is to present new differencing algorithms that (i) operate at a fine granularity (the atomic unit of change), (ii) make no assumptions about the format or alignment of input data, and (iii) in practice use linear time, use constant space, and give good compression. We present new algorithms, which do not always compress optimally but use considerably less time or space than existing algorithms. One new algorithm runs in O ( n ) time and O (1) space in the worst case (where each unit of space contains [log n ] bits), as compared to algorithms that run in O ( n ) time and O ( n ) space or in O ( n 2 ) time and O (1) space. We introduce two new techniques for differential compression and apply these to give additional algorithms that improve compression and time performance. We experimentally explore the properties of our algorithms by running them on actual versioned data. Finally, we present theoretical results that limit the compression power of differencing algorithms that are restricted to making only a single pass over the data. Miklós Ajtai, Randal C. Burns, Ronald Fagin, Darrell D. E. Long, Larry J. Stockmeyer |
J. ACM | 2 |
| 2001 | An Analytical Study of Opportunistic Lease RenewalabstractWe present opportunistic renewal, a lease management protocol designed to keep distributed file systems or distributed shared memories consistent in the presence of a network partition or other computer failures.Our treatment includes an analytical model of the protocol that compares performance with existing lease protocols and quantifies improvements.In addition, this analytical model provides the structure to understand message overhead and availability trade-offs when selecting lease parameters.We include results demonstrating that opportunistic renewal substantially reduces the network overhead associated with lease renewal.As a corollary, opportunistic renewal can reduce the lease length at any given network overhead; e.g., by a factor of 50 at 1% network overhead.Lower overhead makes leasing less intrusive and shorter lease periods allow a system to recover from failure more quickly. Randal C. Burns, Robert M. Rees, Darrell D. E. Long |
ICDCS | 1 |
| 2000 | Safe Caching in a Distributed File System for Network Attached StorageabstractIn a distributed file system built on network attached storage, client computers access data directly from shared storage, rather than submitting I/O requests through a server. Without a server marshaling access to data, if a computer fails or becomes isolated in a network partition while holding locks on cached data objects, those objects become inaccessible to other computers until a locking authority can guarantee that the lock holder will not again directly access these data. We describe a server that acts as the locking authority and implements a lease-based protocol for revoking access to data objects locked by an isolated or failed computer. When a lease expires, the server can be assured that the client no longer acts on locked data, and can safely redistribute locks to other clients. During normal operation, this protocol invokes no message overhead, and uses no memory and performs no computation at the locking authority. Randal C. Burns, Robert M. Rees, Darrell D. E. Long |
IPDPS | 1 |
| 1998 | In-Place Reconstruction of Delta Compressed FilesabstractAbstract results in high latency and low bandwidth to web-enabled clients and prevents the timely delivery of software. We present an algorithm for modifying delta compressed Differential or delta compression [5, 11, compactly enfiles so that the compressed versions may be reconstructed without scratch space. This allows network clients with limited resources to efficiently update software by retrieving delta compressed versions over a network. Delta compression for binary files, compactly encoding a version of data with only the changed bytes from a previous version, may be used to efficiently distribute software over low bandwidth channels, such as the Internet. Traditional methods for rebuilding these delta files require memory or storage space on the target machine for both the old and new version of the file to be reconstructed. With the advent of Randal C. Burns, Darrell D. E. Long |
PODC | 1 |