Wojciech M. Golab

dblp:15/5467 · also Wojciech Golab · DBLP profile ↗
← Back
72ranked-venue papers
22as first author
21since 2021 · last 2026
0000-0002-8891-256XORCID · verified

Domains — the database's venue-derived domains; a paper can count in several

Systems, architecture and hardware · 39 · 18 first-author · 10 since 2021Databases, data management, data science and information retrieval · 14 · 3 since 2021Applied, interdisciplinary, general and emerging computing · 6 · 2 since 2021Artificial intelligence and machine learning · 4 · 1 since 2021Theory of computation · 4 · 3 first-authorComputer networks · 2 · 1 first-author · 1 since 2021Security and privacy · 2 · 1 since 2021
YearPublicationVenuePosition
2026 Revisiting shared registers and leaderless consensus in WAN environments
Hao Tan 0006, Wojciech M. Golab, Vivek Alamuri
Distributed Parallel Databases2
2025 Brief Announcement: Using Detectability to Simplify the Design of Concurrent Algorithms for Persistent Memory
abstract
The emergence of persistent shared memory in recent years has created new opportunities to rethink classic algorithmic problems in the presence of crash-restart failures, and also new challenges as algorithm designers must account for both failures and concurrency. Recent research has explored algorithm design techniques based on powerful base objects that provide special features to simplify recovery from failures as an alternative to the traditional approach of programming more directly using processor instructions. Notably, Friedman, Herlihy, Marathe, and Petrank introduced detectable objects, which allow an application to resolve the outcome of operations that may have been interrupted by a crash. We propose a transformation that replaces low-level memory operations in a conventional algorithm with operations on detectable base objects that leverage persistent memory to enable recovery from crash-restart failures. The transformation is almost universal, and can be used to obtain efficient solutions to complex and previously unsolved problems.
Ahmed Fahmy, Wojciech M. Golab, Neeraj Mittal
PODC2
2025 Brief Announcement: Self-Stabilizing Recoverable Mutual Exclusion
abstract
We formalize and solve a novel variation of the classic mutual exclusion problem that tolerates both process crashes and memory corruptions, and we propose the first solutions based on a novel failure detector.
Wojciech M. Golab, Elad Michael Schiller
PODC1
2024 RMR-Efficient Detectable Objects for Persistent Memory and Their Applications
Sahil Dhoked, Ahmed Fahmy, Wojciech M. Golab, Neeraj Mittal
OPODIS3
2024 DULL: A Fast Scalable Detectable Unrolled Lock-Based Linked List
Ahmed Fahmy, Wojciech M. Golab
OPODIS2
2024 Brief Announcement: A Fast Scalable Detectable Unrolled Lock-Based Linked List
abstract
Persistent memory (PM) has emerged as a promising technology that enables data structure algorithms to preserve their consistent state after recovering from system failures. Detectable data structures have been proposed to resolve the response of the last operation of a crashed process. Designing detectable lock-based data structures is challenging due to the need to preserve the correctness properties of the underlying locks, such as mutual exclusion and deadlock-freedom, across failures. Therefore, lock-based detectable and persistent data structures are not as common as lock-free structures. In this work, we introduce DULL: a fast, scalable and Detectable Unrolled Lock-based Linked list. To the best of our knowledge, DULL is the fastest detectable lock-based linked list and the first detectable strictly-linearizable linked list. Experimental results indicate that DULL is several-fold faster than competitors in update workloads and equally fast in read-only scenarios.
Ahmed Fahmy, Wojciech M. Golab
SPAA2
2024 Antipaxos: Taking interactive consistency to the next level
abstract
Classical Paxos-like consensus protocols limit system scalability due to a single leader and the inability to process conflicting proposals in parallel. We introduce a novel agreement protocol, called Antipaxos, that instead reaches agreement on a collection of proposals using an efficient leaderless fast path when the environment is synchronous and failure-free, and falls back on a more elaborate slow path to handle other cases. We first specify the main safety property of Antipaxos by formalizing a new agreement problem called k-Interactive Consistency (k-IC). Then, we present a solution to this problem in the Byzantine failure model. We prove safety and liveness, and also present an experimental performance evaluation in the Amazon cloud. Our experiments show that Antipaxos achieves several-fold higher failure-free peak throughput than Mir-BFT. The inherent efficiency of our approach stems from the low message complexity of the fast path: agreement on n batches of conflict-prone proposals is achieved using only Θ(n2) messages in one consensus cycle, or Θ(n) amortized messages per batch.
Chunyu Mao, Wojciech M. Golab, Bernard Wong 0001
J. Parallel Distributed Comput.2
2023 Brief Announcement: On Solving Recoverable Mutual Exclusion Under System-Wide Failures
abstract
Recoverable mutual exclusion (RME) is a fault-tolerant variation of Dijkstra's classic mutual exclusion (ME) problem that allows processes to fail by crashing as long as they recover eventually. A growing body of literature on this topic, starting with the problem formulation by Golab and Ramaraju (PODC'16), examines the cost of solving the RME problem, which is quantified by counting expensive shared memory operations called remote memory references (RMRs), under a variety of conditions. Published results show that the RMR complexity of RME algorithms, among other factors, depends crucially on the failure model used: individual process versus system-wide. Recent work by Golab and Hendler (PODC'18) also suggests that explicit failure detection can be helpful in attaining constant RMR solutions to the RME problem in the system-wide failure model. Follow-up work by Jayanti et al. (SPAA'23) shows that such a solution exists even without employing a failure detector, albeit this solution uses a more complex algorithmic approach.
Sahil Dhoked, Wojciech M. Golab, Neeraj Mittal
SPAA2
2023 Brief Announcement: On Implementing Wear Leveling in Persistent Synchronization Structures
abstract
The last decade has witnessed an explosion of research on persistent memory, which combines the low access latency of dynamic random access memory (DRAM) with the durability of secondary storage. Intel’s implementation of persistent memory, called Optane, comes close to realizing the game-changing potential of persistent memory in terms of performance; however, it also suffers from limited endurance and relies on a proprietary wear leveling mechanism to mitigate memory cell wear-out. The traditional embedded approach to wear leveling, in which the storage device itself maps logical addresses to physical addresses, can be fast and energy-efficient, but it is also relatively inflexible and can lead to missed opportunities for optimization. An alternative school of thought, exemplified by "open channel" solid state drives (SSDs), delegates responsibility for wear leveling to software, where it can be tailored to specific applications. In this research, we consider a hypothetical hardware platform where the same paradigm is applied to the persistent memory device, and ask how the wear leveling mechanism can be co-designed with synchronization structures that generate highly skewed memory access patterns. Building on the recent work of Liu and Golab, we implement an improved wear leveling atomic counter by leveraging hardware transactional memory in a novel way. Our solution is close to optimal with respect to both space complexity and measured performance.
Jakeb Chouinard, Kush Kansara, Xialin Liu, Nihal Potdar, Wojciech M. Golab
DISC5
2023 Modular Recoverable Mutual Exclusion Under System-Wide Failures
abstract
Recoverable mutual exclusion (RME) is a fault-tolerant variation of Dijkstra’s classic mutual exclusion (ME) problem that allows processes to fail by crashing as long as they recover eventually. A growing body of literature on this topic, starting with the problem formulation by Golab and Ramaraju (PODC'16), examines the cost of solving the RME problem, which is quantified by counting the expensive shared memory operations called remote memory references (RMRs), under a variety of conditions. Published results show that the RMR complexity of RME algorithms, among other factors, depends crucially on the failure model used: individual process versus system-wide. Recent work by Golab and Hendler (PODC'18) also suggests that explicit failure detection can be helpful in attaining constant RMR solutions to the RME problem in the system-wide failure model. Follow-up work by Jayanti, Jayanti, and Joshi (SPAA'23) shows that such a solution exists even without employing a failure detector, albeit this solution uses a more complex algorithmic approach. In this work, we dive deeper into the study of RMR-optimal RME algorithms for the system-wide failure model, and present contributions along multiple directions. First, we introduce the notion of withdrawing from a lock acquisition rather than resetting the lock. We use this notion to design a withdrawable RME algorithm with optimal O(1) RMR complexity for both cache-coherent (CC) and distributed shared memory (DSM) models in a modular way without using an explicit failure detector. In some sense, our technique marries the simplicity of Golab and Hendler’s algorithm with Jayanti, Jayanti and Joshi’s weaker system model. Second, we present a variation of our algorithm that supports fully dynamic process participation (i.e., both joining and leaving) in the CC model, while maintaining its constant RMR complexity. We show experimentally that our algorithm is substantially faster than Jayanti, Jayanti, and Joshi’s algorithm despite having stronger correctness properties. Finally, we establish an impossibility result for fully dynamic RME algorithms with bounded RMR complexity in the DSM model that are adaptive with respect to space, and provide a wait-free withdraw section.
Sahil Dhoked, Wojciech M. Golab, Neeraj Mittal
DISC2
2022 Brief Announcement: Towards a Theory of Wear Leveling in Persistent Data Structures
abstract
The last decade has witnessed an explosion of research on persistent memory, covering both hardware implementations and software techniques. Research activities in this area are primarily driven by the performance benefits of persistent memory, which behaves like DRAM with respect to latency and yet provides the durability of secondary storage. These benefits can only be realized with efficient solutions to the problem of memory cell wear-out, which is one of the fundamental weaknesses of persistent memory versus DRAM, and has traditionally been addressed in hardware. In this paper, we consider the theoretical foundations of solving this problem in software, which allows for application-specific optimizations. Our main contributions are to formalize the problem, and present a novel software implementation of atomic Fetch-And-Increment (FAI) that internally uses multiple words of persistent memory to distribute wear.
Xialin Liu, Wojciech M. Golab
PODC2
2022 A NUMA-Aware Recoverable Mutex Lock
abstract
The mutual exclusion (ME) problem has been of interest to the scientific community since it was first defined by Dijkstra. Various algorithms have been developed to solve the problem, like the MCS and CLH queue-based locks. The problem was generalized into the recoverable mutual exclusion (RME) problem by Golab and Ramaraju to accommodate the possibility of process crash failures. Since then, multiple RME algorithms have been presented in the literature that vary in design and performance. Furthermore, non-uniform memory access (NUMA) architecture has become mainstream in designing modern distributed systems, stimulating the development of NUMA-aware mutex locks. None of the existing NUMA-aware mutex locks are recoverable to the best of our knowledge. In addition, none of the transformation techniques in the literature, such as flat-combining and cohort-locking, is a black-box transformation. Precisely, each of the existing transformation techniques requires specific characteristics of, and possible modifications to, the underlying NUMA-oblivious lock. In this work, we propose the Recoverable Filter (RF) lock, a black-box transformation approach that exploits memory locality to transform a NUMA-oblivious recoverable mutex lock into a NUMA-aware one. Practical experiments are conducted using two existing RME algorithms, Golab and Hendler's (GH) and Jayanti, Jayanti, and Joshi's (JJJ). The two RME locks are transformed into NUMA-aware locks using the proposed RF and the existing cohort algorithms. Results show that, in multi-socket configurations, our transformation boosts the performance of the NUMA-oblivious RME locks by up to 45%. The RME locks transformed using the proposed RF lock are slower than their non-recoverable cohort variants by up to 9%. Outcomes demonstrate that the overhead of our algorithm is minimal when using a single socket. Moreover, a deeper empirical assessment shows that the gap in performance between GH and JJJ is due to the entry section of JJJ, not its exit section.
Ahmed Fahmy, Wojciech M. Golab
SPAA2
2021 An Implementation of Fake News Prevention by Blockchain and Entropy-based Incentive Mechanism
abstract
Fake news is undoubtedly a significant threat to democratic countries nowadays because existing technologies can quickly and massively produce fake videos, articles, or social media messages based on the rapid development of artificial intelligence and deep learning. Therefore, human assistance is critical if current automatic fake new identification technologies desire to improve accuracy. Given this situation, prior research has proposed to add a quorum, a group of appraisers trusted by users to verify the authenticity of the information, to the fake news prevention systems. This paper proposes a stake-based incentive mechanism to diminish the negative effect of malicious behaviors on a quorum-based fake news prevention system. Moreover, we use Hyperledger Fabric, Schnorr signatures, and human appraisers to implement a practical prototype of a quorum-based fake news prevention system. Then we conduct necessary case analyses and experiments to realize how dishonest participants, crash failures, and scale impact our system. The outcomes of the case analyses and experiments show that our mechanisms are feasible and provide an analytical basis for developing fake news prevention systems.
Chien-Chih Chen, Richards Peter, Wojciech M. Golab
IEEE BigData4
2021 Brief Announcement: Detectable Sequential Specifications for Recoverable Shared Objects
abstract
The recent commercial release of persistent main memory by Intel has sparked intense interest in recoverable concurrent objects. Specifying and implementing such objects is technically challenging on current generation hardware precisely because the top layers of the memory hierarchy (CPU registers and cache) remain volatile, which causes application threads to lose critical execution state during a failure. Friedman, Herlihy, Marathe, and Petrank (DISC'17) recently proposed that this difficulty can be alleviated by making the recoverable objects detectable, meaning that during recovery, they can resolve the status of an operation that was interrupted by a failure. In this paper, we formalize this important concept using a detectable sequential specification (DSS), which augments an object's interface with auxiliary methods that threads use to first declare their need for detectability, and then perform detection if needed after a failure. Compared to prior work on this topic, our DSS-based approach is less reliant on assumptions regarding the system, and more flexible in the sense that it allows applications to request detectability on demand. As a proof of concept, we present a detectable recoverable lock-free queue algorithm and evaluate its performance on a multiprocessor equipped with Intel Optane persistent memory. Our queue outperforms the detectable log queue of Friedman, Herlihy, Marthe, and Petrank (PPoPP'18) by up to 1.7x.
Wojciech M. Golab
PODC2
2021 PHPRX: An Efficient Hash Table for Persistent Memory
abstract
Volatile media have dominated the realm of main memory on servers and desktop computers for decades. In 2019, Intel released the Optane Data Center Persistent Memory Module (DCPMM), which offers the capacity and persistence of block devices while providing the byte addressability and low latency of DRAM. These new memory modules allow programmers to develop data structures that can survive in main memory across crashes and power failures, without relying on secondary power sources such as batteries. This work presents the design of a persistent memory hash table data structure that incorporates several features to maximize efficiency: the locks for concurrency control are kept in volatile DRAM, an embedded memory allocator is used, a parallel table resize operation is implemented, and a mechanism is provided to incrementally expand the underlying memory-mapped file. We compare PHPRX experimentally against the Dash persistent memory hash table published recently by Lu et al., and demonstrate substantial speed-ups on an Intel Xeon server equipped with genuine Intel Optane DCPMM. Our performance advantage holds despite PHPRX using a space-efficient incremental approach to expanding the underlying memory-mapped file, as opposed to the much simpler static allocation approach used by Dash.
Diego Cepeda, Wojciech M. Golab
SPAA2
2021 A Scalable Recoverable Skip List for Persistent Memory
abstract
Interest in recoverable, persistent-memory-resident (PMEM-resident) data structures is growing as availability of Intel Optane Data Center Persistent Memory increases. An interesting use case for inmemory, recoverable data structures is for database indexes, which need high availability and reliability. RECIPE, a popular conversion technique to make existing, proven-correct algorithms recoverable, is limited to certain classes of algorithms and does not prescribe how to reference data stored in relocatable regions of memory. The Untitled Persistent Skip List (UPSkipList) is a PMEM-resident recoverable skip list derived from Herlihy et al.'s lock-free skip list algorithm. It is developed using a new conversion technique that extends the RECIPE algorithm by Lee et al. to work on lock-free algorithms with non-blocking writes and no inherent recovery mechanism. The algorithm is also extended to support concurrent data node splitting to improve performance. Comparison was done against the BzTree of Arulraj et al., as implemented by Lersch et al., which has non-blocking, non-repairing writes implemented using the persistent multi-word CAS (PMwCAS) primitive by Wang et al. Tested with the Yahoo Cloud Serving Benchmark (YCSB), UPSkipList achieves better performance in write-heavy workloads at high levels of concurrency than BzTree, showing that the extension to RECIPE is an effective alternative.
Sakib Chowdhury, Wojciech M. Golab
SPAA2
2021 Sharding Techniques in the Era of Blockchain
abstract
Blockchain is a peer-to-peer ledger that records a growing list of transactions in a tamper-resistant manner using cryptographic hashes. Centralized points of vulnerability are eliminated in blockchain, and so it is considered secure by design under some reasonable assumptions, such as honest majority. But scalability remains a major limitation that can be improved by sharding. Full sharding is one of the approaches in achieving the high performance of blockchain systems. This paper proposes a locality-based full sharding protocol in permissioned blockchains. We introduce a simple and efficient cross-shard transaction handling protocol. A prototype is under development based on Hyperledger Fabric.
Chunyu Mao, Wojciech M. Golab
SRDS2
2021 Detectable Sequential Specifications for Recoverable Shared Objects
abstract
The recent commercial release of persistent main memory by Intel has sparked intense interest in recoverable concurrent objects. Such objects maintain state in persistent memory, and can be recovered directly following a system-wide crash failure, as opposed to being painstakingly rebuilt using recovery state saved in slower secondary storage. Specifying and implementing recoverable objects is technically challenging on current generation hardware precisely because the top layers of the memory hierarchy (CPU registers and cache) remain volatile, which causes application threads to lose critical execution state during a failure. For example, a thread that completes an operation on a shared object and then crashes may have difficulty determining whether this operation took effect, and if so, what response it returned. Friedman, Herlihy, Marathe, and Petrank (DISC'17) recently proposed that this difficulty can be alleviated by making the recoverable objects detectable, meaning that during recovery, they can resolve the status of an operation that was interrupted by a failure. In this paper, we formalize this important concept using a detectable sequential specification (DSS), which augments an object’s interface with auxiliary methods that threads use to first declare their need for detectability, and then perform detection if needed after a failure. Our contribution is closely related to the nesting-safe recoverable linearizability (NRL) framework of Attiya, Ben-Baruch, and Hendler (PODC'18), which follows an orthogonal approach based on ordinary sequential specifications combined with a novel correctness condition. Compared to NRL, our DSS-based approach is more portable across different models of distributed computation, compatible with several existing linearizability-like correctness conditions, less reliant on assumptions regarding the system, and more flexible in the sense that it allows applications to request detectability on demand. On the other hand, application code assumes full responsibility for nesting DSS-based objects. As a proof of concept, we demonstrate the DSS in action by presenting a detectable recoverable lock-free queue algorithm and evaluating its performance on a multiprocessor equipped with Intel Optane persistent memory.
Wojciech M. Golab
DISC2
2021 Deadline-Aware Cost Optimization for Spark
abstract
We present OptEx, a closed-form model of job execution on Apache Spark, a popular parallel processing engine. To the best of our knowledge, OptEx is the first work that analytically models job completion time on Spark. The model can be used to estimate the completion time of a given Spark job on a cloud, with respect to the size of the input dataset, the number of iterations, and the number of nodes comprising the underlying cluster. Experimental results demonstrate that OptEx yields a mean relative error of 6 percent in estimating the job completion time. Furthermore, the model can be applied for estimating the cost-optimal cluster composition for running a given Spark job on a cloud under a completion deadline specified in theSLO(i.e., Service Level Objective). We show experimentally that OptEx is able to correctly estimate the required cluster composition for running a given Spark job under a given SLO deadline with an accuracy of 98 percent. We also provide a tool which can classify Spark jobs into job categories based on bisimilarity analysis on lineage graphs collected from the given jobs.
Subhajit Sidhanta, Wojciech M. Golab, Supratik Mukhopadhyay
IEEE Trans. Big Data2
2021 Optimizing All-to-All Data Transmission in WANs
abstract
All-to-all data transmission is a typical data transmission pattern in both consensus protocols and blockchain systems. Developing an optimization scheme that provides high throughput and low latency data transmission can significantly benefit the performance of those systems. This paper investigates the problem of optimizing all-to-all data transmission in a wide area network (WAN) using overlay multicast. We prove that in a hose network model, using shallow tree overlays with height up to two is sufficient for all-to-all data transmission to achieve the optimal throughput allowed by the available network resources. Upon this foundation, we build ShallowForest, a data plane optimization for consensus protocols and blockchain systems. The goal of ShallowForest is to improve consensus protocols’ resilience to skewed client load distribution. Experiments with skewed client load across replicas in the Amazon cloud demonstrate that ShallowForest can improve the commit throughput of the EPaxos consensus protocol by up to 100% with up to 60% reduction in commit latency.
Hao Tan 0006, Wojciech M. Golab
IEEE Trans. Netw. Serv. Manag.2
2021 Gossip-based visibility control for high-performance geo-distributed transactions
Hua Fan 0002, Wojciech M. Golab
VLDB J.2
2020 Energy-Efficient Energy Analytics Using a General Purpose Graphics Processing Unit
abstract
Smart meters allow energy providers to monitor their customers' power consumption. This fine-grained data stream generates many data points, which hides broader trends in power consumption and makes it difficult for energy providers to make decisions regarding a specific customer or a subset of customers. Since the raw power data has little direct use, various algorithms have been proposed to lower the dimensionality of data, discover trends, study relationships between different features of collected data, and summarize data. These analytical techniques make the data more palatable to the end user. Analyzing smart meter data is computationally intensive as there is a large number of households connected to one energy provider, and each household generates years of data at hourly intervals. To speed up the analysis, clusters of commodity computers have been used. Ironically, such clusters consume substantial energy - studies have shown that about 10% of the world-wide supply of electrical power is consumed by the computing infrastructure. In this paper, we describe the use of a graphics processing unit (GPU) to analyze smart meter data, and compare its performance with a conventional multi-core CPU. We discuss the technical challenges in programming a GPU effectively to process smart meter data, and demonstrate experimentally that this choice of implementation enables substantial improvements in terms of both running time and energy-efficiency as compared to the multi-core CPU.
Sagnik De, Wojciech M. Golab
IEEE BigData2
2020 The Recoverable Consensus Hierarchy
abstract
Herlihy's consensus hierarchy ranks the power of various synchronization primitives for solving consensus in a model where asynchronous processes communicate through shared memory, and may fail by halting. This paper revisits the consensus hierarchy in a model with crash-recovery failures, where the specification of consensus, called recoverable consensus (RC) in this paper, is weakened by allowing non-terminating executions when a process fails infinitely often. Two variations of this model are considered: with independent process failures, and with simultaneous (i.e., system-wide) process failures. We prove several fundamental results: (i) any synchronization primitive at level 2 in the conventional consensus hierarchy remains at level 2 in the RC hierarchy if failures are simultaneous; (ii) any commutative or overwriting primitive (including Test-And-Set, Fetch-And-Add, and Fetch-And-Store) at level 2 in the conventional consensus hierarchy drops to level 1 in the RC hierarchy if failures are independent, unless the number of such failures is bounded; and (iii) there exists a primitive at level 2 in the conventional consensus hierarchy that remains at level 2 in the RC hierarchy. Result (iii) implies that result (ii) cannot be generalized to all primitives at level 2 in the conventional consensus hierarchy. To our knowledge, this collection of results exhibits the first separation between the simultaneous and independent crash-recovery failure models with respect to the computability of consensus.
Wojciech M. Golab
SPAA1
2020 A Closer Look at Quantum Distributed Consensus
abstract
In a PODC 2008 paper, Helm proposed a protocol for solving distributed consensus using quantum techniques, and without exchanging messages in the classical sense. In this protocol, entangled qubits are distributed to the participants at initialization. Each participant then measures its qubit, and outputs a binary value determined by the outcome of the binary measurement. Since Helm's protocol does not provide the essential properties of consensus (agreement and validity) deterministically, we pose the following question: does this quantum protocol offer any advantage at all over classical protocols that provide similar non-deterministic guarantees? We answer this question in the negative by proving an inherent trade-off between the probability of achieving agreement and the probability of achieving validity in the absence of communication. Our result applies to both classical and quantum protocols.
Wojciech M. Golab, Hao Tan 0006
SPAA1
2020 Benchmarking Recoverable Mutex Locks
abstract
Golab and Ramaraju recently formalized the Recoverable Mutual Exclusion (RME) problem -- a fault-tolerant generalization of Dijkstra's mutual exclusion problem. Several solutions to the RME problem have been proposed since its introduction, and the hardware required to evaluate their performance became available recently following Intel's public launch of Optane Data Center Persistent Memory. In this paper, we present the first experimental evaluation of RME algorithms using an Optane-equipped multiprocessor, with a focus on efficient queue locks. Specifically, we compare Golab and Hendler's recoverable queue lock against Jayanti, Jayanti, and Joshi's, and show that the former is up to 2x faster. Furthermore, we measure the performance penalty of Optane-based RME locks versus DRAM-based conventional locks by comparing the two recoverable locks against an implementation of Mellor-Crummey and Scott's queue lock, and observe that the latter is several-fold faster.
Jeffrey Xiao, Wojciech M. Golab
SPAA3
2019 Toward Linearizability Testing for Multi-Word Persistent Synchronization Primitives
abstract
Persistent memory makes it possible to recover in-memory data structures following a failure instead of rebuilding them from state saved in slow secondary storage. Implementing such recoverable data structures correctly is challenging as their underlying algorithms must deal with both parallelism and failures, which makes them especially susceptible to programming errors. Traditional proofs of correctness should therefore be combined with other methods, such as model checking or software testing, to minimize the likelihood of uncaught defects. This research focuses specifically on the algorithmic principles of software testing, particularly linearizability analysis, for multi-word persistent synchronization primitives such as conditional swap operations. We describe an efficient decision procedure for linearizability in this context, and discuss its practical applications in detecting previously-unknown bugs in implementations of multi-word persistent primitives.
Diego Cepeda, Sakib Chowdhury, Raphael Lopez, Wojciech M. Golab
OPODIS6
2019 Tutorial: Specifying, Implementing, and Verifying Algorithms for Persistent Memory
abstract
High-density byte-addressable non-volatile memory became a reality earlier this year when Intel launched the long-awaited Optane persistent memory module. This tutorial is intended for researchers interested in using persistent memory to construct fault-tolerant data structures that can maintain state consistently across power outages and system crashes without relying on conventional secondary storage. A number of practical and theoretical topics will be covered including hardware purchasing considerations, operating system and programming language support for persistent memory, definitions of correctness properties for fault-tolerant data structures, techniques for implementing fault-tolerant concurrency control and memory management, as well as verification of correctness.
Diego Cepeda, Sakib Chowdhury, Wojciech M. Golab
PODC3
2019 The Recoverable Consensus Hierarchy
abstract
Herlihy's consensus hierarchy ranks the power of various synchronization primitives for solving consensus in a model where asynchronous processes communicate through shared memory, and may fail by halting. This paper revisits the consensus hierarchy in a model with crash-recovery failures, where the specification of consensus, called recoverable consensus in this paper, is weakened by allowing non-terminating executions when a process fails infinitely often. Two variations of this model are considered: with independent process failures, and with simultaneous (i.e., system-wide) process failures. We prove two fundamental results: (i) Test-And-Set is at level 2 of the recoverable consensus hierarchy if failures are simultaneous, and similarly for any primitive at level 2 of the traditional consensus hierarchy; and (ii) Test-And-Set drops to level 1 of the hierarchy if failures are independent, unless the number of such failures is bounded. To our knowledge, this is the first separation between the simultaneous and independent crash-recovery failure models with respect to the computability of consensus.
Wojciech M. Golab
PODC1
2019 Dyn-YCSB: Benchmarking Adaptive Frameworks
abstract
We demonstrate Dyn-YCSB, a tool that builds upon YCSB (Yahoo Cloud Serving Benchmark suite) to assist users in simulating dynamic variations in workloads. Dyn-YCSB automatically varies the parameters in YCSB workloads over time according to user-specified time series functions, without requiring users to manually change the workload configuration in individual nodes each time the workload parameters needs to be modified. The dynamic workload variations simulated with Dyn-YCSB can be used to evaluate the adaptability of such frameworks to changing workload characteristics. We demonstrate the ability of Dyn-YCSB to evaluate the adaptability of OptCon, an automated framework, that tunes the consistency settings of Cassandra with respect to the latency and staleness thresholds in an SLA.
Subhajit Sidhanta, Supratik Mukhopadhyay, Wojciech M. Golab
SERVICES3
2019 Recoverable mutual exclusion
Wojciech M. Golab, Aditya Ramaraju
Distributed Comput.1
2019 Ocean Vista: Gossip-Based Visibility Control for Speedy Geo-Distributed Transactions
abstract
Providing ACID transactions under conflicts across globally distributed data is the Everest of transaction processing protocols. Transaction processing in this scenario is particularly costly due to the high latency of cross-continent network links, which inflates concurrency control and data replication overheads. To mitigate the problem, we introduce Ocean Vista - a novel distributed protocol that guarantees strict serializability . We observe that concurrency control and replication address different aspects of resolving the visibility of transactions, and we address both concerns using a multi-version protocol that tracks visibility using version watermarks and arrives at correct visibility decisions using efficient gossip. Gossiping the watermarks enables asynchronous transaction processing and acknowledging transaction visibility in batches in the concurrency control and replication protocols, which improves efficiency under high cross-datacenter network delays. In particular, Ocean Vista can process conflicting transactions in parallel, and supports efficient write-quorum / read-one access using one round trip in the common case. We demonstrate experimentally in a multi-data-center cloud environment that our design outperforms a leading distributed transaction processing engine (TAPIR) more than 10-fold in terms of peak throughput, albeit at the cost of additional latency for gossip. The latency penalty is generally bounded by one wide area network (WAN) round trip time (RTT), and in the best case (i.e., under light load) our system nearly breaks even with TAPIR by committing transactions in around one WAN RTT.
Hua Fan 0002, Wojciech M. Golab
Proc. VLDB Endow.2
2018 Scalable Transaction Processing Using Functors
abstract
Distributed transactions, which access data items at multiple sites atomically, face well-known scalability challenges. To avoid the high overhead, in prior work Fan et al. proposed Epoch-based Concurrency Control (ECC), which makes transactions visible at epoch boundaries, and presented a system that supports high performance read-only and write-only transactions. However, this idea has a clear difficulty to overcome: the common case of a single transaction that does both reading and writing. This paper proposes ALOHA-DB, a scalable distributed transaction processing system. ALOHA-DB uses a novel paradigm of serializable transaction processing using functors, which conceptually resemble futures in modern programming languages. A functor is a placeholder for the value of a key, which can be computed asynchronously in the future in parallel with other functor computations of the same or other transactions. With multi-versioning in ECC, the functor computations only rely on accessing historical versions, and so the traditional locking mechanism is not needed for concurrency control. Functors elevate ECC to a new level: supporting serializable distributed read-write transactions. This combination of techniques never aborts transactions due to read-write or write-write conflicts, but allows transactions to fail due to logic errors or constraint violations. We used functor-enabled ECC to implement ALOHA-DB and evaluated it using TPC-C and YCSB read-write distributed transactions. Experimental results demonstrate that our system's performance on the TPC-C benchmark is nearly 2 million transactions per second over 20 eight-core virtual machines, which outperforms Calvin, a state-of-the-art transaction processing and replication layer, by one to two orders of magnitude.
Hua Fan 0002, Wojciech M. Golab
ICDCS2
2018 Recoverable Mutual Exclusion Under System-Wide Failures
abstract
Recoverable mutual exclusion (RME) is a variation on the classic mutual exclusion (ME) problem that allows processes to crash and recover. The time complexity of RME algorithms is quantified in the same way as for ME, namely by counting remote memory references -- expensive memory operations that traverse the processor-to-memory interconnect. Prior work has established that the RMR complexity of the RME problem for n processes is Θ(log n) for the class of algorithms that use read/write registers and single-word comparison primitives such as Compare-And-Swap (Golab and Ramaraju 2016), O(log n / log log n) for the class of algorithms that use read/write registers and additional single-word read-modify-primitives such as Fetch-And-Store (Golab and Hendler 2017), and Θ(1) for the class of algorithms that use read/write registers and specialized double-word read-modify-write primitives (Golab and Hendler 2017). These complexity bounds hold in a model of computation where processes may fail independently, and where a process that fails while accessing the mutex is required to recover eventually. This body of work leaves open two important questions: (i) what is the tight bound on the RMR complexity of RME for the class of algorithms that use read/write registers and commonly supported single-word read-modify-primitives; and (ii) how is the RMR complexity of RME affected by variations in the failure model? This paper answers both questions partially by showing that RME can be solved using O(1) RMRs per passage in the worst case in a model where failures are system-wide (i.e., all processes crash simultaneously), and processes receive additional information from the environment regarding the occurrence of the failure. The upper bound algorithm we present relies crucially on a novel RMR-efficient barrier that processes use to synchronize recovery actions after each failure. The barrier uses read/write registers and single-word Compare-And-Swap only. Additionally, we present a transformation that can add properties such as critical section re-entry and a strong notion of starvation freedom to any RME algorithm while preserving its asymptotic RMR complexity.
Wojciech M. Golab, Danny Hendler
PODC1
2018 Analyzing linearizability violations in the presence of read-modify-write operations
Hua Fan 0002, Wojciech M. Golab
Inf. Process. Lett.2
2018 Computing k-Atomicity in Polynomial Time
abstract
The $k$-atomicity property can be used to describe the consistency of data operations in large distributed storage systems. The weak consistency guarantees offered by such systems are seen as a necessary compromise in view of Brewer's CAP principle. The $k$-atomicity property requires that every read operation obtains a value that is at most $k$ updates (writes) old and becomes a useful way to quantify weak consistency if $k$ is treated as a variable that can be computed from a history of operations. Specifically, the value of $k$ quantifies how far the history deviates from the atomicity (linearizability) property for read/write registers. We address the problem of computing $k$ indirectly by solving the $k$-atomicity verification problem ($k$-AV): given a history of read/write operations and a positive integer $k$, decide whether the history is $k$-atomic. Gibbons and Korach showed that in general this problem is NP-complete when $k=1$ and hence not solvable in polynomial time unless $P = NP$. In this paper we present two algorithms that solve the $k$-AV problem for any $k \geq 3$ in special cases. Similarly to known solutions for $k = 1$ and $k = 2$, both algorithms assume that all the values written to a given object are distinct. The first algorithm places an additional restriction on the structure of the input history and solves $k$-AV in $O(n^2 + n \cdot k \log k)$ time, where $n$ is the number of operations in the history. The second algorithm does not place any additional restrictions on the input but is efficient only when $k$ is small and when concurrency among write operations is limited. Its time complexity is $O(n^2)$ if both $k$ and our particular measure of write concurrency are bounded by constants.
Wojciech M. Golab, Xiaozhou Li 0001, Alejandro López-Ortiz, Naomi Nishimura
SIAM J. Comput.1
2017 Efficient incremental data analytics with apache spark
abstract
As smart electricity meters are becoming more popular and starting to replace conventional meters worldwide, new area of research for meter data analytics has emerged. Wide spectrum of computations in this context has been applied, ranging from computationally inexpensive tasks such as calculating monthly bills and peak usage, to elaborate computations to provide energy saving feedback to consumers in order to reduce peak energy demand. Examples include model building approaches for usage predictions and recommendations. Although research efforts in this field are progressing, majority of research in this domain still has overlooked the incremental aspects of energy data analytics, or in best cases, researches have not been able to properly utilize the incremental nature of the energy data. We have noticed that incremental approaches can significantly improve performance of smart meter analytics. For example, per-hour readings of a smart meter can efficiently become integrated with the previous readings and result in an incremental re-computation of a particular smart meter task. In this paper, we introduce UW Incremental Spark Analytics (UWISA), our incremental smart meter data platform, which applies efficient incremental techniques for calculating “energy-temperature” model (also called three-line model) [9]. Our platform can achieve better multi-core scalability and speedup of 4.5× (on average) compared to non-incremental implementation and speedup of higher than 2× when compared to previous incremental research for smart meter datasets up to tens of GBs. We also investigate the reasons behind better performance of incremental method when compared to the non-incremental and Spark Streaming approaches.
Sina Gholamian, Wojciech M. Golab, Paul A. S. Ward
IEEE BigData2
2017 ALOHA-KV: high performance read-only and write-only distributed transactions
abstract
There is a trend in recent database research to pursue coordination avoidance and weaker transaction isolation under a long-standing assumption: concurrent serializable transactions under read-write or write-write conflicts require costly synchronization, and thus may incur a steep price in terms of performance. In particular, distributed transactions, which access multiple data items atomically, are considered inherently costly. They require concurrency control for transaction isolation since both read-write and write-write conflicts are possible, and they rely on distributed commitment protocols to ensure atomicity in the presence of failures. This paper presents serializable read-only and write-only distributed transactions as a counterexample to show that concurrent transactions can be processed in parallel with low-overhead despite conflicts.
Hua Fan 0002, Wojciech M. Golab, Charles B. Morrey III
SoCC2
2017 Brief Announcement: A Probabilistic Performance Model and Tuning Framework for Eventually Consistent Distributed Storage Systems
abstract
Replication protocols in distributed storage systems are fundamentally constrained by the finite propagation speed of information, which necessitates trade-offs among performance metrics even in the absence of failures. We make two contributions toward a clearer understanding of such trade-offs. First, we introduce a probabilistic model of eventual consistency that captures precisely the relationship between the workload, the network latency, and the consistency observed by clients. Second, we propose a technique for adaptive tuning of the consistency-latency trade-off that is based partly on measurement and partly on mathematical modeling. Experiments demonstrate that our probabilistic model predicts the behavior of a practical storage system accurately for low levels of throughput, and that our tuning framework provides superior convergence compared to a state-of-the-art solution.
Shankha Subhra Chatterjee, Wojciech M. Golab
PODC2
2017 Recoverable Mutual Exclusion in Sub-logarithmic Time
abstract
Recoverable mutual exclusion (RME) is a variation on the classic mutual exclusion (ME) problem that allows processes to crash and recover. The time complexity of RME algorithms is quantified in the same way as for ME, namely by counting remote memory references -- expensive memory operations that traverse the processor-to-memory interconnect. Prior work on the RME problem established an upper bound of O(log N) RMRs in an asynchronous shared memory model with N processes that communicate using atomic read and write operations, prompting the question whether sub-logarithmic RMR complexity is attainable using common read-modify-write primitives. We answer this question positively in the cache-coherent model by presenting an RME algorithm that incurs O(log N / log log N) RMRs and uses read, write, Fetch-And-Store, and Compare-And-Swap instructions. We also present an O(1) RMRs algorithm that relies on double-word Compare-And-Swap and a double-word variation of Fetch-And-Store. Both algorithms are inspired by Mellor-Crummey and Scott's queue lock.
Wojciech M. Golab, Danny Hendler
PODC1
2017 Self-tuning Eventually-Consistent Data Stores
Shankha Subhra Chatterjee, Wojciech M. Golab
SSS2
2017 Adaptable SLA-Aware Consistency Tuning for Quorum-Replicated Datastores
abstract
Users of distributed datastores that employ quorum-based replication are burdened with the choice of a suitable client-centric consistency setting for each storage operation. The above matching choice is difficult to reason about as it requires deliberating about the tradeoff between the latency and staleness, i.e., how stale (old) the result is. The latency and staleness for a given operation depend on the client-centric consistency setting applied, as well as dynamic parameters such as the current workload and network condition. We present OptCon, a machine learning-based predictive framework, that can automate the choice of client-centric consistency setting under user-specified latency and staleness thresholds given in the service level agreement (SLA). Under a given SLA, OptCon predicts a client-centric consistency setting that is matching, i.e., it is weak enough to satisfy the latency threshold, while being strong enough to satisfy the staleness threshold. While manually tuned consistency settings remain fixed unless explicitly reconfigured, OptCon tunes consistency settings on a per-operation basis with respect to changing workload and network state. Using decision tree learning, OptCon yields 0.14 cross validation error in predicting matching consistency settings under latency and staleness thresholds given in the SLA. We demonstrate experimentally that OptCon is at least as effective as any manually chosen consistency settings in adapting to the SLA thresholds for different use cases. We also demonstrate that OptCon adapts to variations in workload, whereas a given manually chosen fixed consistency setting satisfies the SLA only for a characteristic workload.
Subhajit Sidhanta, Wojciech M. Golab, Supratik Mukhopadhyay, Saikat Basu
IEEE Trans. Big Data2
2017 Smart Meter Data Analytics: Systems, Algorithms, and Benchmarking
abstract
Smart electricity meters have been replacing conventional meters worldwide, enabling automated collection of fine-grained (e.g., every 15 minutes or hourly) consumption data. A variety of smart meter analytics algorithms and applications have been proposed, mainly in the smart grid literature. However, the focus has been on what can be done with the data rather than how to do it efficiently. In this article, we examine smart meter analytics from a software performance perspective. First, we design a performance benchmark that includes common smart meter analytics tasks. These include offline feature extraction and model building as well as a framework for online anomaly detection that we propose. Second, since obtaining real smart meter data is difficult due to privacy issues, we present an algorithm for generating large realistic datasets from a small seed of real data. Third, we implement the proposed benchmark using five representative platforms: a traditional numeric computing platform (Matlab), a relational DBMS with a built-in machine learning toolkit (PostgreSQL/MADlib), a main-memory column store (“System C”), and two distributed data processing platforms (Hive and Spark/Spark Streaming). We compare the five platforms in terms of application development effort and performance on a multicore machine as well as a cluster of 16 commodity servers.
Xiufeng Liu 0001, Lukasz Golab, Wojciech M. Golab, Ihab F. Ilyas, Shichao Jin
ACM Trans. Database Syst.3
2016 OptEx: A Deadline-Aware Cost Optimization Model for Spark
abstract
We present OptEx, a closed-form model of job execution on Apache Spark, a popular parallel processing engine. To the best of our knowledge, OptEx is the first work that analytically models job completion time on Spark. The model can be used to estimate the completion time of a given Spark job on a cloud, with respect to the size of the input dataset, the number of iterations, the number of nodes comprising the underlying cluster. Experimental results demonstrate that OptEx yields a mean relative error of 6% in estimating the job completion time. Furthermore, the model can be applied for estimating the cost optimal cluster composition for running a given Spark job on a cloud under a completion deadline specified in the SLO (i.e.,Service Level Objective). We show experimentally that OptEx is able to correctly estimate the cost optimal cluster composition for running a given Spark job under an SLO deadline with an accuracy of 98%.
Subhajit Sidhanta, Wojciech M. Golab, Supratik Mukhopadhyay
CCGrid2
2016 OptCon: An Adaptable SLA-Aware Consistency Tuning Framework for Quorum-Based Stores
abstract
Users of distributed datastores that employquorum-based replication are burdened with the choice of asuitable client-centric consistency setting for each storage operation. The above matching choice is difficult to reason about asit requires deliberating about the tradeoff between the latencyand staleness, i.e., how stale (old) the result is. The latencyand staleness for a given operation depend on the client-centricconsistency setting applied, as well as dynamic parameters such asthe current workload and network condition. We present OptCon, a novel machine learning-based predictive framework, that canautomate the choice of client-centric consistency setting underuser-specified latency and staleness thresholds given in the servicelevel agreement (SLA). Under a given SLA, OptCon predictsa client-centric consistency setting that is matching, i.e., it isweak enough to satisfy the latency threshold, while being strongenough to satisfy the staleness threshold. While manually tunedconsistency settings remain fixed unless explicitly reconfigured, OptCon tunes consistency settings on a per-operation basis withrespect to changing workload and network state. Using decisiontree learning, OptCon yields 0.14 cross validation error in predictingmatching consistency settings under latency and stalenessthresholds given in the SLA. We demonstrate experimentally thatOptCon is at least as effective as any manually chosen consistencysettings in adapting to the SLA thresholds for different usecases. We also demonstrate that OptCon adapts to variationsin workload, whereas a given manually chosen fixed consistencysetting satisfies the SLA only for a characteristic workload.
Subhajit Sidhanta, Wojciech M. Golab, Supratik Mukhopadhyay, Saikat Basu
CCGrid2
2016 WatCA: The Waterloo consistency analyzer
abstract
Today's online applications depend on fast storage and retrieval of up-to-date data at web scale. To meet this growing demand, the designers of distributed storage systems have devised a rich variety of data replication protocols, offering different trade-offs between consistency, latency, and availability. Understanding the sweet spot, and testing whether a system delivers a particular level of consistency, are challenging problems as consistency itself is difficult to reason about. This demo paper describes an interactive software tool for measuring and visualizing the consistency actually observed by client applications accessing a key-value storage system in real time. The tool can be used to evaluate performance trade-offs in a system with tunable consistency, or to verify the correctness of a storage system that guarantees certain forms of so-called “strong consistency”.
Hua Fan 0002, Shankha Subhra Chatterjee, Wojciech M. Golab
ICDE3
2016 Recoverable Mutual Exclusion: [Extended Abstract]
abstract
Mutex locks have traditionally been the most common mechanism for protecting shared data structures in parallel programs. However, the robustness of such locks against process failures has not been studied thoroughly. Most (user-level) mutex algorithms are designed around the assumption that processes are reliable, meaning that a process may not fail while executing the lock acquisition and release code, or while inside the critical section.
Wojciech M. Golab, Aditya Ramaraju
PODC1
2015 Fine-tuning the consistency-latency trade-off in quorum-replicated distributed storage systems
abstract
NoSQL storage systems are used extensively by web applications and provide an attractive alternative to conventional databases when the need for scalability outweighs the need for transactions. Several of these systems, notably Amazon's Dynamo and its open-source derivatives, provide quorum-based replication and present the application developer with a choice of multiple client-side "consistency levels" that determine the number of replicas accessed by reads and writes. This setting, in turn, affects both the latency and the consistency observed by the client application. Since using a fixed combination of read and write consistency levels for a given application provides only a limited number of discrete options for tuning the consistency-latency trade-off, we investigate techniques that allow more fine-grained tuning as may be required to support consistency guarantees through service level agreements (SLAs). We consider two such techniques, a novel technique that assigns the consistency level on a peroperation basis by choosing randomly between two options (e.g., weak vs. strong consistency) with a tunable probability, and a known technique that uses weak consistency and injects delays into storage operations artificially. We compare and contrast these two techniques experimentally against each other and against combinations of fixed consistency levels using Apache Cassandra deployed in Amazon's EC2 environment.
Marlon McKenzie, Hua Fan 0002, Wojciech M. Golab
IEEE BigData3
2015 Benchmarking Smart Meter Data Analytics
abstract
Smart electricity meters have been replacing conventional meters worldwide, enabling automated collection of fine-grained (every 15 minutes or hourly) consumption data. A variety of smart meter analytics algorithms and applications have been proposed, mainly in the smart grid literature, but the focus thus far has been on what can be done with the data rather than how to do it efficiently. In this paper, we examine smart meter analytics from a software per-formance perspective. First, we propose a performance benchmark that includes common data analysis tasks on smart meter data. Sec-ond, since obtaining large amounts of smart meter data is diffi-cult due to privacy issues, we present an algorithm for generat-ing large realistic data sets from a small seed of real data. Third, we implement the proposed benchmark using five representative platforms: a traditional numeric computing platform (Matlab), a relational DBMS with a built-in machine learning toolkit (Post-greSQL/MADLib), a main-memory column store (“System C”), and two distributed data processing platforms (Hive and Spark). We compare the five platforms in terms of application development effort and performance on a multi-core machine as well as a cluster of 16 commodity servers. We have made the proposed benchmark and data generator freely available online. 1.
Xiufeng Liu 0001, Lukasz Golab, Wojciech M. Golab, Ihab F. Ilyas
EDBT3
2015 Robust Shared Objects for Non-Volatile Main Memory
abstract
Research in concurrent in-memory data structures has focused almost exclusively on models where processes are either reliable, or may fail by crashing permanently. The case where processes may recover from failures has received little attention because recovery from conventional volatile memory is impossible in the event of a system crash, during which both the state of main memory and the private states of processes are lost. Future hardware architectures are likely to include various forms of non-volatile random access memory (NVRAM), creating new opportunities to design robust main memory data structures that can recover from system crashes. In this paper we advance the theoretical foundations of such data structures in two ways. First, we review several known variations of Herlihy and Wing's linearizability property that were proposed in the context of message passing systems but also apply in our NVRAM-based model, we discuss the limitations of these properties with respect to our specific goals, and we propose an alternative correctness condition called recoverable linearizability. Second, we discuss techniques for implementing shared objects that satisfy such properties with a focus on wait-free implementations. Specifically, we demonstrate how to achieve different variations of linearizability in our model by transforming two classic wait-free constructions.
Ryan Berryhill, Wojciech M. Golab, Mahesh Tripunitara
OPODIS2
2015 Computing Weak Consistency in Polynomial Time: [Extended Abstract]
abstract
The k-atomicity property can be used to describe the consistency of data operations in large distributed storage systems. The weak consistency guarantees offered by such systems are seen as a necessary compromise in view of Brewer's CAP principle. The k-atomicity property requires that every read operation obtains a value that is at most k updates (writes) old, and becomes a useful way to quantify weak consistency if k is treated as a variable that can be computed from a history of operations. Specifically, the value of k quantifies how far the history deviates from Lamport's atomicity property for read/write registers. We address the problem of computing k indirectly by solving the k-atomicity verification problem (k-AV): given a history of read/write operations and a positive integer k, decide whether the history is k-atomic. Gibbons and Korach showed that in general this problem is NP-complete when k = 1, and hence not solvable in polynomial time unless P = NP. In this paper we present two algorithms that solve the k-AV problem for any k >= 2 in special cases. Similarly to known solutions for k = 1 and k = 2, both algorithms assume that all the values written to a given object are distinct. The first algorithm places an additional restriction on the structure of the input history and solves k-AV in O(n^2 + n (k log k) time. The second algorithm does not place any additional restrictions on the input but is efficient only when k is small and when concurrency among write operations is limited. Its time complexity is O(n2) if both k and our particular measure of write concurrency are bounded by constants.
Wojciech M. Golab, Xiaozhou Li 0001, Alejandro López-Ortiz, Naomi Nishimura
PODC1
2015 Understanding the Causes of Consistency Anomalies in Apache Cassandra
abstract
A recent paper on benchmarking eventual consistency showed that when a constant workload is applied against Cassandra, the staleness of values returned by read operations exhibits interesting but unexplained variations when plotted against time. In this paper we reproduce this phenomenon and investigate in greater depth the low-level mechanisms that give rise to stale reads. We show that the staleness spikes exhibited by Cassandra are strongly correlated with garbage collection, particularly the "stop-the-world" phase which pauses all application threads in a Java virtual machine. We show experimentally that the staleness spikes can be virtually eliminated by delaying read operations artificially at servers immediately after a garbage collection pause. In our experiments this yields more than a 98% reduction in the number of consistency anomalies that exceed 5ms, and has negligible impact on throughput and latency.
Hua Fan 0002, Aditya Ramaraju, Marlon McKenzie, Wojciech M. Golab, Bernard Wong 0001
Proc. VLDB Endow.4
2014 Client-Centric Benchmarking of Eventual Consistency for Cloud Storage Systems
abstract
Eventually-consistent key-value storage systems sacrifice the ACID semantics of conventional databases to achieve superior latency and availability. However, this means that client applications, and hence end-users, can be exposed to stale data. The degree of staleness observed depends on various tuning knobs set by application developers (customers of key-value stores) and system administrators (providers of key-value stores). Both parties must be cognizant of how these tuning knobs affect the consistency observed by client applications in the interest of both providing the best end-user experience and maximizing revenues for storage providers. Quantifying consistency in a meaningful way is a critical step toward both understanding what clients actually observe, and supporting consistency-aware service level agreements (SLAs) in next generation storage systems. This paper proposes a novel consistency metric called Gamma that captures client-observed consistency. This metric provides quantitative answers to questions regarding observed consistency anomalies, such as how often they occur and how bad they are when they do occur. We argue that Gamma is more useful and accurate than existing metrics. We also apply Gamma to benchmark the popular Cassandra key-value store. Our experiments demonstrate that Gamma is sensitive to both the workload and client-level tuning knobs, and is preferable to existing techniques which focus on worst-case behavior.
Wojciech M. Golab, Muntasir Raihan Rahman, Alvin AuYoung, Kimberly Keeton, Indranil Gupta
ICDCS1
2014 Making objects writable
abstract
We devise a technique for augmenting shared objects in the standard n-process shared memory model with a linearizable Write{} operation, using bounded space and optimal worst-case step complexity. We provide a transformation of any shared object SW supporting only sequential Write{} operations into an object $W$ that supports concurrent Write{} operations. This transformation requires O(n2) SW objects and O(n2) O(log n)-bit registers, and each method (including Write{}) has, up to a constant additive term, the same time complexity as the corresponding method on object $SW$. Our implementation is deterministic, wait-free, and uses only shared registers (supporting atomic read and write operations). To the best of our knowledge, similarly efficient general constructions are not known even if stronger primitives such as CAS or LL/SC are available.
Zahra Aghazadeh, Wojciech M. Golab, Philipp Woelfel
PODC2
2014 Making Sense of Relativistic Distributed Systems
Seth Gilbert, Wojciech M. Golab
DISC2
2013 Client-centric benchmarking of eventual consistency for cloud storage systems
abstract
Eventually consistent storage systems give up the ACID semantics of conventional databases in order to gain better scalability, higher availability, and lower latency. A side-effect of this design decision is that application developers must deal with stale or out of order data. As a result, substantial intellectual effort has been devoted to studying the behavior of eventually consistent systems, in particular finding quantitative answers to the questions "how eventual" and "how consistent"?
Wojciech M. Golab, Muntasir Raihan Rahman, Alvin AuYoung, Kimberly Keeton, Jay J. Wylie, Indranil Gupta
SoCC1
2013 On the k-Atomicity-Verification Problem
abstract
Modern Internet-scale storage systems often provide weak consistency in exchange for better performance and resilience. An important weak consistency property is k-atomicity, which bounds the staleness of values returned by read operations. The k-atomicity-verification problem (or k-AV for short) is the problem of deciding whether a given history of operations is k-atomic. The 1-AV problem is equivalent to verifying atomicity/linearizability, a well-known and solved problem. However, for k 2, no polynomial-time k-AV algorithm is known. This paper makes the following contributions towards solving the k-AV problem. First, we present a simple 2- AV algorithm called LBT, which is likely to be efficient (quasilinear) for histories that arise in practice, although it is less efficient (quadratic) in the worst case. Second, we present a more involved 2-AV algorithm called FZF, which runs efficiently (quasilinear) even in the worst case. To our knowledge, these are the first algorithms that solve the 2-AV problem fully. Third, we show that the weighted k-AV problem, a natural extension of the k-AV problem, is NP-complete.
Wojciech M. Golab, Jeremy Hurwitz, Xiaozhou Li 0001
ICDCS1
2013 Brief announcement: resettable objects and efficient memory reclamation for concurrent algorithms
abstract
We present a new technique for reclaiming memory in concurrent shared memory algorithms with n asynchronous processes. Our methodology can be applied in the same settings as hazard pointers [10], but provides better worst-case guarantees: For the same tasks for which hazard pointers have expected constant amortized complexity, our technique guarantees constant time in the worst-case.
Zahra Aghazadeh, Wojciech M. Golab, Philipp Woelfel
PODC2
2012 RMR-efficient implementations of comparison primitives using read and write operations
Wojciech M. Golab, Vassos Hadzilacos, Danny Hendler, Philipp Woelfel
Distributed Comput.1
2012 Minuet: A Scalable Distributed Multiversion B-Tree
abstract
Data management systems have traditionally been designed to support either long-running analytics queries or short-lived transactions, but an increasing number of applications need both. For example, online games, socio-mobile apps, and e-commerce sites need to not only maintain operational state, but also analyze that data quickly to make predictions and recommendations that improve user experience. In this paper, we present Minuet, a distributed, main-memory B-tree that supports both transactions and copy-on-write snapshots for in-situ analytics. Minuet uses main-memory storage to enable low-latency transactional operations as well as analytics queries without compromising transaction performance. In addition to supporting read-only analytics queries on snapshots, Minuet supports writable clones, so that users can create branching versions of the data. This feature can be quite useful, e.g. to support complex "what-if" analysis or to facilitate wide-area replication. Our experiments show that Minuet outperforms a commercial main-memory database in many ways. It scales to hundreds of cores and TBs of memory, and can process hundreds of thousands of B-tree operations per second while executing long-running scans.
Ben Sowell, Wojciech M. Golab, Mehul A. Shah
Proc. VLDB Endow.2
2011 A complexity separation between the cache-coherent and distributed shared memory models
abstract
We consider asynchronous multiprocessor systems where processes communicate by accessing shared memory. Exchange of information among processes in such a multiprocessor necessitates costly memory accesses called remote memory references (RMRs), which generate communication on the interconnect joining processors and main memory. In this paper we compare two popular shared memory architecture models, namely the ca che-coherent (CC) and distributed shared memory (DSM) models, in terms of their power for solving synchronization problems efficiently with respect to RMRs. The particular problem we consider entails one process sending a signal to a subset of other processes. We show that a variant of this problem can be solved very efficiently with respect to RMRs in the CC model, but not so in the DSM model, even when we consider amortized RMR complexity.To our knowledge, this is the first separation in terms of amortized RMR complexity between the CC and DSM models. It is also the first separation in terms of RMR complexity (for asynchronous systems) that does not rely in any way on wait-freedom---the requirement that a process makes progress in a bounded number of its own steps.
Wojciech M. Golab
PODC1
2011 Analyzing consistency properties for fun and profit
abstract
Motivated by the increasing popularity of eventually consistent key-value stores as a commercial service, we address two important problems related to the consistency properties in a history of operations on a read/write register (i.e., the start time, finish time, argument, and response of every operation). First, we consider how to detect a consistency violation as soon as one happens. To this end, we formulate a specification for online verification algorithms, and we present such algorithms for several well-known consistency properties. Second, we consider how to quantify the severity of the violations, if a history is found to contain consistency violations. We investigate two quantities: one is the staleness of the reads, and the other is the commonality of violations. For staleness, we further consider time-based staleness and operation-count-based staleness. We present efficient algorithms that compute these quantities. We believe that addressing these problems helps both key-value store providers and users adopt data consistency as an important aspect of key-value store offerings.
Wojciech M. Golab, Xiaozhou Li 0001, Mehul A. Shah
PODC1
2011 Linearizable implementations do not suffice for randomized distributed computation
abstract
Linearizability is the gold standard among algorithm designers for deducing the correctness of a distributed algorithm using implemented shared objects from the correctness of the corresponding algorithm using atomic versions of the same objects. We show that linearizability does not suffice for this purpose when processes can exploit randomization, and we discuss the existence of alternative correctness conditions. This paper makes the following contributions: 1. Various examples demonstrate that using well-known linearizable implementations of objects (e.g., snapshots) in place of atomic objects can change the probability distribution of the outcomes that the adversary is able to generate. In some cases, an oblivious adversary can create a probability distribution of outcomes for an algorithm with implemented, linearizable objects, that not even a strong adversary can generate for the same algorithm with atomic objects. 2. A new correctness condition for shared object implementations, called strong inearizability, is defined. We prove that a strong adversary (i.e., one that sees the outcome of each coin flip immediately) gains no additional power when atomic objects are replaced by strongly linearizable implementations. In general, no strictly weaker correctness condition suffices to ensure this. We also show that strong linearizability is a local and composable property. 3. In contrast to the situation for the strong adversary, for a natural weaker adversary (one that cannot see a process' coin flip until its next operation on a shared object) we prove that there is no correspondingly general correctness condition. Specifically, any linearizable implementation of counters called terminating. from atomic registers and load-linked/store-conditional objects, that satisfies a natural locality property, necessarily gives the weak adversary more power than it has with atomic counters.
Wojciech M. Golab, Lisa Higham, Philipp Woelfel
STOC1
2010 Brief announcement: locally-accessible implementations for distributed shared memory multiprocessors
abstract
We consider asynchronous multiprocessors that support the distributed shared memory (DSM) model. Algorithms for such multiprocessors exploit the ability to co-locate shared objects with particular processes in order to reduce the cost of accessing shared memory. When a shared object fits inside a single memory word and operations on it are supported directly through machine instructions, it can be made local to any process simply by fixing its physical address. We show that even if the shared object is not supported in hardware directly, it can always be simulated using a software implementation that behaves as though it is local to some designated process. That is, operations applied by the designated process on the implemented object access only local base objects, which is non-trivial when processes synchronize by busy-waiting. We also discuss time complexity bounds for such locally-accessible implementations.
Wojciech M. Golab
PODC1
2010 Closing the complexity gap between FCFS mutual exclusion and mutual exclusion
Robert Danek, Wojciech M. Golab
Distributed Comput.2
2010 An O(1) RMRs Leader Election Algorithm
abstract
The leader election problem is a fundamental coordination problem. We present leader election algorithms for multiprocessor systems where processes communicate by reading and writing shared memory asynchronously and do not fail. In particular, we consider the cache-coherent (CC) and distributed shared memory (DSM) models of such systems. We present leader election algorithms that perform a constant number of remote memory references (RMRs) in the worst case. Our algorithms use splitter-like objects [J. Anderson and M. Moir, Sci. Comput. Programming, 25 (1995), pp. 1–39; H. Attiya and A. Fouren, Theory Comput. Syst., 31 (2001), pp. 642–664] in a novel way, by organizing active processes into teams that share work. As there is an $\Omega(\log n)$ lower bound on the RMR complexity of mutual exclusion for n processes using reads and writes only [H. Attiya, D. Hendler, and W. Woelfel, in Proceedings of the ACM Symposium on Theory of Computing, ACM, New York, 2008, pp. 217–226], our result separates the mutual exclusion and leader election problems in terms of RMR complexity in both the CC and DSM models. Our result also implies that any algorithm using reads, writes, and one-time test-and-set objects can be simulated by an algorithm using reads and writes with only a constant blowup of the RMR complexity; proving this is easy in the CC model but presents subtle challenges in the DSM model, as we explain later. Anderson, Herman, and Kim raise the question of whether conditional primitives such as test-and-set and compare-and-swap can be used, along with reads and writes, to solve mutual exclusion with better worst-case RMR complexity than is possible using reads and writes only [Distributed Computing, 16 (2003), pp. 75–110]. We provide a negative answer to this question in the case of implementing one-time test-and-set.
Wojciech M. Golab, Danny Hendler, Philipp Woelfel
SIAM J. Comput.1
2008 Closing the complexity gap between mutual exclusion and FCFS mutual exclusion
abstract
We consider the worst-case remote memory reference (RMR) complexity of first-come-first-served (FCFS) mutual exclusion (ME) algorithms for N asynchronous reliable processes that communicate only by reading and writing shared memory. We exhibit an upper bound of O(log N) RMRs for FCFS ME, which is tight, improves on prior results, and matches a lower bound for ME (with or without FCFS).
Robert Danek, Wojciech M. Golab
PODC2
2008 Closing the Complexity Gap between FCFS Mutual Exclusion and Mutual Exclusion
Robert Danek, Wojciech M. Golab
DISC2
2008 A practical scalable distributed B-tree
abstract
Internet applications increasingly rely on scalable data structures that must support high throughput and store huge amounts of data. These data structures can be hard to implement efficiently. Recent proposals have overcome this problem by giving up on generality and implementing specialized interfaces and functionality (e.g., Dynamo [4]). We present the design of a more general and flexible solution: a fault-tolerant and scalable distributed B-tree. In addition to the usual B-tree operations, our B-tree provides some important practical features: transactions for atomically executing several operations in one or more B-trees, online migration of B-tree nodes between servers for load-balancing, and dynamic addition and removal of servers for supporting incremental growth of the system. Our design is conceptually simple. Rather than using complex concurrency and locking protocols, we use distributed transactions to make changes to B-tree nodes. We show how to extend the B-tree and keep additional information so that these transactions execute quickly and efficiently. Our design relies on an underlying distributed data sharing service, Sinfonia [1], which provides fault tolerance and a light-weight distributed atomic primitive. We use this primitive to commit our transactions. We implemented our B-tree and show that it performs comparably to an existing open-source B-tree and that it scales to hundreds of machines. We believe that our approach is general and can be used to implement other distributed data structures easily.
Marcos K. Aguilera, Wojciech M. Golab, Mehul A. Shah
Proc. VLDB Endow.2
2007 Constant-RMR implementations of CAS and other synchronization primitives using read and write operations
abstract
We consider asynchronous multiprocessors where processes communicate only by reading or writing shared memory. We show how to implement consensus, all comparison primitives (such as CAS and TAS), and load-linked/store-conditional using only a constant number of remote memory references (RMRs), in both the cache-coherent and the distributed-shared-memory models of such multiprocessors. Our implementations are blocking, rather than wait-free: they ensure progress provided all processes that invoke the implemented primitive are live.
Wojciech M. Golab, Vassos Hadzilacos, Danny Hendler, Philipp Woelfel
PODC1
2007 Admission control in data transfers over lightpaths
abstract
The availability of optical network infrastructure and appropriate user control software has recently made it possible for scientists to establish end-to-end circuits across multiple management domains in support of large data transfers. These high-performance data paths are typically provisioned over 10 Gigabit optical links, and accessed using ethernet encapsulation at Gigabit and 10 Gigabit rates. The resulting mixture of circuit sizes gives rise to resource conflicts whereby requests to allocate bandwidth partitions are blocked despite vast underutilization of the optical link. In an attempt to remedy this problem, we investigate intelligent admission control policies that consider the long-term effects of admission decisions. Using analytic techniques we show that the greedy policy, which accepts requests to allocate bandwidth partitions whenever sufficient bandwidth exists, is suboptimal in a pertinent scenario. We then consider dynamic online computation of the optimal admission control policy and show that the acceptance ratio of requests to establish end-to-end circuits can be improved by up to 19% on a fifteen-node network where the behaviour of each link is governed by a local optimization effort.
Wojciech M. Golab, Raouf Boutaba
IEEE J. Sel. Areas Commun.1
2006 An O(1) RMRs leader election algorithm
abstract
The leader election problem is a fundamental distributed coordination problem. We present leader election algorithms for the cache-coherent (CC) and distributed shared memory (DSM) models using reads and writes only, for which the number of remote memory references (RMRs) is constant in the worst case.The algorithms use splitter-like objects [6, 8] in a novel way for the efficient partitioning of processes into disjoint sets that share work. As there is an Ω(log n/log log n) lower bound on the RMR complexity of mutual exclusion for n processes using reads and writes only [4], our result separates the mutual exclusion and leader election problems in terms of RMR complexity in both the CC and DSM models.Our result also implies that any algorithm using reads, writes and one-time test-and-set objects can be simulated by an algorithm using reads and writes with only a constant blowup of the RMR complexity. Anderson, Herman and Kim raise the question of whether conditional primitives such as test-and-set and compare-and-swap are stronger than read and write for the implementation of local-spin mutual exclusion [3]. We provide a negative answer to this question, at least for one-time test-and-set.
Wojciech M. Golab, Danny Hendler, Philipp Woelfel
PODC1
2003 Grid-Controlled Lightpaths for High Performance Grid Applications
Raouf Boutaba, Wojciech M. Golab, Youssef Iraqi, Tianshu Li, Bill St. Arnaud
J. Grid Comput.2