EDBT 2026 Demo / reviewers in the wild / expert
Michael Dahlin
dblp:d/MDahlin · also Mike Dahlin
· DBLP profile ↗
60ranked-venue papers
5as first author
0since 2021 · last 2016
—ORCID · none
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 20 · 3 first-authorSoftware engineering, systems software and programming languages · 16 · 2 first-authorComputer networks · 10 · 1 first-authorDatabases, data management, data science and information retrieval · 10Security and privacy · 5Applied, interdisciplinary, general and emerging computing · 3Artificial intelligence and machine learning · 1Graphics, computer vision, multimedia, augmented reality and games · 1Theory of computation · 1
Expertise — from the expertise taxonomy: the topics of the expert's papers under the CCF categories. A weight counts papers with recency: 1 for a paper about the topic, 0.3 when the topic is its context, halved every five years.
| Computer architecture, parallel and distributed computing, and storage systems
37 papers |
Distributed systems · 66% Storage systems · 17% Cloud and datacenter computing · 7% | |
| Databases, data mining, and information retrieval
3 papers |
Data stream processing · 62% Query processing and optimization · 21% Database system architecture and tuning · 16% | |
| Network and information security
6 papers |
Privacy and data protection · 47% Network security · 27% Systems and software security · 22% | |
| Computer networks
7 papers |
Transport protocols and congestion control · 36% Network measurement and analytics · 28% Content delivery and video streaming · 15% |
Topics — the 30 heaviest of 86, each with the papers that count most for it
| Topic | Weight | Papers | Last | Evidence papers |
|---|---|---|---|---|
Distributed systems
fault tolerance |
0.7 | 11 | 2012 | All about Eve: Execute-Verify Replication for Multi-Core Servers · OSDI 2012 Zyzzyva: Speculative Byzantine fault tolerance · ACM Trans. Comput. Syst. 2009 Upright cluster services · SOSP 2009 |
Distributed systems › fault tolerance
byzantine fault tolerance |
0.6 | 7 | 2011 | Depot: Cloud Storage with Minimal Trust · ACM Trans. Comput. Syst. 2011 Zyzzyva: Speculative Byzantine fault tolerance · ACM Trans. Comput. Syst. 2009 Upright cluster services · SOSP 2009 |
Distributed systems
replication |
0.5 | 7 | 2012 | All about Eve: Execute-Verify Replication for Multi-Core Servers · OSDI 2012 Dual-Quorum: A Highly Available and Consistent Replication System for Edge Services · IEEE Trans. Dependable Secur. Comput. 2010 PRACTI Replication · NSDI 2006 |
Storage systems
storage reliability |
0.3 | 4 | 2011 | Depot: Cloud Storage with Minimal Trust · OSDI 2010 SafeStore: A Durable and Practical Storage System · USENIX ATC 2007 TAPER: Tiered Approach for Eliminating Redundancy in Replica Synchronization · FAST 2005 |
Distributed systems › replication
state machine replication |
0.3 | 4 | 2009 | Zyzzyva: Speculative Byzantine fault tolerance · ACM Trans. Comput. Syst. 2009 Zyzzyva: speculative byzantine fault tolerance · SOSP 2007 BAR fault tolerance for cooperative services · SOSP 2005 |
Cloud and datacenter computing
cloud storage |
0.2 | 2 | 2011 | Depot: Cloud Storage with Minimal Trust · ACM Trans. Comput. Syst. 2011 Depot: Cloud Storage with Minimal Trust · OSDI 2010 |
Distributed systems › replication
speculative replication |
0.2 | 2 | 2009 | Zyzzyva: Speculative Byzantine fault tolerance · ACM Trans. Comput. Syst. 2009 Zyzzyva: speculative byzantine fault tolerance · SOSP 2007 |
Storage systems › data redundancy
replicated storage |
0.2 | 2 | 2012 | Gnothi: Separating Data and Metadata for Efficient and Available Storage Replication · USENIX ATC 2012 PRACTI Replication · NSDI 2006 |
Storage systems
distributed storage |
0.1 | 2 | 2009 | PADS: A Policy Architecture for Distributed Storage Systems · NSDI 2009 TAPER: Tiered Approach for Eliminating Redundancy in Replica Synchronization · FAST 2005 |
Distributed systems › consistency models
causal consistency |
0.1 | 1 | 2011 | Depot: Cloud Storage with Minimal Trust · ACM Trans. Comput. Syst. 2011 |
Distributed systems › fault tolerance
fault-tolerant distributed systems |
0.1 | 1 | 2011 | Depot: Cloud Storage with Minimal Trust · ACM Trans. Comput. Syst. 2011 |
Distributed systems › replication › replica control
quorum consensus |
0.1 | 1 | 2010 | Dual-Quorum: A Highly Available and Consistent Replication System for Edge Services · IEEE Trans. Dependable Secur. Comput. 2010 |
Distributed systems › consistency models
strong consistency |
0.1 | 1 | 2010 | Dual-Quorum: A Highly Available and Consistent Replication System for Edge Services · IEEE Trans. Dependable Secur. Comput. 2010 |
Distributed systems › replication
data replication |
0.1 | 2 | 2005 | Improving Availability and Performance with Application-Specific Data Replication · IEEE Trans. Knowl. Data Eng. 2005 Application specific data replication for edge services · WWW 2003 |
Distributed systems
edge computing |
0.1 | 2 | 2005 | Improving Availability and Performance with Application-Specific Data Replication · IEEE Trans. Knowl. Data Eng. 2005 Application specific data replication for edge services · WWW 2003 |
Query processing and optimization
approximate query processing |
0.1 | 1 | 2009 | Self-Tuning, Bandwidth-Aware Monitoring for Dynamic Data Streams · ICDE 2009 |
Algorithmic game theory and mechanism design
incentive mechanism |
0.1 | 1 | 2008 | FlightPath: Obedience vs. Choice in Cooperative Services · OSDI 2008 |
Storage systems › file systems
distributed file system |
0.1 | 5 | 2009 | Upright cluster services · SOSP 2009 Volume Leases for Consistency in Large-Scale Systems · IEEE Trans. Knowl. Data Eng. 1999 Serverless Network File Systems · SOSP 1995 |
Network security › anonymity networks
anonymous communication |
0.1 | 1 | 2006 | BAR Gossip · OSDI 2006 |
Distributed systems
gossip protocols |
0.1 | 1 | 2006 | BAR Gossip · OSDI 2006 |
Performance modeling and evaluation
storage performance evaluation |
0.1 | 1 | 2014 | Exalt: Empowering Researchers to Evaluate Large-Scale Storage Systems · NSDI 2014 |
Data stream processing › continuous query processing
continuous query monitoring |
0.1 | 1 | 2005 | INSIGHT: a distributed monitoring system for tracking continuous queries · SOSP 2005 |
Data stream processing
distributed monitoring |
0.1 | 1 | 2005 | INSIGHT: a distributed monitoring system for tracking continuous queries · SOSP 2005 |
Electronic design automation › logic synthesis › logic optimization
redundancy removal |
0.1 | 1 | 2005 | TAPER: Tiered Approach for Eliminating Redundancy in Replica Synchronization · FAST 2005 |
Distributed systems › replication › replica consistency
replica synchronization |
0.1 | 1 | 2005 | TAPER: Tiered Approach for Eliminating Redundancy in Replica Synchronization · FAST 2005 |
Distributed systems › peer-to-peer systems
distributed hash table |
0.0 | 1 | 2004 | A scalable distributed information management system · SIGCOMM 2004 |
Distributed systems › distributed information systems
distributed information management |
0.0 | 1 | 2004 | A scalable distributed information management system · SIGCOMM 2004 |
Distributed systems
peer-to-peer systems |
0.0 | 1 | 2004 | A scalable distributed information management system · SIGCOMM 2004 |
Distributed systems
distributed caching |
0.0 | 2 | 2002 | Coordinated Placement and Replacement for Large-Scale Distributed Caches · IEEE Trans. Knowl. Data Eng. 2002 Bandwidth constrained placement in a WAN · PODC 2001 |
Distributed systems
consistency models |
0.0 | 1 | 2003 | Application specific data replication for edge services · WWW 2003 |
Methods — techniques the papers use, named apart from their topics
simulation · 0.2protocol design · 0.2volume leases · 0.2analytical evaluation · 0.2temporal batching · 0.2hierarchical bandwidth allocation · 0.2arithmetic filtering · 0.2benchmarking · 0.2fault tolerance · 0.2game theory · 0.2aggregation with precision bounds · 0.1distributed object architecture · 0.1performance evaluation · 0.1byzantine fault tolerance protocols · 0.1gossip protocol · 0.1trace-based simulation · 0.0connectivity traces · 0.0resource isolation · 0.0
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2016 | Privacy Preserving Collaboration in Bring-Your-Own-AppsabstractEnterprise environments limit personal device usage for corporate data within a small set of enterprise provided apps or by using a whitelist of third-party apps. Both these options provide employees with limited app features, and a whitelist can be cumbersome to manage. Deepak Goel, Edmund L. Wong, Asim Kadav, Michael Dahlin |
SoCC | 5 |
| 2014 | Exalt: Empowering Researchers to Evaluate Large-Scale Storage Systems
Yang Wang 0009, Manos Kapritsos, Lara Schmidt, Lorenzo Alvisi, Michael Dahlin |
NSDI | 5 |
| 2014 | Lazy Means Smart: Reducing Repair Bandwidth Costs in Erasure-coded Distributed StorageabstractErasure coding schemes provide higher durability at lower storage cost, and thus constitute an attractive alternative to replication in distributed storage systems, in particular for storing rarely accessed "cold" data. These schemes, however, require an order of magnitude higher recovery bandwidth for maintaining a constant level of durability in the face of node failures. In this paper we propose lazy recovery, a technique to reduce recovery bandwidth demands down to the level of replicated storage. The key insight is that a careful adjustment of recovery rate substantially reduces recovery bandwidth, while keeping the impact on read performance and data durability low. We demonstrate the benefits of lazy recovery via extensive simulation using a realistic distributed storage configuration and published component failure parameters. For example, when applied to the commonly used RS(14, 10) code, lazy recovery reduces repair bandwidth by up to 76% even below replication, while increasing the amount of degraded stripes by 0.1 percentage points. Lazy recovery works well with a variety of erasure coding schemes, including the recently introduced bandwidth efficient codes, achieving up to a factor of 2 additional bandwidth savings. Mark Silberstein, Lakshmi Ganesh, Yang Wang 0009, Lorenzo Alvisi, Michael Dahlin |
SYSTOR | 5 |
| 2013 | Robustness in the Salus Scalable Block Store
Yang Wang 0009, Manos Kapritsos, Zuocheng Ren, Prince Mahajan, Jeevitha Kirubanandam, Lorenzo Alvisi, Michael Dahlin |
NSDI | 7 |
| 2013 | πBox: A Platform for Privacy-Preserving Apps
Edmund L. Wong, Deepak Goel, Michael Dahlin, Vitaly Shmatikov |
NSDI | 4 |
| 2012 | All about Eve: Execute-Verify Replication for Multi-Core Servers
Manos Kapritsos, Yang Wang 0009, Vivien Quéma, Allen Clement, Lorenzo Alvisi, Michael Dahlin |
OSDI | 6 |
| 2012 | Gnothi: Separating Data and Metadata for Efficient and Available Storage Replication
Yang Wang 0009, Lorenzo Alvisi, Michael Dahlin |
USENIX ATC | 3 |
| 2011 | Regret Freedom Isn't Free
Edmund L. Wong, Isaac Levy, Lorenzo Alvisi, Allen Clement, Michael Dahlin |
OPODIS | 5 |
| 2011 | Depot: Cloud Storage with Minimal TrustabstractThis article describes the design, implementation, and evaluation of Depot, a cloud storage system that minimizes trust assumptions. Depot tolerates buggy or malicious behavior byany numberof clients or servers, yet it provides safety and liveness guarantees to correct clients. Depot provides these guarantees using a two-layer architecture. First, Depot ensures that the updates observed by correct nodes are consistently ordered under Fork-Join-Causal consistency (FJC). FJC is a slight weakening of causal consistency that can be both safe and live despite faulty nodes. Second, Depot implements protocols that use this consistent ordering of updates to provide other desirable consistency, staleness, durability, and recovery properties. Our evaluation suggests that the costs of these guarantees are modest and that Depot can tolerate faults and maintain good availability, latency, overhead, and staleness even when significant faults occur. Prince Mahajan, Srinath Setty, Allen Clement, Lorenzo Alvisi, Michael Dahlin, Michael Walfish |
ACM Trans. Comput. Syst. | 6 |
| 2010 | Depot: Cloud Storage with Minimal Trust
Prince Mahajan, Srinath Setty, Allen Clement, Lorenzo Alvisi, Michael Dahlin, Michael Walfish |
OSDI | 6 |
| 2010 | Dual-Quorum: A Highly Available and Consistent Replication System for Edge ServicesabstractThis paper introduces dual-quorum replication, a novel data replication algorithm designed to support Internet edge services. Edge services allow clients to access Internet services via distributed edge servers that operate on a shared collection of underlying data. Although it is generally difficult to share data while providing high availability, good performance, and strong consistency, replication algorithms designed for specific access patterns can offer nearly ideal trade-offs among these metrics. In this paper, we focus on the key problem of sharing read/write data objects across a collection of edge servers when the references to each object 1) tend not to exhibit high concurrency across multiple nodes and 2) tend to exhibit bursts of read-dominated or write-dominated behavior. Dual-quorum replication combines volume leases and quorum-based techniques to achieve excellent availability, response time, and consistency for such workloads. In particular, through both analytical and experimental evaluations, we show that the dual-quorum protocol can (for the workloads of interest) approach the optimal performance and availability of Read-One/Write-All-Asynchronously (ROWA-A) epidemic algorithms without suffering the weak consistency guarantees and resulting design complexity inherent in ROWA-A systems. Michael Dahlin, Jiandan Zheng, Lorenzo Alvisi, Arun Iyengar |
IEEE Trans. Dependable Secur. Comput. | 2 |
| 2009 | Self-Tuning, Bandwidth-Aware Monitoring for Dynamic Data StreamsabstractWe present SMART, a self-tuning, bandwidth-aware monitoring system that maximizes result precision of continuous aggregate queries over dynamic data streams. While prior approaches minimize bandwidth cost under fixed precision constraints, they may still overload a monitoring system during traffic bursts. To facilitate practical deployment of monitoring systems, SMART therefore bounds the worst-case bandwidth cost for overload resilience. The primary challenge for SMART is how to dynamically select updates at each node to maximize query precision while keeping per-node monitoring bandwidth below a specified budget. To address this challenge, SMARTpsilas hierarchical algorithm (1) allocates bandwidth budgets in an ear-optimal manner to maximize global precision and (2) self-tunes bandwidth settings to improve precision under dynamic workloads. Our prototype implementation of SMART provides key solutions to (a) prioritize pending updates for multi-attribute queries, (b) build bounded fan-in, load-aware aggregation trees to improve accuracy, and (c) combine temporal batching with arithmetic filtering to reduce load and to quantify result staleness. Our evaluation using simulations and a network monitoring application shows that SMART incurs low overheads, improves accuracy by up to an order of magnitude compared to uniform bandwidth allocation, and performs close to the optimal algorithm under modest bandwidth budgets. Navendu Jain, Praveen Yalagandula, Michael Dahlin |
ICDE | 3 |
| 2009 | PADS: A Policy Architecture for Distributed Storage Systems
Nalini Moti Belaramani, Jiandan Zheng, Amol Nayate, Robert Soulé, Michael Dahlin, Robert Grimm 0001 |
NSDI | 5 |
| 2009 | Making Byzantine Fault Tolerant Systems Tolerate Byzantine Faults
Allen Clement, Edmund L. Wong, Lorenzo Alvisi, Michael Dahlin, Mirco Marchetti |
NSDI | 4 |
| 2009 | Upright cluster servicesabstractThe UpRight library seeks to make Byzantine fault tolerance (BFT) a simple and viable alternative to crash fault tolerance for a range of cluster services. We demonstrate UpRight by producing BFT versions of the Zookeeper lock service and the Hadoop Distributed File System (HDFS). Our design choices in UpRight favor simplifying adoption by existing applications; performance is a secondary concern. Despite these priorities, our BFT Zookeeper and BFT HDFS implementations have performance comparable with the originals while providing additional robustness. Allen Clement, Manos Kapritsos, Yang Wang 0009, Lorenzo Alvisi, Michael Dahlin, Taylor L. Riché |
SOSP | 6 |
| 2009 | Zyzzyva: Speculative Byzantine fault toleranceabstractA longstanding vision in distributed systems is to build reliable systems from unreliable components. An enticing formulation of this vision is Byzantine Fault-Tolerant (BFT) state machine replication, in which a group of servers collectively act as a correct server even if some of the servers misbehave or malfunction in arbitrary (“Byzantine”) ways. Despite this promise, practitioners hesitate to deploy BFT systems, at least partly because of the perception that BFT must impose high overheads. In this article, we present Zyzzyva, a protocol that uses speculation to reduce the cost of BFT replication. In Zyzzyva, replicas reply to a client's request without first running an expensive three-phase commit protocol to agree on the order to process requests. Instead, they optimistically adopt the order proposed by a primary server, process the request, and reply immediately to the client. If the primary is faulty, replicas can become temporarily inconsistent with one another, but clients detect inconsistencies, help correct replicas converge on a single total ordering of requests, and only rely on responses that are consistent with this total order. This approach allows Zyzzyva to reduce replication overheads to near their theoretical minima and to achieve throughputs of tens of thousands of requests per second, making BFT replication practical for a broad range of demanding services. Ramakrishna Kotla, Lorenzo Alvisi, Michael Dahlin, Allen Clement, Edmund L. Wong |
ACM Trans. Comput. Syst. | 3 |
| 2008 | BAR primerabstractByzantine and rational behaviors are increasingly recognized as unavoidable realities in todaypsilas cooperative services. Yet, how to design BAR-tolerant protocols and rigorously prove them strategy proof remains somewhat of a mystery: existing examples tend either to focus on unrealistically simple problems or to want in rigor. The goal of this paper is to demystify the process by presenting the full algorithmic development cycle that, starting from the classic synchronous repeated terminating reliable broadcast (R-TRB) problem statement, leads to a provably BAR-tolerant solution. We show i) how to express R-TRB as a game; ii) why the strategy corresponding to the optimal Byzantine fault tolerant algorithm of Dolev and strong does not guarantee safety when non-Byzantine players behave rationally; iii) how to derive a BAR-tolerant R-TRB protocol: iv) how to prove rigorously that the protocol ensures safety in the presence of non-Byzantine rational players. Allen Clement, Harry C. Li, Jeff Napper, Jean-Philippe Martin, Lorenzo Alvisi, Michael Dahlin |
DSN | 6 |
| 2008 | Network Imprecision: A New Consistency Metric for Scalable Monitoring
Navendu Jain, Prince Mahajan, Dmitry Kit, Praveen Yalagandula, Michael Dahlin |
OSDI | 5 |
| 2008 | FlightPath: Obedience vs. Choice in Cooperative Services
Harry C. Li, Allen Clement, Mirco Marchetti, Manos Kapritsos, Luke Robison, Lorenzo Alvisi, Michael Dahlin |
OSDI | 7 |
| 2007 | Machine Learning for On-Line Hardware Reconfiguration
Jonathan Wildstrom, Peter Stone 0001, Emmett Witchel, Michael Dahlin |
IJCAI | 4 |
| 2007 | Theory of BAR gamesabstractNo abstract available. Allen Clement, Jeff Napper, Harry C. Li, Jean-Philippe Martin, Lorenzo Alvisi, Michael Dahlin |
PODC | 6 |
| 2007 | Zyzzyva: speculative byzantine fault toleranceabstractWe present Zyzzyva, a protocol that uses speculation to reduce the cost and simplify the design of Byzantine fault tolerant state machine replication. In Zyzzyva, replicas respond to a client's request without first running an expensive three-phase commit protocol to reach agreement on the order in which the request must be processed. Instead, they optimistically adopt the order proposed by the primary and respond immediately to the client. Replicas can thus become temporarily inconsistent with one another, but clients detect inconsistencies, help correct replicas converge on a single total ordering of requests, and only rely on responses that are consistent with this total order. This approach allows Zyzzyva to reduce replication overheads to near their theoretical minimal. Ramakrishna Kotla, Lorenzo Alvisi, Michael Dahlin, Allen Clement, Edmund L. Wong |
SOSP | 3 |
| 2007 | SafeStore: A Durable and Practical Storage System
Ramakrishna Kotla, Lorenzo Alvisi, Michael Dahlin |
USENIX ATC | 3 |
| 2007 | STAR: Self-Tuning Aggregation for Scalable Monitoring
Navendu Jain, Michael Dahlin, Dmitry Kit, Prince Mahajan, Praveen Yalagandula |
VLDB | 2 |
| 2006 | PRACTI Replication
Nalini Moti Belaramani, Michael Dahlin, Amol Nayate, Arun Venkataramani, Praveen Yalagandula, Jiandan Zheng |
NSDI | 2 |
| 2006 | BAR Gossip
Harry C. Li, Allen Clement, Edmund L. Wong, Jeff Napper, Indrajit Roy 0001, Lorenzo Alvisi, Michael Dahlin |
OSDI | 7 |
| 2005 | TAPER: Tiered Approach for Eliminating Redundancy in Replica Synchronization
Navendu Jain, Michael Dahlin, Renu Tewari |
FAST | 2 |
| 2005 | Dual-Quorum Replication for Edge Services
Michael Dahlin, Jiandan Zheng, Lorenzo Alvisi, Arun Iyengar |
Middleware | 2 |
| 2005 | BAR fault tolerance for cooperative servicesabstractThis paper describes a general approach to constructing cooperative services that span multiple administrative domains. In such environments, protocols must tolerate both Byzantine behaviors when broken, misconfigured, or malicious nodes arbitrarily deviate from their specification and rational behaviors when selfish nodes deviate from their specification to increase their local benefit. The paper makes three contributions: (1) It introduces the BAR (Byzantine, Altruistic, Rational) model as a foundation for reasoning about cooperative services; (2) It proposes a general three-level architecture to reduce the complexity of building services under the BAR model; and (3) It describes an implementation of BAR-B the first cooperative backup service to tolerate both Byzantine users and an unbounded number of rational users. At the core of BAR-B is an asynchronous replicated state machine that provides the customary safety and liveness guarantees despite nodes exhibiting both Byzantine and rational behaviors. Our prototype provides acceptable performance for our application: our BAR-tolerant state machine executes 15 requests per second, and our BAR-B backup service can back up 100MB of data in under 4 minutes. Amitanand S. Aiyer, Lorenzo Alvisi, Allen Clement, Michael Dahlin, Jean-Philippe Martin, Carl Porth |
SOSP | 4 |
| 2005 | INSIGHT: a distributed monitoring system for tracking continuous queriesabstractA distributed monitoring framework can serve as an important building block for constructing large-scale data aggregation and continuous event monitoring applications, such as IP traffic monitoring (DDoS attacks), network anomaly detection (Internet worms), accounting and bandwidth provisioning (hot spots, flash crowds), sensor monitoring and control, and grid resource monitoring. At the core of these applications is a distributed query engine that aggregates information and performs continuous tracking of queries over collections of physically-distributed and rapidly-updating data streams. The underlying aim is to provide a global view of information in the system at a reasonable cost and within a specified precision bound. To achieve this objective, a distributed monitoring system should (a) scale to a large number of streams and query attributes, (b) incur minimal communication overhead for aggregating query results, (c) be time responsive for quickly identifying anomalies, and (d) be able to bound the inaccuracy of the computed value for the aggregate function. Navendu Jain, Praveen Yalagandula, Michael Dahlin |
SOSP | 3 |
| 2005 | Using Bloom Filters to Refine Web Search Results
Navendu Jain, Michael Dahlin, Renu Tewari |
WebDB | 2 |
| 2005 | Improving Availability and Performance with Application-Specific Data ReplicationabstractThe emerging edge services architecture promises to improve the availability and performance of Web services by replicating servers at geographically distributed sites. A key challenge in such systems is data replication and consistency, so that edge server code can manipulate shared data without suffering the availability and performance penalties that would be incurred by accessing a traditional centralized database. This work explores using a distributed object architecture to build an edge service data replication system for an e-commerce application, the TPC-W benchmark, which simulates an online bookstore. We take advantage of application-specific semantics to design distributed objects that each manages a specific subset of shared information using simple and effective consistency models. Our experimental results show that by slightly relaxing consistency within individual distributed objects, our application realizes both high availability and excellent performance. For example, in one experiment, we find that our object-based edge server system provides five times better response time over a traditional centralized cluster architecture and a factor of nine improvement over an edge service system that distributes code but retains a centralized database. Michael Dahlin, Amol Nayate, Jiandan Zheng, Arun Iyengar |
IEEE Trans. Knowl. Data Eng. | 2 |
| 2004 | High Throughput Byzantine Fault ToleranceabstractThis paper argues for a simple change to Byzantine fault tolerant (BFT) state machine replication libraries. Traditional BFT state machine replication techniques provide high availability and security but fail to provide high throughput. This limitation stems from the fundamental assumption of generalized state machine replication techniques that all replicas execute requests sequentially in the same total order to ensure consistency across replicas. We propose a high throughput Byzantine fault tolerant architecture that uses application-specific information to identify and concurrently execute independent requests. Our architecture thus provides a general way to exploit application parallelism in order to provide high throughput without compromising correctness. Although this approach is extremely simple, it yields dramatic practical benefits. When sufficient application concurrency and hardware resources exist, CBASE, our system prototype, provides orders of magnitude improvements in throughput over BASE, a traditional BFT architecture. CBASE-FS, a Byzantine fault tolerant file system that uses CBASE, achieves twice the throughput of BASE-FS for the IOZone micro-benchmarks even in a configuration with modest available hardware parallelism. Ramakrishna Kotla, Michael Dahlin |
DSN | 2 |
| 2004 | Transparent Information Dissemination
Amol Nayate, Michael Dahlin, Arun Iyengar |
Middleware | 2 |
| 2004 | A scalable distributed information management systemabstractWe present a Scalable Distributed Information Management System (SDIMS) that aggregates information about large-scale networked systems and that can serve as a basic building block for a broad range of large-scale distributed applications by providing detailed views of nearby information and summary views of global information. To serve as a basic building block, a SDIMS should have four properties: scalability to many nodes and attributes, flexibility to accommodate a broad range of applications, administrative isolation for security and availability, and robustness to node and network failures. We design, implement and evaluate a SDIMS that (1) leverages Distributed Hash Tables (DHT) to create scalable aggregation trees, (2) provides flexibility through a simple API that lets applications control propagation of reads and writes, (3) provides administrative isolation through simple extensions to current DHT algorithms, and (4) achieves robustness to node and network reconfigurations through lazy reaggregation, on-demand reaggregation, and tunable spatial replication. Through extensive simulations and micro-benchmark experiments, we observe that our system is an order of magnitude more scalable than existing approaches, achieves isolation properties at the cost of modestly increased read latency in comparison to flat DHTs, and gracefully handles failures. Praveen Yalagandula, Michael Dahlin |
SIGCOMM | 2 |
| 2003 | Separating agreement from execution for byzantine fault tolerant servicesabstractWe describe a new architecture for Byzantine fault tolerant state machine replication that separates agreement that orders requests from execution that processes requests. This separation yields two fundamental and practically significant advantages over previous architectures. First, it reduces replication costs because the new architecture can tolerate faults in up to half of the state machine replicas that execute requests. Previous systems can tolerate faults in at most a third of the combined agreement/state machine replicas. Second, separating agreement from execution allows a general privacy firewall architecture to protect confidentiality through replication. In contrast, replication in previous systems hurts confidentiality because exploiting the weakest replica can be sufficient to compromise the system. We have constructed a prototype and evaluated it running both microbenchmarks and an NFS server. Overall, we find that the architecture adds modest latencies to unreplicated systems and that its performance is competitive with existing Byzantine fault tolerant systems. Jian Yin 0002, Jean-Philippe Martin, Arun Venkataramani, Lorenzo Alvisi, Michael Dahlin |
SOSP | 5 |
| 2003 | Application specific data replication for edge servicesabstractThe emerging edge services architecture promises to improve the availability and performance of web services by replicating servers at geographically distributed sites. A key challenge in such systems is data replication and consistency so that edge server code can manipulate shared data without incurring the availability and performance penalties that would be incurred by accessing a traditional centralized database. This paper explores using a distributed object architecture to build an edge service system for an e-commerce application, an online bookstore represented by the TPC-W benchmark. We take advantage of application specific semantics to design distributed objects to manage a specific subset of shared information using simple and effective consistency models. Our experimental results show that by slightly relaxing consistency within individual distributed objects, we can build an edge service system that is highly available and efficient. For example, in one experiment we find that our object-based edge server system provides a factor of five improvement in response time over a traditional centralized cluster architecture and a factor of nine improvement over an edge service system that distributes code but retains a centralized database. Michael Dahlin, Amol Nayate, Jiandan Zheng, Arun Iyengar |
WWW | 2 |
| 2003 | Emulations between QSM, BSP and LogP: a framework for general-purpose parallel algorithm design
Vijaya Ramachandran, Brian Grayson, Michael Dahlin |
J. Parallel Distributed Comput. | 3 |
| 2003 | End-to-end WAN service availabilityabstractThis paper seeks to understand how network failures affect the availability of service delivery across wide-area networks (WANs) and to evaluate classes of techniques for improving end-to-end service availability. Using several large-scale connectivity traces, we develop a model of network unavailability that includes key parameters such as failure location and failure duration. We then use trace-based simulation to evaluate several classes of techniques for coping with network unavailability. We find that caching alone is seldom effective at insulating services from failures but that the combination of mobile extension code and prefetching can improve average unavailability by as much as an order of magnitude for classes of service whose semantics support disconnected operation. We find that routing-based techniques may provide significant improvements but that the improvements of many individual techniques are limited because they do not address all significant categories of network failures. By combining the techniques we examine, some systems may be able to reduce average unavailability by as much as one or two orders of magnitude. Michael Dahlin, Bharat Chandra, Amol Nayate |
IEEE/ACM Trans. Netw. | 1 |
| 2002 | Small Byzantine Quorum SystemsabstractIn this paper we present two protocols for asynchronous Byzantine quorum systems (BQS) built on top of reliable channels-one for self-verifying data and the other for any data. Our protocols tolerate f Byzantine failures with f fewer servers than existing solutions by eliminating nonessential work in the write protocol and by using read and write quorums of different sizes. Since engineering a reliable network layer on an unreliable network is difficult, two other possibilities must be explored. The first is to strengthen the model by allowing synchronous networks that use time-outs to identify failed links or machines. We consider running synchronous and asynchronous Byzantine quorum protocols over synchronous networks and conclude that, surprisingly, "self-timing" asynchronous Byzantine protocols may offer significant advantages for many synchronous networks when network time-outs are long. We show how to extend an existing Byzantine quorum protocol to eliminate its dependency on reliable networking and to handle message loss and retransmission explicitly. Jean-Philippe Martin, Lorenzo Alvisi, Michael Dahlin |
DSN | 3 |
| 2002 | TCP Nice: A Mechanism for Background Transfers
Arun Venkataramani, Ravi Kokku, Michael Dahlin |
OSDI | 3 |
| 2002 | Minimal Byzantine Storage
Jean-Philippe Martin, Lorenzo Alvisi, Michael Dahlin |
DISC | 3 |
| 2002 | The potential costs and benefits of long-term prefetching for content distribution
Arun Venkataramani, Praveen Yalagandula, Ravi Kokku, Sadia Sharif, Michael Dahlin |
Comput. Commun. | 5 |
| 2002 | Coordinated Placement and Replacement for Large-Scale Distributed CachesabstractIn a large-scale information system such as a digital library or the Web, a set of distributed caches can improve their effectiveness by coordinating their data placement decisions. Using simulation, we examine three practical cooperative placement algorithms, including one that is provably close to optimal, and we compare these algorithms to the optimal placement algorithm and several cooperative and noncooperative replacement algorithms. We draw five conclusions from these experiments: 1) cooperative placement can significantly improve performance compared to local replacement algorithms, particularly when the size of individual caches is limited compared to the universe of objects; 2) although the amortizing placement algorithm is only guaranteed to be within 14 times the optimal, in practice it seems to provide an excellent approximation of the optimal; 3) in a cooperative caching scenario, the recent greedy-dual local replacement algorithm performs much better than the other local replacement algorithms; 4) our hierarchical-greedy-dual replacement algorithm yields further improvements over the greedy-dual algorithm especially when there are idle caches in the system; and 5) a key challenge to coordinated placement algorithms is generating good predictions of access patterns based on past accesses. Madhukar R. Korupolu, Michael Dahlin |
IEEE Trans. Knowl. Data Eng. | 2 |
| 2002 | Engineering web cache consistencyabstractServer-driven consistency protocols can reduce read latency and improve data freshness for a given network and server overhead, compared to the traditional consistency protocols that rely on client polling. Server-driven consistency protocols appear particularly attractive for large-scale dynamic Web workloads because dynamically generated data can change rapidly and unpredictably. However, there have been few reports on engineering server-driven consistency for such workloads. This article reports our experience in engineering server-driven consistency for a sporting and event Web site hosted by IBM, one of the most popular sites on the Internet for the duration of the event. We also examine an e-commerce site for a national retail store. Our study focuses on scalability and cachability of dynamic content. To assess scalability, we measure both the amount of state that a server needs to maintain to ensure consistency and the bursts of load in sending out invalidation messages when a popular object is modified. We find that server-driven protocols can cap the size of the server's state to a given amount without significant performance costs, and can smooth the bursts of load with minimal impact on the consistency guarantees. To improve performance, we systematically investigate several design issues for which prior research has suggested widely different solutions, including whether servers should send invalidations to idle clients. Finally, we quantify the performance impact of caching dynamic data with server-driven consistency protocols and the benefits of server-driven consistency protocols for large-scale dynamic Web services. We find that (i) caching dynamically generated data can increase cache hit rates by up to 10%, compared to the systems that do not cache dynamically generated data; and (ii) server-driven consistency protocols can increase cache hit rates by a factor of 1.5-3 for large-scale dynamic Web services, compared to client polling protocols. We have implemented a prototype of a server-driven consistency protocol based on our findings by augmenting the popular Squid cache. Jian Yin 0002, Lorenzo Alvisi, Michael Dahlin, Arun Iyengar |
ACM Trans. Internet Techn. | 3 |
| 2001 | Bandwidth constrained placement in a WANabstractIn this paper, we examine the bandwidth-constrained placement problem, focusing on trade-offs appropriate for wide area network (WAN) environments. The goal is to place copies of objects at a collection of distributed caches to minimize expected access times from distributed clients to those objects subject to a maximum bandwidth constraint at each cache. We develop a simple algorithm to generate a bandwidth-constrained placement by hierarchically refining an initial per-cache greedy placement. We prove that this hierarchical algorithm generates a placement whose expected access time is within a constant factor of the optimal placement's expected access time. We then proceed to extend this algorithm to compute close to optimal placement strategies for dynamic environments. Arun Venkataramani, Phoebe Weidmann, Michael Dahlin |
PODC | 3 |
| 2001 | Resource management for scalable disconnected access to Web servicesabstractDisconnected operation, in which a client accesses a service without relying on network connectivity, is crucial for improving availability, supporting mobility, and providing responsive performance. Because manyweb services are not cachable, disconnected access to web services may require mobile service code to execute in client caches. Unfortunately, (a) this code is untrusted, (b) this code mayhave nearly limitless resource demands due to prefetching, and (c) a large number of competing code modules must coexist. Thus, resource managementisakey problem both for preventing denial of service attacks and for providing good performance across many services. This paper addresses the feasibility of meeting the resource management needs of an environment where service code is shipped to clients, proxies, or content distribution intermediaries. It rst examines the requirements of such a system and then develops a resource-management strategy to meet these requirements by (a) providing isolation across services to prevent denial of service attacks, (b) automatically providing appropriate allocations to dierent services to provide good global performance, and (c) requiring no hand tuning across a wide range of system congurations and workloads. 1. Bharat Chandra, Michael Dahlin, Amjad-Ali Khoja, Amol Nayate, Asim Razzaq, Anil Sewani |
WWW | 2 |
| 2001 | Engineering server-driven consistency for large scale dynamic Web servicesabstractMany researchers have shown that server-driven consistency protocols can potentially reduce read latency. Server-driven consistency protocols are particularly attractive for large-scale dynamic web workloads because dynamically generated data can change rapidly and unpredictably. However, there have been no reports on engineering server-driven consistency for such a workload. This paper reports our experience in engineering server-driven consistency for a Sporting and Eventweb site hosted by IBM, one of the most popular web sites on the Internet for the duration of the event. Our study focuses on scalability and cachability of dynamic content. To assess scalability, we measure both the amount of state that a server needs to maintain to ensure consistency and the bursts of load that a server sustains to send out invalidation messages when a popular object is modified. We find that it is possible to limit the size of the server's state without significant performance costs and that bursts of load can be smoothed out with minimal impact on the consistency guarantees. To improve performance, we systematically investigate several design issues for which prior research has suggested widely different solutions, including how long servers should send invalidations to idle clients. Finally, we quantify the performance impact of caching dynamic data with server-driven consistency protocols and find that it can reduce read latency by more than 10%. We have implemented a prototype of a server-driven consistency protocol based on our findings on top of the popular Squid cache. Jian Yin 0002, Lorenzo Alvisi, Michael Dahlin, Arun Iyengar |
WWW | 3 |
| 2000 | Interpreting Stale Load InformationabstractIn this paper, we examine the problem of balancing load in a large-scale distributed system when information about server loads may be stale. It is well-known that sending each request to the machine with the apparent lowest load can behave badly in such systems, yet this technique is common in practice. Other systems use round-robin or random selection algorithms that entirely ignore load information or that only use a small subset of the load information. Rather than risk extremely bad performance on one hand or ignore the chance to use load information to improve performance on the other, we develop strategies that interpret load information based on its age. Through simulation, we examine several simple algorithms that use such load interpretation strategies under a range of workloads. Our experiments suggest that by properly interpreting load information, systems can: 1) match the performance of the most aggressive algorithms when load information is fresh relative to the job arrival rate, 2) outperform the best of the other algorithms we examine by as much as 60 percent when information is moderately old, 3) significantly outperform random load distribution when information is older still, and 4) avoid pathological behavior even when information is extremely old. Michael Dahlin |
IEEE Trans. Parallel Distributed Syst. | 1 |
| 1999 | Interpreting Stale Load InformationabstractIn this paper we examine the problem of balancing load in a large-scale distributed system when information about server loads may be stale. It is well known that sending each request to the machine with the apparent lowest load can behave badly in such systems, yet this technique is common in practice. Other systems use round-robin or random selection algorithms that entirely ignore load information or that only use a small subset of the load information. Rather than risk extremely bad performance on one hand or ignore the chance to use load information to improve performance on the other, we develop strategies that interpret load information based on its age. Through simulation, we examine several simple algorithms that use such load interpretation strategies under a range of workloads. Our experiments suggest that by properly interpreting load information, systems can (1) match the performance of the most aggressive algorithms when load information is fresh relative to the job arrival rate, (2) outperform the best of the other algorithms we examine by as much as 60% when information is moderately old, (3) significantly outperform random load distribution when information is older still, and (4) avoid pathological behavior even when information is extremely old. Michael Dahlin |
ICDCS | 1 |
| 1999 | Design Considerations for Distributed Caching on the InternetabstractWe describe the design and implementation of an integrated architecture for cache systems that scale to hundreds or thousands of caches with thousands to millions of users. Rather than simply try to maximize hit rates, we take an end-to-end approach to improving response time by also considering hit times and miss times. We begin by studying several Internet caches and workloads, and we derive three core design principles for large scale distributed caches: minimize the number of hops to locate and access data on both hits and misses; share data among many users and scale to many caches; and cache data close to clients. Our strategies for addressing these issues are built around a scalable, high-performance data-location service that tracks where objects are replicated. We describe how to construct such a service and how to use this service to provide direct access to remote data and push-based data replication. We evaluate our system through trace-driven simulation and find that these strategies together provide response time speedups of 1.27 to 2.43 compared to a traditional three-level cache hierarchy for a range of trace workloads and simulated environments. Renu Tewari, Michael Dahlin, Harrick M. Vin, Jonathan S. Kay |
ICDCS | 2 |
| 1999 | Emulations Between QSM, BSP, and LogP: A Framework for General-Purpose Parallel Algorithm Design
Vijaya Ramachandran, Brian Grayson, Michael Dahlin |
SODA | 3 |
| 1999 | Volume Leases for Consistency in Large-Scale SystemsabstractThis article introduces volume leases as a mechanism for providing server-driven cache consistency for large-scale, geographically distributed networks. Volume leases retain the good performance, fault tolerance, and server scalability of the semantically weaker client-driven protocols that are now used on the Web. Volume leases are a variation of object leases, which were originally designed for distributed file systems. However, whereas traditional object leases amortize overheads over long lease periods, volume leases exploit spatial locality to amortize overheads across multiple objects in a volume. This approach allows systems to maintain good write performance even in the presence of failures. Using trace-driven simulation, we compare three volume lease algorithms against four existing cache consistency algorithms and show that our new algorithms provide strong consistency while maintaining scalability and fault-tolerance. For a trace-based workload of Web accesses, we find that volumes can reduce message traffic at servers by 40 percent compared to a standard lease algorithm, and that volumes can considerably reduce the peak load at servers when popular objects are modified. Jian Yin 0002, Lorenzo Alvisi, Michael Dahlin, Calvin Lin |
IEEE Trans. Knowl. Data Eng. | 3 |
| 1998 | WebOS: Operating System Services for Wide Area ApplicationsabstractDemonstrates the power of providing a common set of operating system services to wide-area applications, including mechanisms for naming, persistent storage, remote process execution, resource management, authentication and security. On a single machine, application developers can rely on the local operating system to provide these abstractions. In the wide area, however, application developers are forced to build these abstractions themselves or to do without. This ad-hoc approach often results in individual programmers implementing non-optimal solutions, wasting both programmer effort and system resources. To address these problems, we are building a system, WebOS, that provides the basic operating systems services needed to build applications that are geographically distributed, highly available, incrementally scalable and dynamically reconfigurable. Experience with a number of applications developed under WebOS indicates that it simplifies system development and improves resource utilization. In particular, we use WebOS to implement Rent-A-Server to provide dynamic replication of overloaded Web services across the wide area in response to client demands. Amin Vahdat, Thomas E. Anderson, Michael Dahlin, Eshwar Belani, David E. Culler, Paul Eastham, Chad Yoshikawa |
HPDC | 3 |
| 1998 | Using Leases to Support Server-Driven Consistency in Large-Scale SystemsabstractThe paper introduces volume leases as a mechanism for providing cache consistency for large scale, geographically distributed networks. Volume leases are a variation of leases, which were originally designed for distributed file systems. Using trace driven simulation, we compare two new algorithms against four existing cache consistency algorithms and show that our new algorithms provide strong consistency while maintaining scalability and fault tolerance. For a trace based workload of Web accesses, we find that volumes can reduce message traffic at servers by 40% compared to a standard lease algorithm, and that volumes can considerably reduce the peak load at servers when popular objects are modified. Jian Yin 0002, Lorenzo Alvisi, Michael Dahlin, Calvin Lin |
ICDCS | 3 |
| 1998 | The CRISIS Wide Area Security Architecture
Eshwar Belani, Amin Vahdat, Thomas E. Anderson, Michael Dahlin |
USENIX Security Symposium | 4 |
| 1996 | Serverless Network File SystemsabstractWe propose a new paradigm for network file system design:serverless network file systems. While traditional network file systems rely on a central server machine, a serverless system utilizes workstations cooperating as peers to provide all file system services. Any machine in the system can store, cache, or control any block of data. Our approach uses this location independence, in combination with fast local area networks, to provide better performance and scalability than traditional file systems. Furthermore, because any machine in the system can assume the responsibilities of a failed component, our serverless design also provides high availability via redundatn data storage. To demonstrate our approach, we have implemented a prototype serverless network file system called xFS. Preliminary performance measurements suggest that our architecture achieves its goal of scalability. For instance, in a 32-node xFS system with 32 active clients, each client receives nearly as much read or write throughput as it would see if it were the only active client. Thomas E. Anderson, Michael Dahlin, Jeanna Matthews, David A. Patterson 0001, Drew S. Roselli, Randolph Y. Wang |
ACM Trans. Comput. Syst. | 2 |
| 1995 | Serverless Network File SystemsabstractArticle Serverless network file systems Share on Authors: T. E. Anderson Computer Science Division, University of California at Berkeley Computer Science Division, University of California at BerkeleyView Profile , M. D. Dahlin Computer Science Division, University of California at Berkeley Computer Science Division, University of California at BerkeleyView Profile , J. M. Neefe Computer Science Division, University of California at Berkeley Computer Science Division, University of California at BerkeleyView Profile , D. A. Patterson Computer Science Division, University of California at Berkeley Computer Science Division, University of California at BerkeleyView Profile , D. S. Roselli Computer Science Division, University of California at Berkeley Computer Science Division, University of California at BerkeleyView Profile , R. Y. Wang Computer Science Division, University of California at Berkeley Computer Science Division, University of California at BerkeleyView Profile Authors Info & Claims SOSP '95: Proceedings of the fifteenth ACM symposium on Operating systems principlesDecember 1995 Pages 109–126https://doi.org/10.1145/224056.224066Online:03 December 1995Publication History 255citation3,156DownloadsMetricsTotal Citations255Total Downloads3,156Last 12 Months71Last 6 weeks6 Get Citation AlertsNew Citation Alert added!This alert has been successfully added and will be sent to:You will be notified whenever a record that you have chosen has been cited.To manage your alert preferences, click on the button below.Manage my AlertsNew Citation Alert!Please log in to your account Save to BinderSave to BinderCreate a New BinderNameCancelCreateExport CitationPublisher SiteGet Access Thomas E. Anderson, Michael Dahlin, Jeanna Matthews, David A. Patterson 0001, Drew S. Roselli, Randolph Y. Wang |
SOSP | 2 |
| 1994 | Cooperative Caching: Using Remote Client Memory to Improve File System Performance
Michael Dahlin, Randolph Y. Wang, Thomas E. Anderson, David A. Patterson 0001 |
OSDI | 1 |
| 1994 | A Quantitative Analysis of Cache Policies for Scalable Network File SystemsabstractCurrent network file system protocols rely heavily on a central server to coordinate file activity among client workstations. This central server can become a bottleneck that limits scalability for environments with large numbers of clients. In central server systems such as NFS and AFS, all client writes, cache misses, and coherence messages are handled by the server. To keep up with this workload, expensive server machines are needed, configured with high-performance CPUs, memory systems, and I/O channels. Since the server stores all data, it must be physically capable of connecting to many disks. This reliance on a central server also makes current systems inappropriate for wide area network use where the network bandwidth to the server may be limited. Michael Dahlin, Clifford Mather, Randolph Y. Wang, Thomas E. Anderson, David A. Patterson 0001 |
SIGMETRICS | 1 |