EDBT 2026 Demo / reviewers in the wild / expert
Dahlia Malkhi
dblp:m/DahliaMalkhi · also Dalia Malki
· DBLP profile ↗
128ranked-venue papers
37as first author
13since 2021 · last 2026
0000-0002-7038-7250ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 57 · 12 first-author · 4 since 2021Security and privacy · 22 · 13 first-author · 2 since 2021Theory of computation · 13 · 4 first-authorComputer networks · 6 · 1 since 2021Databases, data management, data science and information retrieval · 5 · 1 first-author · 1 since 2021Software engineering, systems software and programming languages · 2 · 1 first-authorArtificial intelligence and machine learning · 1 · 1 since 2021Graphics, computer vision, multimedia, augmented reality and games · 1 · 1 since 2021Applied, interdisciplinary, general and emerging computing · 1
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | Inequality in the Age of PseudonymityabstractInequality measures such as the Gini coefficient are used to inform and motivate policymaking, and are increasingly applied to digital platforms. We analyze how measures fare in pseudonymous settings that are common in the digital age. One key challenge of such environments is the ability of actors to create fake identities under fictitious false names, also known as ``Sybils.'' While some actors may do so to preserve their privacy, we show that this can hamper inequality measurements: it is impossible for measures satisfying the literature's canonical set of desired properties to assess the inequality of an economy that may harbor Sybils. We characterize the class of all Sybil-proof measures, and prove that they must satisfy relaxed version of the aforementioned properties. Furthermore, we show that the structure imposed restricts the ability to assess inequality at a fine-grained level. We then apply our results to prove that popular measures are not Sybil-proof, with the famous Gini coefficient being but one example out of many. Finally, we examine dynamics leading to the creation of Sybils in digital and traditional settings. Aviv Yaish, Nir Chemaya, Dahlia Malkhi, Will Cong |
AAAI | 3 |
| 2025 | BFTBrain: Adaptive BFT Consensus with Reinforcement Learning
Chenyuan Wu, Haoyun Qin, Mohammad Javad Amiri, Boon Thau Loo, Dahlia Malkhi, Ryan Marcus |
NSDI | 5 |
| 2025 | Brief Announcement: Carry the Tail in Consensus ProtocolsabstractWe present Carry-the-Tail, the first deterministic atomic broadcast protocol in partial synchrony that, after GST, simultaneously guarantees two desirable properties: (i) a constant fraction of commits are proposed by non-faulty leaders against tail-forking attacks, and (ii) optimal, worst-case quadratic communication under a cascade of faulty leaders. The solution also guarantees linear amortized communication, i.e., the steady-state is linear. Combining these two desirable properties was not simultaneously achieved previously: on one hand, prior atomic broadcast solutions achieve per-view linear word communication complexity. However, they face a significant degradation in throughput under tail-forking attack. On the other hand, existing solutions to tail-forking attacks require either quadratic communication steps or computationally-prohibitive SNARK generation. The key technical contribution is Carry, a practical drop-in mechanism for streamlined protocols in the HotStuff family. Carry guarantees good performance against tail-forking and removes most leader-induced stalls, while retaining linear traffic and protocol simplicity. Carry-the-Tail implements the Carry mechanism on HotStuff-2. Suyash Gupta 0001, Dakai Kang, Dahlia Malkhi, Mohammad Sadoghi |
DISC | 3 |
| 2025 | HotStuff-1: Linear Consensus with One-Phase SpeculationabstractThis paper introduces HotStuff-1, a BFT consensus protocol that improves the latency of HotStuff-1 by two network hops while maintaining linear communication complexity against faults. Furthermore, HotStuff-1 incorporates an incentive-compatible leader rotation design that motivates leaders to propose transactions promptly. HotStuff-1 achieves a reduction of two network hops by speculatively sending clients early finality confirmations, after one phase of the protocol. Introducing speculation into streamlined protocols is challenging because, unlike stable-leader protocols, these protocols cannot stop the consensus and recover from failures. Thus, we identify prefix speculation dilemma in the context of streamlined protocols; HotStuff-1 is the first streamlined protocol to resolve it. HotStuff-1 embodies an additional mechanism, slotting , that thwarts delays caused by (1) rationally-incentivized leaders and (2) malicious leaders inclined to sabotage others' progress. The slotting mechanism allows leaders to dynamically drive as many decisions as allowed by network transmission delays before view timers expire, thus mitigating both threats. Dakai Kang, Suyash Gupta 0001, Dahlia Malkhi, Mohammad Sadoghi |
Proc. ACM Manag. Data | 3 |
| 2024 | BBCA-Chain: Low Latency, High Throughput BFT Consensus on a DAG
Dahlia Malkhi, Chrysoula Stathakopoulou, Maofan Yin |
FC (1) | 1 |
| 2024 | Lumiere: Making Optimal BFT for Partial Synchrony PracticalabstractThe view synchronization problem lies at the heart of many Byzantine Fault Tolerant (BFT) State Machine Replication (SMR) protocols in the partial synchrony model, since these protocols are usually based on views. Liveness is guaranteed if honest processors spend a sufficiently long time in the same view during periods of synchrony, and if the leader of the view is honest. Ensuring that these conditions occur, known as Byzantine View Synchronization (BVS), has turned out to be the performance bottleneck of many BFT SMR protocols. Andy Lewis-Pye, Dahlia Malkhi, Oded Naor, Kartik Nayak |
PODC | 2 |
| 2023 | Towards Practical Sleepy BFTabstractBitcoin's longest-chain protocol pioneered consensus under dynamic participation, also known as sleepy consensus, where nodes do not need to be permanently active. However, existing solutions for sleepy consensus still face two major issues, which we address in this work. First, existing sleepy consensus protocols have high latency (either asymptotically or concretely). We tackle this problem and achieve 4Δ latency (Δ is the bound on network delay) in the best case, which is comparable to classic BFT protocols without dynamic participation support. Second, existing protocols have to assume that the set of corrupt participants remains fixed throughout the lifetime of the protocol due to a problem we call costless simulation. We resolve this problem and support growing participation of corrupt nodes. Our new protocol also offers several other important advantages, including support for arbitrary fluctuation of honest participation as well as an efficient recovery mechanism for new active nodes. Dahlia Malkhi, Atsuki Momose, Ling Ren 0001 |
CCS | 1 |
| 2023 | Block-STM: Scaling Blockchain Execution by Turning Ordering Curse to a Performance BlessingabstractBlock-STM is a parallel execution engine for smart contracts, built around the principles of Software Transactional Memory. Transactions are grouped in blocks, and every execution of the block must yield the same deterministic outcome. Block-STM further enforces that the outcome is consistent with executing transactions according to a preset order, leveraging this order to dynamically detect dependencies and avoid conflicts during speculative transaction execution. At the core of Block-STM is a novel, low-overhead collaborative scheduler of execution and validation tasks. Rati Gelashvili, Alexander Spiegelman, Zhuolun Xiang, George Danezis, Zekun Li 0009, Dahlia Malkhi, Yu Xia 0005, Runtian Zhou |
PPoPP | 6 |
| 2021 | Strengthened Fault Tolerance in Byzantine Fault Tolerant ReplicationabstractByzantine fault tolerant (BFT) state machine replication (SMR) is an important building block for constructing permissioned blockchain systems. In contrast to Nakamoto Consensus where any block obtains higher assurance as buried deeper in the blockchain, in BFT SMR, any committed block is secure has a fixed resilience threshold. In this paper, we investigate strengthened fault tolerance (SFT) in BFT SMR under partial synchrony, which provides stronger resilience guarantees during an optimistic period when the network is synchronous and the number of Byzantine faults is small. Moreover, the committed blocks can tolerate more than one-third (up to two-thirds) corruptions even after the optimistic period. Compared to the prior best solution FBFT which requires quadratic message complexity, our solution maintains the linear message complexity of state-of-the-art BFT SMR protocols and requires only marginal bookkeeping overhead. We implement our solution over the open-source Diem project, and give experimental results that demonstrate its efficiency under real-world scenarios. Zhuolun Xiang, Dahlia Malkhi, Kartik Nayak, Ling Ren 0001 |
ICDCS | 2 |
| 2021 | Twins: BFT Systems Made Robust
Shehar Bano, Alberto Sonnino, Andrey Chursin, Dmitri Perelman, Zekun Li 0009, Avery Ching, Dahlia Malkhi |
OPODIS | 7 |
| 2021 | RainBlock: Faster Transaction Processing in Public Blockchains
Soujanya Ponnapalli, Aashaka Shah, Souvik Banerjee, Dahlia Malkhi, Amy Tai, Vijay Chidambaram, Michael Wei |
USENIX ATC | 4 |
| 2021 | Brief Announcement: Twins - BFT Systems Made RobustabstractTwins is an effective strategy for generating test scenarios with Byzantine [Lamport et al., 1982] nodes in order to find flaws in Byzantine Fault Tolerant (BFT) systems. Twins finds flaws in the design or implementation of BFT protocols that may cause correctness issues. The main idea of Twins is the following: running twin instances of a node that use correct, unmodified code and share the same network identity and credentials allows to emulate most interesting Byzantine behaviors. Because a twin executes normal, unmodified node code, building Twins only requires a thin wrapper over an existing distributed system designed for Byzantine tolerance. To emulate material, interesting scenarios with Byzantine nodes, it instantiates one or more twin copies of the node, giving the twins the same identities and network credentials as the original node. To the rest of the system, the node and all its twins appear indistinguishable from a single node behaving in a "questionable" manner. This approach generates many interesting Byzantine behaviors, including equivocation, double voting, and losing internal state, while forgoing uninteresting behavior scenarios that can be filtered at the transport layer, such as producing semantically invalid messages. Building on configurations with twin nodes, Twins systematically generates scenarios with Byzantine nodes via enumeration over protocol rounds and communication patterns among nodes. Despite this being inherently exponential, one new flaw and several known flaws were materialized by Twins in the arena of BFT consensus protocols. In all cases, protocols break within fewer than a dozen protocol rounds, hence it is realistic for the Twins approach to expose the problems. In two of these cases, it took the community more than a decade to discover protocol flaws that Twins would have surfaced within minutes. Additionally, Twins has been incorporated into the continuous release testing process of a production setting (DiemBFT) in which it can execute 44M Twins-generated scenarios daily. Shehar Bano, Alberto Sonnino, Andrey Chursin, Dmitri Perelman, Zekun Li 0009, Avery Ching, Dahlia Malkhi |
DISC | 7 |
| 2021 | Tech Transfer Stories and Takeaways (Invited Talk)abstractIn this talk, I will share impressions from several industrial research project experiences that reached production and became part of successful products. I will go through four stories of how these systems transpired and their journey to impact. All of the stories are in the distributed computing arena, and more specifically, they revolve around the state-machine-replication paradigm. Yet, I hope that the take-aways from the experience of building foundations for these systems may be of interest and value to everyone, no matter the discipline. Dahlia Malkhi |
DISC | 1 |
| 2020 | Asynchronous Distributed Key Generation for Computationally-Secure Randomness, Consensus, and Threshold SignaturesabstractIn this paper, we present the first Asynchronous Distributed Key Generation (ADKG) algorithm which is also the first distributed key generation algorithm that can generate cryptographic keys with a dual (f,2f+1)-threshold (where f is the number of faulty parties). As a result, using our ADKG we remove the trusted setup assumption that the most scalable consensus algorithms make. In order to create a DKG with a dual (f,2f+1)- threshold we first answer in the affirmative the open question posed by Cachin et al. [7] on how to create an Asynchronous Verifiable Secret Sharing (AVSS) protocol with a reconstruction threshold of f+1 Eleftherios Kokoris-Kogias, Dahlia Malkhi, Alexander Spiegelman |
CCS | 2 |
| 2020 | ACE: Abstract Consensus Encapsulation for Liveness Boosting of State Machine ReplicationabstractWith the emergence of attack-prone cross-organization systems, providing asynchronous state machine replication (SMR) solutions is no longer a theoretical concern. This paper presents ACE, a framework for the design of such fault tolerant systems. Leveraging a known paradigm for randomized consensus solutions, ACE wraps existing practical solutions and real-life systems, boosting their liveness under adversarial conditions and, at the same time, promoting load balancing and fairness. Boosting is achieved without modifying the overall design or the engineering of these solutions. ACE is aimed at boosting the prevailing approach for practical fault tolerance. This approach, often named partial synchrony, is based on a leader-based paradigm: a good leader makes progress and a bad leader does no harm. The partial synchrony approach focuses on safety and forgoes liveness under targeted and dynamic attacks. Specifically, an attacker might block specific leaders, e.g., through a denial of service, to prevent progress. ACE provides boosting by running waves of parallel leaders and selecting a winning leader only retroactively, achieving boosting at a linear communication cost increase. ACE is agnostic to the fault model, inheriting it s failure model from the wrapped solution assumptions. As our evaluation shows, an asynchronous Byzantine fault tolerance (BFT) replication system built with ACE around an existing partially synchronous BFT protocol demonstrates reasonable slow-down compared with the base BFT protocol during faultless synchronous scenarios, yet exhibits significant speedup while the system is under attack. Alexander Spiegelman, Arik Rinberg, Dahlia Malkhi |
OPODIS | 3 |
| 2020 | Sync HotStuff: Simple and Practical Synchronous State Machine ReplicationabstractSynchronous solutions for Byzantine Fault Tolerance (BFT) can tolerate up to minority faults. In this work, we present Sync HotStuff, a surprisingly simple and intuitive synchronous BFT solution that achieves consensus with a latency of 2Δ in the steady state (where Δ is a synchronous message delay upper bound). In addition, Sync HotStuff ensures safety in a weaker synchronous model in which the synchrony assumption does not have to hold for all replicas all the time. Moreover, Sync HotStuff has optimistic responsiveness, i.e., it advances at network speed when less than one-quarter of the replicas are not responding. Borrowing from practical partially synchronous BFT solutions, Sync HotStuff has a two-phase leader-based structure, and has been fully prototyped under the standard synchrony assumption. When tolerating a single fault, Sync HotStuff achieves a throughput of over 280 Kops/sec under typical network performance, which is comparable to the best known partially synchronous solution. Ittai Abraham, Dahlia Malkhi, Kartik Nayak, Ling Ren 0001, Maofan Yin |
SP | 2 |
| 2020 | Introduction to the Special Section on USENIX ATC 2019abstractNo abstract available. Dahlia Malkhi, Dan Tsafrir |
ACM Trans. Storage | 1 |
| 2019 | Efficient Verifiable Secret Sharing with Share Recovery in BFT ProtocolsabstractByzantine fault tolerant state machine replication (SMR) provides powerful integrity guarantees, but fails to provide any privacy guarantee whatsoever. A natural way to add such privacy guarantees is to secret-share state instead of fully replicating it. Such a com- bination would enable simple solutions to difficult problems, such as a fair exchange or a distributed certification authority. However, incorporating secret shared state into traditional Byzantine fault tolerant (BFT) SMR protocols presents unique challenges. BFT protocols often use a network model that has some degree of asynchrony, making verifiable secret sharing (VSS) unsuitable. However, full asynchronous VSS (AVSS) is unnecessary as well since the BFT algorithm provides a broadcast channel. We first present the VSS with share recovery problem, which is the subproblem of AVSS required to incorporate secret shared state into a BFT engine. Then, we provide the first VSS with share recovery solution, KZG-VSSR, in which a failure-free sharing incurs only a constant number of cryptographic operations per replica. Finally, we show how to efficiently integrate any instantiation of VSSR into a BFT replication protocol while incurring only constant overhead. Instantiating VSSR with prior AVSS protocols would require a quadratic communication cost for a single shared value and incur a linear overhead when incorporated into BFT replication. We demonstrate our end-to-end solution via a a private key-value store built using BFT replication and two instantiations of VSSR, KZG-VSSR and Ped-VSSR, and present its evaluation. Soumya Basu 0003, Alin Tomescu, Ittai Abraham, Dahlia Malkhi, Michael K. Reiter, Emin Gün Sirer |
CCS | 4 |
| 2019 | Flexible Byzantine Fault ToleranceabstractThis paper introduces Flexible BFT, a new approach for BFT consensus solution design revolving around two pillars, stronger resilience and diversity. The first pillar, stronger resilience, involves a new fault model called alive-but-corrupt faults. Alive-but-corrupt replicas may arbitrarily deviate from the protocol in an attempt to break safety of the protocol. However, if they cannot break safety, they will not try to prevent liveness of the protocol. Combining alive-but-corrupt faults into the model, Flexible BFT is resilient to higher corruption levels than possible in a pure Byzantine fault model. The second pillar, diversity, designs consensus solutions whose protocol transcript is used to draw different commit decisions under diverse beliefs. With this separation, the same Flexible BFT solution supports synchronous and asynchronous beliefs, as well as varying resilience threshold combinations of Byzantine and alive-but-corrupt faults. At a technical level, Flexible BFT achieves the above results using two new ideas. First, it introduces a synchronous BFT protocol in which only the commit step requires to know the network delay bound and thus replicas execute the protocol without any synchrony assumption. Second, it introduces a notion called Flexible Byzantine Quorums by dissecting the roles of different quorums in existing consensus protocols. Dahlia Malkhi, Kartik Nayak, Ling Ren 0001 |
CCS | 1 |
| 2019 | SBFT: A Scalable and Decentralized Trust InfrastructureabstractSBFT is a state of the art Byzantine fault tolerant state machine replication system that addresses the challenges of scalability, decentralization and global geo-replication. SBFT is optimized for decentralization and is experimentally evaluated on a deployment of more than 200 active replicas withstanding a malicious adversary controlling f=64 replicas. Our experiments show how the different algorithmic ingredients of SBFT contribute to its performance and scalability. The results show that SBFT simultaneously provides almost 2x better throughput and about 1.5x better latency relative to a highly optimized system that implements the PBFT protocol. To achieve this performance improvement, SBFT uses a combination of four ingredients: using collectors and threshold signatures to reduce communication to linear, using an optimistic fast path, reducing client communication and utilizing redundant servers for the fast path. SBFT is the first system to implement a correct dual-mode view change protocol that allows to efficiently run either an optimistic fast path or a fallback slow path without incurring a view change to switch between modes. Guy Golan-Gueta, Ittai Abraham, Shelly Grossman, Dahlia Malkhi, Benny Pinkas, Michael K. Reiter, Dragos-Adrian Seredinschi, Orr Tamir, Alin Tomescu |
DSN | 4 |
| 2019 | FairLedger: A Fair Blockchain Protocol for Financial InstitutionsabstractFinancial institutions are currently looking into technologies for permissioned blockchains. A major effort in this direction is Hyperledger, an open source project hosted by the Linux Foundation and backed by a consortium of over a hundred companies. A key component in permissioned blockchain protocols is a byzantine fault tolerant (BFT) consensus engine that orders transactions. However, currently available BFT solutions in Hyperledger (as well as in the literature at large) are inadequate for financial settings; they are not designed to ensure fairness or to tolerate selfish behavior that arises when financial institutions strive to maximize their own profit. We present FairLedger, a permissioned blockchain BFT protocol, which is fair, designed to deal with rational behavior, and, no less important, easy to understand and implement. The secret sauce of our protocol is a new communication abstraction, called detectable all-to-all (DA2A), which allows us to detect participants (byzantine or rational) that deviate from the protocol, and punish them. We implement FairLedger in the Hyperledger open source project, using Iroha framework, one of the biggest projects therein. To evaluate FairLegder's performance, we also implement it in the PBFT framework and compare the two protocols. Our results show that in failure-free scenarios FairLedger achieves better throughput than both Iroha's implementation and PBFT in wide-area settings. Kfir Lev-Ari, Alexander Spiegelman, Idit Keidar, Dahlia Malkhi |
OPODIS | 4 |
| 2019 | Asymptotically Optimal Validated Asynchronous Byzantine AgreementabstractWe provide a new protocol for Validated Asynchronous Byzantine Agreement in the authenticated setting. Validated (multi-valued) Asynchronous Byzantine Agreement is a key building block in constructing Atomic Broadcast and fault-tolerant state machine replication in the asynchronous setting. Our protocol has optimal resilience of ƒ < n/3 Byzantine failures and asymptotically optimal expected O(1) running time to reach agreement. Honest parties in our protocol send only an expected O(n2) messages where each message contains a value and a constant number of signatures. Hence our total expected communication is O(n2) words. The best previous result of Cachin et al. from 2001 solves Validated Byzantine Agreement with optimal resilience and O(1) expected time but with O(n3) expected word communication. Our work addresses an open question of Cachin et al. from 2001 and improves the expected word communication from O(n3) to asymptotically optimal O(n2). Ittai Abraham, Dahlia Malkhi, Alexander Spiegelman |
PODC | 2 |
| 2019 | HotStuff: BFT Consensus with Linearity and ResponsivenessabstractWe present HotStuff, a leader-based Byzantine fault-tolerant replication protocol for the partially synchronous model. Once network communication becomes synchronous, HotStuff enables a correct leader to drive the protocol to consensus at the pace of actual (vs. maximum) network delay--a property called responsiveness---and with communication complexity that is linear in the number of replicas. To our knowledge, HotStuff is the first partially synchronous BFT replication protocol exhibiting these combined properties. Its simplicity enables it to be further pipelined and simplified into a practical, concise protocol for building large-scale replication services. Maofan Yin, Dahlia Malkhi, Michael K. Reiter, Guy Golan-Gueta, Ittai Abraham |
PODC | 2 |
| 2018 | Stable and Consistent Membership at Scale with Rapid
Lalith Suresh 0001, Dahlia Malkhi, Parikshit Gopalan, Ivan Porto Carreiro, Zeeshan Lokhandwala |
USENIX ATC | 2 |
| 2018 | State Machine Replication Is More Expensive Than Consensus
Karolos Antoniadis, Rachid Guerraoui, Dahlia Malkhi, Dragos-Adrian Seredinschi |
DISC | 3 |
| 2018 | Introduction to the Special Issue on the Award Papers of USENIX ATC 2019abstractNo abstract available. Dahlia Malkhi, Dan Tsafrir |
ACM Trans. Comput. Syst. | 1 |
| 2017 | vCorfu: A Cloud-Scale Object Store on a Shared Log
Michael Wei, Amy Tai, Christopher J. Rossbach, Ittai Abraham, Maithem Munshed, Medhavi Dhawan, Jim Stabile, Udi Wieder, Scott Fritchie, Steven Swanson, Michael J. Freedman, Dahlia Malkhi |
NSDI | 12 |
| 2017 | Solida: A Blockchain Protocol Based on Reconfigurable Byzantine ConsensusabstractThe decentralized cryptocurrency Bitcoin has experienced great success but also encountered many challenges. One of the challenges has been the long confirmation time. Another challenge is the lack of incentives at certain steps of the protocol, raising concerns for transaction withholding, selfish mining, etc. To address these challenges, we propose Solida, a decentralized blockchain protocol based on reconfigurable Byzantine consensus augmented by proof-of-work. Solida improves on Bitcoin in confirmation time, and provides safety and liveness assuming the adversary control less than (roughly) one-third of the total mining power. Ittai Abraham, Dahlia Malkhi, Kartik Nayak, Ling Ren 0001, Alexander Spiegelman |
OPODIS | 2 |
| 2017 | Dynamic Reconfiguration: Abstraction and Optimal Asynchronous SolutionabstractProviding clean and efficient foundations and tools for reconfiguration is a crucial enabler for distributed system management today. This work takes a step towards developing such foundations. It considers classic fault-tolerant atomic objects emulated on top of a static set of fault-prone servers, and turns them into dynamic ones. The specification of a dynamic object extends the corresponding static (non-dynamic) one with an API for changing the underlying set of fault-prone servers. Thus, in a dynamic model, an object can start in some configuration and continue in a different one. Its liveness is preserved through the reconfigurations it undergoes, tolerating a versatile set of faults as it shifts from one configuration to another. In this paper we present a general abstraction for asynchronous reconfiguration, and exemplify its usefulness for building two dynamic objects: a read/write register and a max-register. We first define a dynamic model with a clean failure condition that allows an administrator to reconfigure the system and switch off a server once the reconfiguration operation removing it completes. We then define the Reconfiguration abstraction and show how it can be used to build dynamic registers and max-registers. Finally, we give an optimal asynchronous algorithm implementing the Reconfiguration abstraction, which in turn leads to the first asynchronous (consensus-free) dynamic register emulation with optimal complexity. More concretely, faced with n requests for configuration changes, the number of configurations that the dynamic register is implemented over is n; and the complexity of each client operation is O(n). Alexander Spiegelman, Idit Keidar, Dahlia Malkhi |
DISC | 3 |
| 2017 | Apache REEF: Retainable Evaluator Execution FrameworkabstractResource Managers like YARN and Mesos have emerged as a critical layer in the cloud computing system stack, but the developer abstractions for leasing cluster resources and instantiating application logic are very low level. This flexibility comes at a high cost in terms of developer effort, as each application must repeatedly tackle the same challenges (e.g., fault tolerance, task scheduling and coordination) and reimplement common mechanisms (e.g., caching, bulk-data transfers). This article presents REEF, a development framework that provides a control plane for scheduling and coordinating task-level (data-plane) work on cluster resources obtained from a Resource Manager. REEF provides mechanisms that facilitate resource reuse for data caching and state management abstractions that greatly ease the development of elastic data processing pipelines on cloud platforms that support a Resource Manager service. We illustrate the power of REEF by showing applications built atop: a distributed shell application, a machine-learning framework, a distributed in-memory caching system, and a port of the CORFU system. REEF is currently an Apache top-level project that has attracted contributors from several institutions and it is being used to develop several commercial offerings such as the Azure Stream Analytics service. Byung-Gon Chun, Tyson Condie, Yingda Chen, Carlo Curino, Chris Douglas, Matteo Interlandi, Beomyeol Jeon, Joo Seong Jeong, Gyewon Lee, Yunseong Lee, Tony Majestro, Dahlia Malkhi, Sergiy Matusevych, Brandon Myers, Mariia Mykhailova, Shravan M. Narayanamurthy, Joseph Noor, Raghu Ramakrishnan 0001, Sriram Rao, Russell Sears, Beysim Sezgin, Taegeon Um, Julia Wang, Markus Weimer, Youngseok Yang |
ACM Trans. Comput. Syst. | 14 |
| 2016 | Silver: A Scalable, Distributed, Multi-versioning, Always Growing (Ag) File System
Michael Wei, Christopher J. Rossbach, Ittai Abraham, Udi Wieder, Steven Swanson, Dahlia Malkhi, Amy Tai |
HotStorage | 6 |
| 2016 | Flexible Paxos: Quorum Intersection RevisitedabstractDistributed consensus is integral to modern distributed systems. The widely adopted Paxos algorithm uses two phases, each requiring majority agreement, to reliably reach consensus. In this paper, we demonstrate that Paxos, which lies at the foundation of many production systems, is conservative. Specifically, we observe that each of the phases of Paxos may use non-intersecting quorums. Majority quorums are not necessary as intersection is required only across phases. Using this weakening of the requirements made in the original formulation, we propose Flexible Paxos, which generalizes over the Paxos algorithm to provide flexible quorums. We show that Flexible Paxos is safe, e cient and easy to utilize in existing distributed systems. We discuss far reaching implications of this result. For example, improved availability results from reducing the size of second phase quorums by one when the system size is even, while keeping majority quorums in the first phase. Another example is improved throughput of replication by using much smaller phase 2 quorums, while increasing the leader election (phase 1) quorums. Finally, non intersecting quorums in either first or second phases may enhance the efficiency of both. Heidi Howard, Dahlia Malkhi, Alexander Spiegelman |
OPODIS | 2 |
| 2016 | Replex: A Scalable, Highly Available Multi-Index Data Store
Amy Tai, Michael Wei, Michael J. Freedman, Ittai Abraham, Dahlia Malkhi |
USENIX ATC | 5 |
| 2015 | Dynamic Reconfiguration: A Tutorial (Tutorial)abstractA key challenge for distributed systems is the problem of reconfiguration. Clearly, any production storage system that provides data reliability and availability for long periods must be able to reconfigure in order to remove failed or old servers and add healthy or new ones. This is far from trivial since we do not want the reconfiguration management to be centralized or cause a system shutdown. In this tutorial we look into existing reconfigurable storage algorithms. We propose a common model and failure condition capturing their guarantees. We define a reconfiguration problem around which dynamic object solutions may be designed. To demonstrate its strength, we use it to implement dynamic atomic storage. We present a generic framework for solving the reconfiguration problem, show how to recast existing algorithms in terms of this framework, and compare among them. Alexander Spiegelman, Idit Keidar, Dahlia Malkhi |
OPODIS | 3 |
| 2015 | Distributed Resource Discovery in Sub-Logarithmic TimeabstractWe present a new distributed algorithm for the resource discovery problem introduced by Harchol-Balter, Leighton, and Levin in PODC'99. The resource discovery problem consists of a synchronous network with n machines in which at any timestep any machine v can PUSH or PULL a message to/from any other machine u whose (IP) address is known to v. Messages can contain addresses which then change the "topology". The goal of a distributed resource discovery problem is to enable all machines to learn the addresses of all other machines as fast as possible while keeping the number of messages sent low. Bernhard Haeupler, Dahlia Malkhi |
PODC | 2 |
| 2015 | REEF: Retainable Evaluator Execution FrameworkabstractResource Managers like Apache YARN have emerged as a critical layer in the cloud computing system stack, but the developer abstractions for leasing cluster resources and instantiating application logic are very low-level. This flexibility comes at a high cost in terms of developer effort, as each application must repeatedly tackle the same challenges (e.g., fault-tolerance, task scheduling and coordination) and re-implement common mechanisms (e.g., caching, bulk-data transfers). This paper presents REEF, a development framework that provides a control-plane for scheduling and coordinating task-level (data-plane) work on cluster resources obtained from a Resource Manager. REEF provides mechanisms that facilitate resource re-use for data caching, and state management abstractions that greatly ease the development of elastic data processing work-flows on cloud platforms that support a Resource Manager service. REEF is being used to develop several commercial offerings such as the Azure Stream Analytics service. Furthermore, we demonstrate REEF development of a distributed shell application, a machine learning algorithm, and a port of the CORFU [4] system. REEF is also currently an Apache Incubator project that has attracted contributors from several instititutions. Markus Weimer, Yingda Chen, Byung-Gon Chun, Tyson Condie, Carlo Curino, Chris Douglas, Yunseong Lee, Tony Majestro, Dahlia Malkhi, Sergiy Matusevych, Brandon Myers, Shravan M. Narayanamurthy, Raghu Ramakrishnan 0001, Sriram Rao, Russell Sears, Beysim Sezgin, Julia Wang |
SIGMOD Conference | 9 |
| 2015 | Elastic Configuration Maintenance via a Parsimonious Speculating Snapshot Solution
Eli Gafni, Dahlia Malkhi |
DISC | 2 |
| 2014 | Optimal gossip with direct addressingabstractGossip algorithms spread information in distributed networks by nodes repeatedly forwarding information to a few random contacts. By their very nature, gossip algorithms tend to be distributed and fault tolerant. If done right, they can also be fast and message-efficient. A common model for gossip communication is the random phone call model, in which in each synchronous round each node can PUSH or PULL information to or from a random other node. For example, Karp et al. [FOCS 2000] gave algorithms in this model that spread a message to all nodes in Θ(log n) rounds while sending only O(log log n) messages per node on average. They also showed that at least Θ(log n) rounds are necessary in this model and that algorithms achieving this round-complexity need to send ω(1) messages per node on average. Recently, Avin and Elsasser [DISC 2013], studied the random phone call model with the natural and commonly used assumption of direct addressing. Direct addressing allows nodes to directly contact nodes whose ID (e.g., IP address) was learned before. They show that in this setting, one can "break the log n barrier" and achieve a gossip algorithm running in O(√log n) rounds, albeit while using O(√log n) messages per node. Bernhard Haeupler, Dahlia Malkhi |
PODC | 2 |
| 2013 | Tango: distributed data structures over a shared logabstractDistributed systems are easier to build than ever with the emergence of new, data-centric abstractions for storing and computing over massive datasets. However, similar abstractions do not exist for storing and accessing meta-data. To fill this gap, Tango provides developers with the abstraction of a replicated, in-memory data structure (such as a map or a tree) backed by a shared log. Tango objects are easy to build and use, replicating state via simple append and read operations on the shared log instead of complex distributed protocols; in the process, they obtain properties such as linearizability, persistence and high availability from the shared log. Tango also leverages the shared log to enable fast transactions across different objects, allowing applications to partition state across machines and scale to the limits of the underlying log without sacrificing consistency. Mahesh Balakrishnan 0001, Dahlia Malkhi, Ted Wobber, Ming Wu 0007, Vijayan Prabhakaran, Michael Wei, John D. Davis, Sriram Rao, Tao Zou 0002, Aviad Zuck |
SOSP | 2 |
| 2013 | Beyond block I/O: implementing a distributed shared log in hardwareabstractThe basic block I/O interface used for interacting with storage devices hasn't changed much in 30 years. With the advent of very fast I/O devices based on solid-state memory, it becomes increasingly attractive to make many devices directly and concurrently available to many clients. However, when multiple clients share media at fine grain, retaining data consistency is problematic: SCSI, IDE, and their descendants don't offer much help. We propose an interface to networked storage that reduces an existing software implementation of a distributed shared log to hardware. Our system achieves both scalable throughput and strong consistency, while obtaining significant benefits in cost and power over the software implementation. Michael Wei, John D. Davis, Ted Wobber, Mahesh Balakrishnan 0001, Dahlia Malkhi |
SYSTOR | 5 |
| 2013 | CORFU: A distributed shared logabstractCORFU is a global log which clients can append-to and read-from over a network. Internally, CORFU is distributed over a cluster of machines in such a way that there is no single I/O bottleneck to either appends or reads. Data is fully replicated for fault tolerance, and a modest cluster of about 16--32 machines with SSD drives can sustain 1 million 4-KByte operations per second. The CORFU log enabled the construction of a variety of distributed applications that require strong consistency at high speeds, such as databases, transactional key-value stores, replicated state machines, and metadata services. Mahesh Balakrishnan 0001, Dahlia Malkhi, John D. Davis, Vijayan Prabhakaran, Michael Wei, Ted Wobber |
ACM Trans. Comput. Syst. | 2 |
| 2012 | CORFU: A Shared Log Design for Flash Clusters
Mahesh Balakrishnan 0001, Dahlia Malkhi, Vijayan Prabhakaran, Ted Wobber, Michael Wei, John D. Davis |
NSDI | 2 |
| 2012 | Dynamic Reconfiguration of Primary/Backup Clusters
Alexander Shraer, Benjamin C. Reed, Dahlia Malkhi, Flavio Paiva Junqueira |
USENIX ATC | 3 |
| 2012 | Efficient distributed approximation algorithms via probabilistic tree embeddings
Maleq Khan, Fabian Kuhn, Dahlia Malkhi, Gopal Pandurangan, Kunal Talwar |
Distributed Comput. | 3 |
| 2011 | DISC 2011 Invited Lecture by Dahlia Malkhi: Going beyond Paxos
Mahesh Balakrishnan 0001, Dahlia Malkhi, Vijayan Prabhakaran, Ted Wobber |
DISC | 2 |
| 2011 | Dynamic atomic storage without consensusabstractThis article deals with the emulation of atomic read/write (R/W) storage in dynamic asynchronous message passing systems. In static settings, it is well known that atomic R/W storage can be implemented in a fault-tolerant manner even if the system is completely asynchronous, whereas consensus is not solvable. In contrast, all existing emulations of atomic storage in dynamic systems rely on consensus or stronger primitives, leading to a popular belief that dynamic R/W storage is unattainable without consensus. In this article, we specify the problem of dynamic atomic read/write storage in terms of the interface available to the users of such storage. We discover that, perhaps surprisingly, dynamic R/W storage is solvable in a completely asynchronous system: we present DynaStore, an algorithm that solves this problem. Our result implies that atomic R/W storage is in fact easier than consensus, even in dynamic systems. Marcos K. Aguilera, Idit Keidar, Dahlia Malkhi, Alexander Shraer |
J. ACM | 3 |
| 2010 | Differential RAID: rethinking RAID for SSD reliabilityabstractSSDs exhibit very different failure characteristics compared to hard drives. In particular, the Bit Error Rate (BER) of an SSD climbs as it receives more writes. As a result, RAID arrays composed from SSDs are subject to correlated failures. By balancing writes evenly across the array, RAID schemes can wear out devices at similar times. When a device in the array fails towards the end of its lifetime, the high BER of the remaining devices can result in data loss. We propose Diff-RAID, a parity-based redundancy solution that creates an age differential in an array of SSDs. Diff-RAID distributes parity blocks unevenly across the array, leveraging their higher update rate to age devices at different rates. To maintain this age differential when old devices are replaced by new ones, Diff-RAID reshuffles the parity distribution on each drive replacement. We evaluate Diff-RAID's reliability by using real BER data from 12 flash chips on a simulator and show that it is more reliable than RAID-5, in some cases by multiple orders of magnitude. We also evaluate Diff-RAID's performance using a software implementation on a 5-device array of 80 GB Intel X25-M SSDs and show that it offers a trade-off between throughput and reliability. Mahesh Balakrishnan 0001, Asim Kadav, Vijayan Prabhakaran, Dahlia Malkhi |
EuroSys | 4 |
| 2010 | Fast Asynchronous Consensus with Optimal Resilience
Ittai Abraham, Marcos K. Aguilera, Dahlia Malkhi |
DISC | 3 |
| 2010 | Brief Announcement: Flash-Log - A High Throughput Log
Mahesh Balakrishnan 0001, Philip A. Bernstein, Dahlia Malkhi, Vijayan Prabhakaran, Colin W. Reid |
DISC | 3 |
| 2010 | Strong-Diameter Decompositions of Minor Free Graphs
Ittai Abraham, Cyril Gavoille, Dahlia Malkhi, Udi Wieder |
Theory Comput. Syst. | 3 |
| 2010 | Differential RAID: Rethinking RAID for SSD reliabilityabstractSSDs exhibit very different failure characteristics compared to hard drives. In particular, the bit error rate (BER) of an SSD climbs as it receives more writes. As a result, RAID arrays composed from SSDs are subject to correlated failures. By balancing writes evenly across the array, RAID schemes can wear out devices at similar times. When a device in the array fails towards the end of its lifetime, the high BER of the remaining devices can result in data loss. We propose Diff-RAID, a parity-based redundancy solution that creates an age differential in an array of SSDs. Diff-RAID distributes parity blocks unevenly across the array, leveraging their higher update rate to age devices at different rates. To maintain this age differential when old devices are replaced by new ones, Diff-RAID reshuffles the parity distribution on each drive replacement. We evaluate Diff-RAID's reliability by using real BER data from 12 flash chips on a simulator and show that it is more reliable than RAID-5, in some cases by multiple orders of magnitude. We also evaluate Diff-RAID's performance using a software implementation on a 5-device array of 80 GB Intel X25-M SSDs and show that it offers a trade-off between throughput and reliability. Mahesh Balakrishnan 0001, Asim Kadav, Vijayan Prabhakaran, Dahlia Malkhi |
ACM Trans. Storage | 4 |
| 2009 | RPC Chains: Efficient Client-Server Communication in Geodistributed Systems
Yee Jiun Song, Marcos K. Aguilera, Ramakrishna Kotla, Dahlia Malkhi |
NSDI | 4 |
| 2009 | Dynamic atomic storage without consensusabstractThis paper deals with the emulation of atomic read/write (R/W) storage in dynamic asynchronous message passing systems. In static settings, it is well known that atomic R/W storage can be implemented in a fault-tolerant manner even if the system is completely asynchronous, whereas consensus is not solvable. In contrast, all existing emulations of atomic storage in dynamic systems rely on consensus or stronger primitives, leading to a popular belief that dynamic R/W storage is unattainable without consensus. Marcos K. Aguilera, Idit Keidar, Dahlia Malkhi, Alexander Shraer |
PODC | 3 |
| 2009 | Vertical paxos and primary-backup replicationabstractNo abstract available. Leslie Lamport, Dahlia Malkhi, Lidong Zhou |
PODC | 2 |
| 2009 | Compact Multicast Routing
Ittai Abraham, Dahlia Malkhi, David Ratajczak |
DISC | 2 |
| 2009 | Virtual Ring Routing Trends
Dahlia Malkhi, Siddhartha Sen 0001, Kunal Talwar, Renato F. Werneck, Udi Wieder |
DISC | 1 |
| 2009 | Chasing the Weakest System Model for Implementing Ω and ConsensusabstractAguilera et al. and Malkhi et al. presented two system models, which are weaker than all previously proposed models where the eventual leader election oracle Ω can be implemented, and thus, consensus can also be solved. The former model assumes unicast steps and at least one correct process with f outgoing eventually timely links, whereas the latter assumes broadcast steps and at least one correct process with f bidirectional but moving eventually timely links. Consequently, those models are incomparable. In this paper, we show that Ω can also be implemented in a system with at least one process with f outgoing moving eventually timely links, assuming either unicast or broadcast steps. It seems to be the weakest system model that allows to solve consensus via Ω-based algorithms known so far. We also provide matching lower bounds for the communication complexity of Ω in this model, which are based on an interesting “stabilization property” of infinite runs. Those results reveal a fairly high price to be paid for this further relaxation of synchrony properties. Martin Hutle, Dahlia Malkhi, Ulrich Schmid 0001, Lidong Zhou |
IEEE Trans. Dependable Secur. Comput. | 2 |
| 2008 | Efficient distributed approximation algorithms via probabilistic tree embeddingsabstractWe present a uniform approach to design efficient distributed approximation algorithms for various network optimization problems. Our approach is randomized and based on a probabilistic tree embedding due to Fakcharoenphol, Rao, and Talwar (FRT embedding). We show how to efficiently compute an (implicit) FRT embedding in a decentralized manner and how to use the embedding to obtain expected O(log n)-approximate distributed algorithms for the generalized Steiner forest problem, the minimum routing cost spanning tree problem, and the $k$-source shortest paths problem in arbitrary networks. The time complexities of our algorithms are within a polylogarithmic factor of the optimum. Maleq Khan, Fabian Kuhn, Dahlia Malkhi, Gopal Pandurangan, Kunal Talwar |
PODC | 3 |
| 2008 | On spreading recommendations via social gossipabstractThis paper introduces and analyzes a variant of distributed gossip which is motivated by the sharing of recommendations in a social network. The social settings bear two implications on gossip. First, rumors fade after a few hops, and so does our gossip mechanism. Second, users require a rumor to be substantiated by multiple, independent sources in order to adopt it. Consequently, in our social gossip a message is adopted only when it is received over a threshold of independent paths. Social gossip is a new, highly relevant and practically motivated variant of distributed gossip, whose analysis contributes to the fundamental theory of distributed algorithms. Yaacov Fernandess, Dahlia Malkhi |
SPAA | 2 |
| 2008 | A unifying framework of rating users and data items in peer-to-peer and social networks
Danny Bickson, Dahlia Malkhi |
Peer-to-Peer Netw. Appl. | 2 |
| 2008 | Compact name-independent routing with minimum stretchabstractGiven a weighted undirected network with arbitrary node names, we present a compact routing scheme, using a Õ (√n,) space routing table at each node, and routing along paths of stretch 3, that is, at most thrice as long as the minimum cost paths. This is optimal in a very strong sense. It is known that no compact routing using o ( n ) space per node can route with stretch below 3. Also, it is known that any stretch below 5 requires Ω(√ n ,)space per node. Ittai Abraham, Cyril Gavoille, Dahlia Malkhi, Noam Nisan, Mikkel Thorup |
ACM Trans. Algorithms | 3 |
| 2007 | Peer-to-Peer RatingabstractTraditional instant messaging applications rely on central server infrastructure to broker user information. The cost and complexity of this infrastructure makes it difficult for developers to build and deploy lightweight presence and instant messaging systems within their own applications. In this paper, we describe P2P-IM, a peer-to-peer instant messaging client that does not rely on any hosted server infrastructure. The system provides the rich facilities available from traditional client-server systems but enables easy deployment and integration with existing applications. The solution provides simplified identity generation, connectivity, and rich per-application data publication. Danny Bickson, Dahlia Malkhi, Lidong Zhou |
Peer-to-Peer Computing | 2 |
| 2007 | Reconstructing approximate tree metricsabstractWe introduce a novel measure called ε-four-pointscondition (ε-4PC), which assigns a value ε ∈ [0,1] to every metric space quantifying how close the metric is to a tree metric. Data-sets taken from real Internet measurements indicate remarkable closeness of Internet latencies to tree metrics based on this condition. We study embeddings of ε-4PC metric spaces into trees and prove tight upper and lower bounds. Specifically, we show that there are constants c1 and c2 such that, (1) every metric (X,d) which satisfies the ε-4PC can be embedded into a tree with distortion (1+ε)c1log|X|, and (2) for every ε ∈: [0,1] and any number of nodes, there is a metric space (X,d) satisfying the ε-4PC that does not embed into a tree with distortion less than (1+ε)c2log|X|. In addition, we prove a lower bound on approximate distance labelings of ε-4PC metrics, and give tight bounds for tree embeddings with additive error guarantees. Ittai Abraham, Mahesh Balakrishnan 0001, Fabian Kuhn, Dahlia Malkhi, Venugopalan Ramasubramanian, Kunal Talwar |
PODC | 4 |
| 2007 | Strong-diameter decompositions of minor free graphsabstractWe provide the first sparse covers and probabilistic partitions for graphs excluding a fixed minor that have strong diameter bounds; i.e. each set of the cover/partition has a small diameter as an induced sub-graph. Using these results we provide improved distributed name-independent routing schemes. Specifically, given a graph excluding a minor on r vertices and a parameter ρ > 0 we obtain the flowing results: (1) a polynomial algorithm that constructs a set of clusters such that each cluster has a strong-diameter of O(r2ρ) and each vertex belongs to 2O(r)r! clusters; (2) a name-independent routing scheme with a stretch of O(r2) and tables of size 2O(r)r! log4n bits; (3) a randomized algorithm that partitions the graph such that each cluster has strong-diameter O(r6r ρ) and the probability an edge (u, v) is cut is O(r d(u, v)/ρ). Ittai Abraham, Cyril Gavoille, Dahlia Malkhi, Udi Wieder |
SPAA | 3 |
| 2007 | Concise version vectors in WinFS
Dahlia Malkhi, Douglas B. Terry |
Distributed Comput. | 1 |
| 2007 | Addendum to "Scalable secure storage when half the system is faulty" [Inform. Comput 174 (2)(2002) 203-213]
Noga Alon, Haim Kaplan, Michael Krivelevich, Dahlia Malkhi, Julien P. Stern |
Inf. Comput. | 4 |
| 2007 | Wait-free regular storage from Byzantine components
Ittai Abraham, Gregory V. Chockler, Idit Keidar, Dahlia Malkhi |
Inf. Process. Lett. | 4 |
| 2007 | On collaborative content distribution using multi-message gossip
Yaacov Fernandess, Dahlia Malkhi |
J. Parallel Distributed Comput. | 2 |
| 2006 | Routing in Networks with Low Doubling DimensionabstractThis paper studies compact routing schemes for networks with low doubling dimension. Two variants are explored, name-independent routing and labeled routing. The key results obtained for this model are the following. First, we provide the first name-independent solution. Specifically, we achieve constant stretch and polylogarithmic storage. Second, we obtain the first truly scale-free solutions, namely, the network’s aspect ratio is not a factor in the stretch. Scale-free schemes are given for three problem models: name-independent routing on graphs, labeled routing on metric spaces, and labeled routing on graphs. Third, we prove a lower bound requiring linear storage for stretch \gt 3 schemes. This has the important ramification of separating for the first time the name-independent problem model from the labeled model for these networks, since compact stretch-1+e labeled schemes are known to be possible. Ittai Abraham, Cyril Gavoille, Andrew V. Goldberg, Dahlia Malkhi |
ICDCS | 4 |
| 2006 | On collaborative content distribution using multi-message gossipabstractWe study epidemic schemes in the context of collaborative data delivery. In this context, multiple chunks of data reside at different nodes, and the challenge is to simultaneously deliver all chunks to all nodes. Here we explore the inter-operation between the gossip of multiple, simultaneous message-chunks. In this setting, interacting nodes must select which chunk, among many, to exchange in every communication round. We provide an efficient solution that possesses the inherent robustness and scalability of gossip. Our approach maintains the simplicity of gossip, and has low message, connections and computation overhead. Because our approach differs from solutions proposed by network coding, we are able to provide insight into the tradeoffs and analysis of the problem of collaborative content distribution. We formally analyze the performance of the algorithm, demonstrating its efficiency with high probability. Yaacov Fernandess, Dahlia Malkhi |
IPDPS | 2 |
| 2006 | On space-stretch trade-offs: lower boundsabstractOne of the fundamental trade-offs in compact routing schemes is between the space used to store the routing table on each node and the stretch factor of the routing scheme -- the ratio between the cost of the route induced by the scheme and the cost of a minimum cost path between the same pair. Using a distributed Kolmogorov Complexity argument, we give a lower bound for the name-independent model that applies even to single-source schemes and does not require a girth conjecture. For any integer k ≥ 1 we prove that any routing scheme for networks with arbitrary weights and arbitrary node names (even a single-source routing scheme) with maximum stretch strictly less than 2k + 1 requires Ω((n log n)1/k)-bit routing tables. We extend our results to lower bound the average-stretch, showing that for any integer k ≥ 1 any name-independent routing scheme with (n/(9k))1/k-bit routing tables has average-stretch of at least k/4 + 7/8. This result is in sharp contrast to recent results on the average-stretch of labeled routing schemes. Ittai Abraham, Cyril Gavoille, Dahlia Malkhi |
SPAA | 3 |
| 2006 | On space-stretch trade-offs: upper boundsabstractInternational audience Ittai Abraham, Cyril Gavoille, Dahlia Malkhi |
SPAA | 3 |
| 2006 | Brief Announcement: Chasing the Weakest System Model for Implementing Omega and Consensus
Martin Hutle, Dahlia Malkhi, Ulrich Schmid 0001, Lidong Zhou |
SSS | 2 |
| 2006 | Byzantine disk paxos: optimal resilience with byzantine shared memory
Ittai Abraham, Gregory V. Chockler, Idit Keidar, Dahlia Malkhi |
Distributed Comput. | 4 |
| 2005 | Hold Your Sessions: An Attack on Java Session-Id Generation
Zvi Gutterman, Dahlia Malkhi |
CT-RSA | 2 |
| 2005 | Name independent routing for growth bounded networksabstractA weighted undirected network is Δ growth-bounded if the number of nodes at distance 2r around any given node is at most Δ times the number of nodes at distance r around the node. Given a weighted undirected network with arbitrary node names and ε > 0, we present a routing scheme that routes along paths of stretch 1+ε and uses with high probability only O(1/εO (log Δ)log5n) bit routing tables per node. Ittai Abraham, Dahlia Malkhi |
SPAA | 2 |
| 2005 | Compact Routing for Graphs Excluding a Fixed Minor
Ittai Abraham, Cyril Gavoille, Dahlia Malkhi |
DISC | 3 |
| 2005 | Papillon: Greedy Routing in Rings
Ittai Abraham, Dahlia Malkhi, Gurmeet Singh Manku |
DISC | 2 |
| 2005 | Omega Meets Paxos: Leader Election and Stability Without Eventual Timely Links
Dahlia Malkhi, Florian Oprea, Lidong Zhou |
DISC | 1 |
| 2005 | Concise Version Vectors in WinFS
Dahlia Malkhi, Douglas B. Terry |
DISC | 1 |
| 2005 | Probabilistic quorums for dynamic systems
Ittai Abraham, Dahlia Malkhi |
Distributed Comput. | 2 |
| 2005 | Active Disk Paxos with infinitely many processes
Gregory V. Chockler, Dahlia Malkhi |
Distributed Comput. | 2 |
| 2004 | Byzantine disk paxos: optimal resilience with byzantine shared memoryabstractWe present Byzantine Disk Paxos, an asynchronous shared-memory consensus protocol that uses a collection of n > 3t disks, t of which may fail by becoming non-responsive or arbitrarily corrupted. We give two constructions of this protocol; that is, we construct two different building blocks, each of which can be used, along with a leader oracle, to solve consensus. One building block is a shared wait-free safe register. The second building block is a regular register that satisfies a weaker termination (liveness) condition than wait freedom: its write operations are wait-free, whereas its read operations are guaranteed to return only in executions with a finite number of writes. We call this termination condition finite writes (FW), and show that consensus is solvable with FW-terminating registers and a leader oracle. We construct each of these reliable registers from n > 3t base registers, t of which can be non-responsive or Byzantine. All the previous wait-free constructions in this model used at least 4t+1 fault-prone registers, and we are not familiar with any prior FW-terminating constructions in this model. Ittai Abraham, Gregory V. Chockler, Idit Keidar, Dahlia Malkhi |
PODC | 4 |
| 2004 | Compact routing on euclidian metricsabstractWe consider the problem of designing a compact communication network that supports efficient routing in an Euclidean plane. Our network design and routing scheme achieves 1+ε stretch, logarithmic diameter, and constant out degree. This improves upon the best known result so far that requires a logarithmic out-degree. Furthermore, our scheme is asymptotically optimal in Euclidean metrics whose diameter is polynomial. Ittai Abraham, Dahlia Malkhi |
PODC | 2 |
| 2004 | LAND: stretch (1 + epsilon) locality-aware networks for DHTs
Ittai Abraham, Dahlia Malkhi, Oren Dobzinski |
SODA | 2 |
| 2004 | Compact name-independent routing with minimum stretchabstractGiven a weighted undirected network with arbitrary node names, we present a compact routing scheme, using a O(√n) space routing table at each node, and routing along paths of stretch 3, that is, at most thrice as long as the shortest paths. This is optimal in a very strong sense. It is known that no compact routing using o(n) space per node can route with stretch below 3. Also, it is known that any stretch below 5 requires Ω(√n) space per node. Ittai Abraham, Cyril Gavoille, Dahlia Malkhi, Noam Nisan, Mikkel Thorup |
SPAA | 3 |
| 2004 | Fairplay - Secure Two-Party Computation System
Dahlia Malkhi, Noam Nisan, Benny Pinkas, Yaron Sella |
USENIX Security Symposium | 1 |
| 2004 | Routing with Improved Communication-Space Trade-Off
Ittai Abraham, Cyril Gavoille, Dahlia Malkhi |
DISC | 3 |
| 2003 | Probabilistic Quorums for Dynamic Systems
Ittai Abraham, Dahlia Malkhi |
DISC | 2 |
| 2003 | Objects shared by Byzantine processes
Dahlia Malkhi, Michael Merritt, Michael K. Reiter, Gadi Taubenfeld |
Distributed Comput. | 1 |
| 2003 | Estimating network size from local information
Keren Horowitz, Dahlia Malkhi |
Inf. Process. Lett. | 2 |
| 2003 | Diffusion without false rumors: on propagating updates in a Byzantine environment
Dahlia Malkhi, Yishay Mansour, Michael K. Reiter |
Theor. Comput. Sci. | 1 |
| 2002 | Active disk paxos with infinitely many processesabstractWe present an improvement to the Disk Paxos protocol by Gafni and Lamport which utilizes extended functionality and flexibility provided by Active Disks and supports unmediated concurrent data access by an unlimited number of processes. The solution facilitates coordination by an infinite number of clients using finite shared memory. It is based on a collection of read-modify-write objects with faults, that emulate a new, reliable shared memory abstraction called a ranked register. The required read-modify-write objects are readily available in Active Disks and in Object Storage Device controllers, making our solution suitable for state-of-the-art Storage Area Network (SAN) environments. Gregory V. Chockler, Dahlia Malkhi |
PODC | 2 |
| 2002 | Viceroy: a scalable and dynamic emulation of the butterflyabstractWe propose a family of constant-degree routing networks of logarithmic diameter, with the additional property that the addition or removal of a node to the network requires no global coordination, only a constant number of linkage changes in expectation, and a logarithmic number with high probability. Our randomized construction improves upon existing solutions, such as balanced search trees, by ensuring that the congestion of the network is always within a logarithmic factor of the optimum with high probability. Our construction derives from recent advances in the study of peer-to-peer lookup networks, where rapid changes require efficient and distributed maintenance, and where the lookup efficiency is impacted both by the lengths of paths to requested data and the presence or elimination of bottlenecks in the network. Dahlia Malkhi, Moni Naor, David Ratajczak |
PODC | 1 |
| 2002 | From Byzantine Agreement to Practical SurvivabilityabstractOnly a decade ago, issues of replication, high availability and load balancing were the focus of small, closely coupled cluster projects. Consequently, techniques for cluster management and small replication systems are abundant. However, the advent of the Internet led to wide spread and highly decentralized access of services and content that bring issues of scale and ubiquitous deployment. In particular, the need to maintain copies of replicated data consistent grows beyond the limits of any local cluster. Consequently, researchers have been looking at ways to improve scalability, survivability and dynamism of replication technology. Additionally, there are a number of recent application domains that exhibit new and challenging models for information replication. For example, advances in storage technology permit processes to share information by directly accessing data on disks that are connected to a storage area network (SAN), thereby avoiding going through a file system service. This form of direct data sharing necessitates coordination among processes contending for access to data, and presents new building blocks for doing it. New needs are also re-shaped by novel services such as Jini, a global resource discovery and location tool that allows anonymous and transient clients to be serviced; by Java-spaces, a universal shared data space; by Oceanstore, an eternal storage archive that is built of peers that have an economical incentive to cooperate; by Publius, an anonymous and survivable publishing archive; and others. Many other peer-to-peer (P2P) systems offer the potential of a truly survivable settings, but on the other hand, pose challenges of scale, dynamism and trust issues. Dahlia Malkhi |
SRDS | 1 |
| 2002 | Scalable Secure Storage When Half the System Is Faulty
Noga Alon, Haim Kaplan, Michael Krivelevich, Dahlia Malkhi, Julien P. Stern |
Inf. Comput. | 4 |
| 2001 | Backoff Protocols for Distributed Mutual Exclusion and OrderingabstractPresents a simple and efficient protocol for mutual exclusion in synchronous message-passing distributed systems subject to failures. Our protocol borrows design principles from prior work in backoff protocols for multiple access channels such as the Ethernet. Our protocol is adaptive in that the expected amortized system response time - informally, the average time a process waits before entering the critical section - is a function only of the number of clients currently contending and is independent of the maximum number of processes that might contend. In particular, in the contention-free case, a process can enter the critical section after only one round-trip message delay. We use this protocol to derive a protocol for ordering operations on a replicated object in an asynchronous distributed system subject to failures. This protocol is always safe, is probabilistically live during periods of stability and is suitable for deployment in practical systems. Gregory V. Chockler, Dahlia Malkhi, Michael K. Reiter |
ICDCS | 2 |
| 2001 | Efficient Update Diffusion in Byzantine EnvironmentsabstractWe present a protocol for diffusion of updates among replicas in a distributed system where up to b replicas may suffer Byzantine failures. Our algorithm ensures that no correct replica accepts spurious updates introduced by faulty replicas, by requiring that a replica accepts an update only after receiving it from at least b+1 distinct replicas (or directly from the update source). Our algorithm diffuses updates more efficiently than previous such algorithms and, by exploiting additional information available in some practical settings, sometimes more efficiently than known lower bounds predict. Dahlia Malkhi, Ohad Rodeh, Michael K. Reiter, Yaron Sella |
SRDS | 1 |
| 2001 | Optimal Unconditional Information Diffusion
Dahlia Malkhi, Elan Pavlov, Yaron Sella |
DISC | 1 |
| 2001 | Probabilistic Quorum Systems
Dahlia Malkhi, Michael K. Reiter, Avishai Wool, Rebecca N. Wright |
Inf. Comput. | 1 |
| 2001 | Fault Detection for Byzantine Quorum SystemsabstractIn this paper, we explore techniques to detect Byzantine server failures in asynchronous replicated data services. Our goal is to detect arbitrary failures of data servers in a system where each client accesses the replicated data at only a subset (quorum) of servers in each operation. In such a system, some correct servers can be out-of-date after a write and can therefore, return values other than the most up-to-date value in response to a client's read request, thus complicating the task of determining the number of faulty servers in the system at any point in time. We initiate the study of detecting server failures in this context, and propose two statistical approaches for estimating the risk posed by faulty servers based on responses to read requests. Lorenzo Alvisi, Dahlia Malkhi, Evelyn Tumlin Pierce, Michael K. Reiter |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2001 | On k-Set Consensus Problems in Asynchronous SystemsabstractIn this paper, we investigate the k-set consensus problem in asynchronous distributed systems. In this problem, each participating process begins the protocol with an input value and by the end of the protocol must decide on one value so that at most k total values are decided by all correct processes. We extend previous work by exploring several variations of the problem definition and model, including for the first time investigation of Byzantine failures. We show that the precise definition of the validity requirement, which characterizes what decision values are allowed as a function of the input values and whether failures occur, is crucial to the solvability of the problem. For example, we show that allowing default decisions in case of failures makes the problem solvable for most values of k despite a minority of failures, even in face of the most severe type of failures (Byzantine). We introduce six validity conditions for this problem (all considered in various contexts in the literature), and demarcate the line between possible and impossible for each case. In many cases, this line is different from the one of the originally studied k-set consensus problem. Roberto De Prisco, Dahlia Malkhi, Michael K. Reiter |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2000 | Dynamic Byzantine Quorum SystemsabstractByzantine quorum systems enhance the availability and efficiency of fault-tolerant replicated services when servers may suffer Byzantine failures. An important limitation of Byzantine quorum systems is their dependence on a static threshold limit on the number of server faults. The correctness of the system is only guaranteed if the actual number of faults is lower than the the threshold at all times. However, a threshold chosen for the worst case wastes expensive replication in the common situation where the number of faults averages well below the worst case. In this paper, we present protocols for dynamically raising and lowering the resilience threshold of a quorum-based Byzantine fault-tolerant data service in response to current information on the number of server failures. Using such protocols, a system can operate in an efficient low-threshold mode with relatively small quorums in the absence of faults, increasing and decreasing the quorum size (and thus the tolerance) as faults appear and are dealt with, respectively. Lorenzo Alvisi, Evelyn Tumlin Pierce, Dahlia Malkhi, Michael K. Reiter, Rebecca N. Wright |
DSN | 3 |
| 2000 | Scalable Secure Storage when Half the System Is Faulty
Noga Alon, Haim Kaplan, Michael Krivelevich, Dahlia Malkhi, Julien P. Stern |
ICALP | 4 |
| 2000 | Objects Shared by Byzantine Processes
Dahlia Malkhi, Michael Merritt, Michael K. Reiter, Gadi Taubenfeld |
DISC | 1 |
| 2000 | Secure Reliable Multicast Protocols in a WAN
Dahlia Malkhi, Michael Merritt, Ohad Rodeh |
Distributed Comput. | 1 |
| 2000 | The Load and Availability of Byzantine Quorum SystemsabstractReplicated services accessed via quorums enable each access to be performed at only a subset (quorum) of the servers and achieve consistency across accesses by requiring any two quorums to intersect. Recently, b-masking quorum systems, whose intersections contain at least 2b+1 servers, have been proposed to construct replicated services tolerant of b-arbitrary (Byzantine) server failures. In this paper we consider a hybrid fault model allowing benign failures in addition to the Byzantine ones. We present four novel constructions for b-masking quorum systems in this model, each of which has optimal load (the probability of access of the busiest server) or optimal availability (probability of some quorum surviving failures). To show optimality we also prove lower bounds on the load and availability of any b-masking quorum system in this model. Dahlia Malkhi, Michael K. Reiter, Avishai Wool |
SIAM J. Comput. | 1 |
| 2000 | An Architecture for Survivable Coordination in Large Distributed SystemsabstractCoordination among processes in a distributed system can be rendered very complex in a large-scale system where messages may be delayed or lost and when processes may participate only transiently or behave arbitrarily, e.g. after suffering a security breach. In this paper, we propose a scalable architecture to support coordination in such extreme conditions. Our architecture consists of a collection of persistent data servers that implement simple shared data abstractions for clients, without trusting the clients or even the servers themselves. We show that, by interacting with these untrusted servers, clients can solve distributed consensus, a powerful and fundamental coordination primitive. Our architecture is very practical, and we describe the implementation of its main components in a system called Fleet. Dahlia Malkhi, Michael K. Reiter |
IEEE Trans. Knowl. Data Eng. | 1 |
| 2000 | Secure Execution of Java Applets Using a Remote PlaygroundabstractMobile code presents a number of threats to machines that execute it. We introduce an approach for protecting machines and the resources they hold from mobile code and describe a system based on our approach for protecting host machines from Java 1.1 applets. In our approach, each Java applet downloaded to the protected domain is rerouted to a dedicated machine (or set of machines), the playground, at which it is executed. Prior to execution, the applet is transformed to use the downloading user's Web browser as a graphics terminal for its input and output, and so the user has the illusion that the applet is running on his own machine. In reality, however, mobile code runs only in the sanitized environment of the playground, where user files cannot be mounted and from which only limited network connections are accepted by machines in the protected domain. Our playground thus provides a second level of defense against mobile code that circumvents language-based defenses. This paper presents the design and implementation of a playground for Java 1.1 applets and discusses extensions of it for other forms of mobile code, including Java 1.2. Dahlia Malkhi, Michael K. Reiter |
IEEE Trans. Software Eng. | 1 |
| 1999 | On k-Set Consensus Problems in Asynchronous SystemsabstractIn this paper we investigate the k-set consensus problem in asynchronous, message-passing distributed systems.In this problem, each participating process begins the protocol with an input value and by the end of the protocol must decide on one value so that at most k different values are decided by all correct processes.We extend previous work by exploring several variations of the problem definition and model, including for the first time investigation of Byzantine failures.We show that the precise definition of the validity requirement, which characterizes what decision values are allowed as a function of the input values and whether failures occur, is crucial to the solvability of the problem.For example, we show that allowing default decisions in case of failures makes the problem solvable for most values of k despite a minority of failures, even for the most severe type of failures (Byiantine).We introduce six validity conditions for this problem (all considered in various contexts in the literature), and demarcate the line between possible and impossible for each case.In many cases this line is different from the one of the originally studied k-set consensus problem. Roberto De Prisco, Dahlia Malkhi, Michael K. Reiter |
PODC | 2 |
| 1999 | On Diffusing Updates in a Byzantine EnvironmentabstractWe study how to efficiently diffuse updates to a large distributed system of data replicas, some of which may exhibit arbitrary (Byzantine) failures. We assume that strictly fewer than t replicas fail, and that each update is initially received by at least t correct replicas. The goal is to diffuse each update to all correct replicas while ensuring that correct replicas accept no updates generated spuriously by faulty replicas. To achieve reliable diffusion, each correct replica accepts an update only after receiving it from at least t others. We provide the first analysis of epidemic-style protocols for such environments. This analysis is fundamentally different from known analyses for the benign case due to our treatment of fully Byzantine failure-which, among other things, precludes the use of digital signatures for authenticating forwarded updates. We propose two epidemic-style diffusion algorithms and two measures that characterize the efficiency of diffusion algorithms in general. We characterize both of our algorithms according to these measures, and also prove lower bounds with regards to these measures that show that our algorithms are close to optimal. Dahlia Malkhi, Yishay Mansour, Michael K. Reiter |
SRDS | 1 |
| 1998 | Probabilistic Byzantine Quorum SystemsabstractIn this paper we present probabilistic masking quorum systems, a technique for replicating data that can mask, with high probability, the arbitrary (Byzantine) failure of data servers from clients. This technique generalizes previous work on probabilistic quorum systems to mask Byzantine server failures in their full generality, and improves over previous masking quorum systems by offering better data availability and access efficiency. We define probabilistic masking quorum systems, demonstrate a novel access protocol for implementing replicated data with them, and prove general and tight lower bounds on the performance that they can achieve. We also present a probabilistic masking quorum construction that outperforms strict masking constructions in measures of both availability and efficiency. Dahlia Malkhi, Michael K. Reiter, Avishai Wool, Rebecca N. Wright |
PODC | 1 |
| 1998 | Secure Execution of Java Applets using a Remote PlaygroundabstractMobile code presents a number of threats to machines that execute it. We introduce an approach for protecting machines and the resources they hold from mobile code, and describe a system based on our approach for protecting host machines from Java 1.1 applets. In our approach, each Java applet downloaded to the protected domain is rerouted to a dedicated machine (or set of machines), the playground, at which it is executed. Prior to execution, the applet is transformed to use the downloading user's Web browser as a graphics terminal for its input and output, and so the user has the illusion that the applet is running on her own machine. In reality, however, mobile code runs only in the sanitized environment of the playground, where user files cannot be mounted and from which only limited network connections are accepted by machines in the protected domain. Our playground thus provides a second level of defense against mobile code that circumvents language based defenses. Dahlia Malkhi, Michael K. Reiter, Aviel D. Rubin |
S&P | 1 |
| 1998 | Secure and Scalable Replication in PhalanxabstractPhalanx is a software system for building a persistent, survivable data repository that supports shared data abstractions (e.g., variables, mutual exclusion) for clients. Phalanx implements data abstraction that ensures useful properties without trusting the servers supporting these abstractions or the clients accessing them, i.e., Phalanx can survive even the arbitrarily malicious corruption of clients and (some number of) servers. At the core of the system are survivable replication techniques that enable efficient scaling to hundreds of Phalanx servers. In this paper we describe the implementation of some of the data abstractions provided by Phalanx, discuss their ability to scale to large systems, and describe an example application. Dahlia Malkhi, Michael K. Reiter |
SRDS | 1 |
| 1998 | Survivable Consensus ObjectsabstractReaching consensus among multiple processes in a distributed system is fundamental to coordinating distributed actions. We present a new approach to building survivable consensus objects in a system consisting of a (possibly large) collection of persistent object servers and a transient population of clients. Our consensus object implementation requires minimal support from servers, but at the same time enables clients to reach coordinated decisions despite the arbitrary (Byzantine) failure of any number of clients and up to a threshold number of servers. Dahlia Malkhi, Michael K. Reiter |
SRDS | 1 |
| 1998 | Byzantine Quorum Systems
Dahlia Malkhi, Michael K. Reiter |
Distributed Comput. | 1 |
| 1998 | Auditable Metering with Lightweight SecurityabstractIn this work we suggest a new mechanism for metering the popularity of Web sites: the compact metering scheme. Our approach does not rely on client authentication or on a third party. Instead, we suggest the notion of a timing function, a computation Matthew K. Franklin, Dahlia Malkhi |
J. Comput. Secur. | 2 |
| 1997 | Unreliable Intrusion Detection in Distributed ComputationsabstractDistributed coordination is difficult, especially when the system may suffer intrusions that corrupt some component processes. We introduce the abstraction of a failure detector that a process can use to (imperfectly) detect the corruption (Byzantine failure) of another process. In general, our failure detectors can be unreliable, both by reporting a correct process to be faulty or by reporting a faulty process to be correct. However, we show that if these detectors satisfy certain plausible properties, then the well known distributed consensus problem can be solved. We also present a randomized protocol using failure detectors that solves the consensus problem if either the requisite properties of failure detectors hold or if certain highly probable events eventually occur. This work can be viewed as a generalization of benign failure detectors popular in the distributed computing literature. Dahlia Malkhi, Michael K. Reiter |
CSFW | 1 |
| 1997 | Secure Reliable Multicast Protocols in a WANabstractA secure reliable multicast protocol enables a process to send a message to a group of recipients such that all honest destinations receive the same message, despite the malicious efforts of fewer than a third of them, including the sender. This has been shown to be a useful tool in building secure distributed services, albeit with a cost that typically grows linearly with the size of the system. For very large networks, for which such a cost may be too prohibitive, we present two approaches for bringing the cost down: First, we show a protocol whose cost is on the order of the number of tolerated failures. Secondly, we show how relaxing the consistency requirement to a selected probability level of guarantee can bring down the associated cost to a constant. Dahlia Malkhi, Michael Merritt, Ohad Rodeh |
ICDCS | 1 |
| 1997 | Failure Detectors in Omission Failure EnvironmentsabstractNo abstract available. Danny Dolev, Roy Friedman 0001, Idit Keidar, Dahlia Malkhi |
PODC | 4 |
| 1997 | The Load and Availability of Byzantine Quorum SystemsabstractReplicated services accessed via quorurmcnable each access to be performed at only a subset (quorum) of the servers, and achieve consistency across accesses by requiring any two quorums to intersect.Recently, bmasking quorum systems, whose intersections contain at least 2b+l servers, have been proposed to construct replicated services tolerant of barbitrary (B ymntine) server failures.In this paper we consider a hybrid fault model allowing benign failures in addition to the Byzantine ones.We present four novel constructions for bmasking quorum systems in this model, each of which has optimal load (the probability of access of the busiest server) or optimal availability (probabllit y of some quorum surviving failures).To show optimalit y we also prove lower bounds on the load and availabilityy of any bmasking quorum system in this model.I&mission to make digilnlflmrd copies of all or piIIIof(hi~nu}icri:ll fix personal or classroom use is grw!cd without ~cc prm idcd III;IIIIICcopies are not mwlc or distrihukd I'orpmlil or wmncrciid :IdvmIIogc.theCXWrlght notice, (he title of the pohlic:it ion mMlils d:IIcnppcw.xnd WIicc is given that urpyrighl is hy pwmission OI"IIW ACM, INC. '1"0 copy otherwise.to republish, In post on wrws or to rcdistrihu[c 10 lists.ruplircs specilic pem~ission andlor 13c 1997 I'OD(' 97 .Srlnta 13m+ora [ '.4 1 I*Y.4 Dahlia Malkhi, Michael K. Reiter, Avishai Wool |
PODC | 1 |
| 1997 | Probabilistic Quorum SystemsabstractServices replicated using a quorum system allow operations to be performed at only a subset (quorum) of the servers, and ensure consistency among operations by requiring that any two quorums intersect. In this paper we explore the consequences of requiring this intersection property to hold only with very high probability. We show that doing so can offer dramatic improvements in the performance and availability of the service, both for services tolerant of benign server failures and services tolerant of arbitrary (Byzantine) ones. We also prove a lower bound on the performance that can be achieved with this technique. 1 Introduction Quorums are tools for increasing the availability and efficiency of replicated services. A quorum system is a set of subsets of servers, every pair of which intersect. Intuitively, the intersection property guarantees that if a "write" operation is performed at one quorum, and later a "read" operation at another quorum, then there is some server that obse... Dahlia Malkhi, Michael K. Reiter, Rebecca N. Wright |
PODC | 1 |
| 1997 | Byzantine Quorum SystemsabstractQuorum systems are well-known tools for ensuring the consistency and availability of replicated data despite the benign failure of data repositories. In this paper we consider the arbitrary (Byzantine) failure of data repositories and present the first study of quorum system requirements and constructions that ensure data availability and consistency despite these failures. We also consider the load associated with our quorum systems, i.e., the minimal access probability of the busiest server. For services subject to arbitrary failures, we demonstrate quorum systems over n servers with a load of O( 1 p n ), thus meeting the lower bound on load for benignly faulttolerant quorum systems. We explore several variations of our quorum systems and extend our constructions to cope with arbitrary client failures. 1 Introduction A well known way to enhance the availability and efficiency of replicated data is by using quorums . A quorum system for a universe of data servers is a collection o... Dahlia Malkhi, Michael K. Reiter |
STOC | 1 |
| 1997 | A High-Throughput Secure Reliable Multicast ProtocolabstractA (secure) reliable multicast protocol enables a process to multicast a message to a group of processes in a way that ensures that all honest destination-group members receive the same message, even if some group members and the multicast initiator a Dahlia Malkhi, Michael K. Reiter |
J. Comput. Secur. | 1 |
| 1996 | A High-Throughput Secure Reliable Multicast ProtocolabstractA reliable multicast protocol enables a process to multicast a message to a group of processes in a way that ensures that all honest destination-group members receive the same message, even if some group members and the multicast initiator are maliciously faulty. Reliable multicast has been shown to be useful for building multiparty cryptographic protocols and secure distributed services. We present a high-throughput reliable multicast protocol that tolerates the malicious behavior of up to fewer than one-third of the group members. Our protocol achieves high-throughput using a novel technique for chaining multicasts, whereby the cost of ensuring agreement on each multicast message is amortized over many multicasts. This is coupled with a novel flow-control mechanism that yields low multicast latency. Dahlia Malkhi, Michael K. Reiter |
CSFW | 1 |
| 1996 | A Framework for Partitionable Membership Service (Abstract)abstractNo abstract available. Danny Dolev, Dahlia Malkhi, Ray Strong |
PODC | 2 |
| 1994 | Uniform Actions in Asynchronous Distributed Systems (Extended Abstract)abstractWe devetop necessary conditions for the development of asynchronous distributed sofiware that will perform uniform actions (’evenis that if performed by any pro-cess, must be performed at all processes). The pa-per focuses on dynamic uniformity, which differs from ihe classical problems in that processes continually leave and join the ongoing computation. It relates the problem to asynchronous Consensus, and shows that Consensus is a harder problem. We provide a rigorous characterization of the framework upon which several existing distributed programming environments are based. And, our work shows that progress is some-times possible in a primary-partition model even when consensus is not. 1 Dahlia Malkhi, Kenneth P. Birman, Aleta Ricciardi, André Schiper |
PODC | 1 |
| 1993 | On Distributed Algorithms in a Broadcast Domain
Danny Dolev, Dahlia Malkhi |
ICALP | 2 |