EDBT 2026 Demo / reviewers in the wild / expert
Lidong Zhou
dblp:z/LidongZhou
· DBLP profile ↗
75ranked-venue papers
5as first author
17since 2021 · last 2026
0000-0002-7258-3116ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 31 · 3 first-author · 5 since 2021Software engineering, systems software and programming languages · 31 · 10 since 2021Computer networks · 6 · 1 since 2021Security and privacy · 4 · 2 first-authorArtificial intelligence and machine learning · 3 · 2 since 2021Databases, data management, data science and information retrieval · 1Graphics, computer vision, multimedia, augmented reality and games · 1
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | SuperBench: A Proactive Validation System for Improving Reliability of Cloud AI InfrastructureabstractReliability in cloud AI infrastructure is crucial for cloud service providers, prompting the widespread use of hardware redundancies. However, these redundancies can inadvertently lead to hidden degradation, known as “gray failure”, for AI workloads, significantly affecting end-to-end performance and concealing performance issues, which complicates root cause analysis for failures and regressions. We introduce SuperBench, a proactive validation system for AI infrastructure that mitigates hidden degradation caused by hardware redundancies and enhances overall reliability. SuperBench features a comprehensive benchmark suite, capable of evaluating individual hardware components and representing most real AI workloads. It comprises a Validator that learns benchmark criteria to pinpoint defective components clearly. Additionally, SuperBench incorporates a Selector to balance validation time and issue-related penalties, enabling optimal timing for validation execution with a tailored subset of benchmarks. Through testbed evaluation and simulation, we demonstrate that SuperBench can increase the mean time between incidents by up to 22.61×. SuperBench has been successfully deployed in Azure production, validating hundreds of thousands of GPUs every year. Yifan Xiong 0001, Ziyue Yang 0002, Guoshuai Zhao 0001, Dong Zhong, Boris Pinzur, Jie Zhang 0048, Yang Wang 0053, Hossein Pourreza, Jeff Baxter, Kushal Datta, Prabhat Ram, Luke Melton, Joe Chau, Peng Cheng 0005, Yongqiang Xiong, Lidong Zhou |
ACM Trans. Comput. Syst. | 20 |
| 2026 | Introduction to the Special Section on USENIX OSDI 2025
Yuanyuan Zhou 0001, Lidong Zhou |
ACM Trans. Storage | 2 |
| 2025 | Automated Proof Generation for Rust Code via Self-EvolutionabstractEnsuring correctness is crucial for code generation. Formal verification offers a
definitive assurance of correctness, but demands substantial human effort in proof
construction and hence raises a pressing need for automation. The primary obsta-
cle lies in the severe lack of data—there is much fewer proofs than code snippets
for Large Language Models (LLMs) to train upon. In this paper, we introduce
SAFE, a framework that overcomes the lack of human-written proofs to enable
automated proof generation of Rust code. SAFE establishes a self-evolving cycle
where data synthesis and fine-tuning collaborate to enhance the model capability,
leveraging the definitive power of a symbolic verifier in telling correct proofs from
incorrect ones. SAFE also re-purposes the large number of synthesized incorrect
proofs to train the self-debugging capability of the fine-tuned models, empowering
them to fix incorrect proofs based on the verifier’s feedback. SAFE demonstrates
superior efficiency and precision compared to GPT-4o. Through tens of thousands
of synthesized proofs and the self-debugging mechanism, we improve the capa-
bility of open-source models, initially unacquainted with formal verification, to
automatically write proofs for Rust code. This advancement leads to a signifi-
cant improvement in performance, achieving a 52.52% accuracy rate in a bench-
mark crafted by human experts, a significant leap over GPT-4o’s performance of
14.39%. Tianyu Chen 0006, Shan Lu 0001, Yeyun Gong, Chenyuan Yang, Xuheng Li, Md Rakib Hossain Misu, Hao Yu 0016, Nan Duan 0001, Peng Cheng 0005, Fan Yang 0024, Shuvendu K. Lahiri, Tao Xie 0001, Lidong Zhou |
ICLR | 14 |
| 2024 | Amanda: Unified Instrumentation Framework for Deep Neural NetworksabstractThe success of deep neural networks (DNNs) has sparked efforts to analyze (e.g., tracing) and optimize (e.g., pruning) them. These tasks have specific requirements and ad-hoc implementations in current execution backends like TensorFlow/PyTorch, which require developers to manage fragmented interfaces and adapt their codes to diverse models. In this study, we propose a new framework called Amanda to streamline the development of these tasks. We formalize the implementation of these tasks as neural network instrumentation, which involves introducing instrumentation into the operator level of DNNs. This allows us to abstract DNN analysis and optimization tasks as instrumentation tools on various DNN models. We build Amanda with two levels of APIs to achieve a unified, extensible, and efficient instrumentation design. The user-level API provides a unified operator-grained instrumentation API for different backends. Meanwhile, internally, we design a set of callback-centric APIs for managing and optimizing the execution of original and instrumentation codes in different backends. Through these design principles, the Amanda framework can accommodate a broad spectrum of use cases, such as tracing, profiling, pruning, and quantization, across different backends (e.g., TensorFlow/PyTorch) and execution modes (graph/eager mode). Moreover, our efficient execution management ensures that the performance overhead is typically kept within 5%. Yue Guan 0003, Yuxian Qiu, Jingwen Leng, Fan Yang 0024, Shuo Yu 0006, Yunxin Liu 0001, Yu Feng 0007, Yuhao Zhu 0001, Lidong Zhou, Yun Liang 0001, Chen Zhang 0001, Chao Li 0009, Minyi Guo |
ASPLOS (1) | 9 |
| 2024 | nnScaler: Constraint-Guided Parallelization Plan Generation for Deep Learning Training
Youshan Miao, Quanlu Zhang, Fan Yang 0024, Cheng Li 0001, Saeed Maleki, Yilei Yang, Weijiang Xu, Mao Yang 0004, Lidong Zhou |
OSDI | 14 |
| 2024 | SuperBench: Improving Cloud AI Infrastructure Reliability with Proactive Validation
Yifan Xiong 0001, Ziyue Yang 0002, Guoshuai Zhao 0001, Dong Zhong, Boris Pinzur, Jie Zhang 0048, Yang Wang 0053, Hossein Pourreza, Jeff Baxter, Kushal Datta, Prabhat Ram, Luke Melton, Joe Chau, Peng Cheng 0005, Yongqiang Xiong, Lidong Zhou |
USENIX ATC | 20 |
| 2023 | SiloD: A Co-design of Caching and Scheduling for Deep Learning ClustersabstractDeep learning training on cloud platforms usually follows the tradition of the separation of storage and computing. The training executes on a compute cluster equipped with GPUs/TPUs while reading data from a separate cluster hosting the storage service. To alleviate the potential bottleneck, a training cluster usually leverages its local storage as a cache to reduce the remote IO from the storage cluster. However, existing deep learning schedulers do not manage storage resources thus fail to consider the diverse caching effects across different training jobs. This could degrade scheduling quality significantly. Zhenhua Han, Zhi Yang 0001, Quanlu Zhang, Mingxia Li, Fan Yang 0024, Qianxi Zhang, Binyang Li, Yuqing Yang 0001, Lili Qiu, Lidong Zhou |
EuroSys | 12 |
| 2023 | On Modular Learning of Distributed Systems for Predicting End-to-End Latency
Chieh-Jan Mike Liang, Zilin Fang, Yuqing Xie 0005, Fan Yang 0024, Zhao Lucis Li, Li Lyna Zhang, Mao Yang 0004, Lidong Zhou |
NSDI | 8 |
| 2023 | Welder: Scheduling Deep Learning Memory Access via Tile-graph
Yining Shi 0001, Zhi Yang 0001, Jilong Xue, Lingxiao Ma, Yuqing Xia, Ziming Miao, Yuxiao Guo 0001, Fan Yang 0024, Lidong Zhou |
OSDI | 9 |
| 2023 | Optimizing Dynamic Neural Networks with Brainstorm
Weihao Cui, Zhenhua Han, Lingji Ouyang, Yichuan Wang 0002, Ningxin Zheng, Lingxiao Ma, Yuqing Yang 0001, Fan Yang 0024, Jilong Xue, Lili Qiu, Lidong Zhou, Quan Chen 0002, Haisheng Tan, Minyi Guo |
OSDI | 11 |
| 2023 | VBASE: Unifying Online Vector Similarity Search and Relational Queries via Relaxed Monotonicity
Qianxi Zhang, Shuotao Xu, Qi Chen 0009, Guoxin Sui, Jiadong Xie 0002, Zhizhen Cai, Yaoqi Chen, Yinxuan He, Yuqing Yang 0001, Fan Yang 0024, Mao Yang 0004, Lidong Zhou |
OSDI | 12 |
| 2023 | PIT: Optimization of Dynamic Sparse Deep Learning Models via Permutation Invariant TransformationabstractDynamic sparsity, where the sparsity patterns are unknown until runtime, poses a significant challenge to deep learning. The state-of-the-art sparsity-aware deep learning solutions are restricted to pre-defined, static sparsity patterns due to significant overheads associated with preprocessing. Efficient execution of dynamic sparse computation often faces the misalignment between the GPU-friendly tile configuration for efficient execution and the sparsity-aware tile shape that minimizes coverage wastes (non-zero values in tensor). Ningxin Zheng, Huiqiang Jiang, Quanlu Zhang, Zhenhua Han, Lingxiao Ma, Yuqing Yang 0001, Fan Yang 0024, Chengruidong Zhang, Lili Qiu, Mao Yang 0004, Lidong Zhou |
SOSP | 11 |
| 2022 | SparTA: Deep-Learning Model Sparsity via Tensor-with-Sparsity-Attribute
Ningxin Zheng, Quanlu Zhang, Lingxiao Ma, Yuqing Yang 0001, Fan Yang 0024, Yang Wang 0053, Mao Yang 0004, Lidong Zhou |
OSDI | 9 |
| 2022 | ROLLER: Fast and Efficient Tensor Compilation for Deep Learning
Hongyu Zhu 0003, Yijia Diao, Shanbin Ke, Chen Zhang 0001, Jilong Xue, Lingxiao Ma, Yuqing Xia, Fan Yang 0024, Mao Yang 0004, Lidong Zhou, Asaf Cidon, Gennady Pekhimenko |
OSDI | 13 |
| 2021 | LayoutLMv2: Multi-modal Pre-training for Visually-rich Document UnderstandingabstractYang Xu, Yiheng Xu, Tengchao Lv, Lei Cui, Furu Wei, Guoxin Wang, Yijuan Lu, Dinei Florencio, Cha Zhang, Wanxiang Che, Min Zhang, Lidong Zhou. Proceedings of the 59th Annual Meeting of the Association for Computational Linguistics and the 11th International Joint Conference on Natural Language Processing (Volume 1: Long Papers). 2021. Yang Xu 0049, Yiheng Xu, Tengchao Lv, Lei Cui 0001, Furu Wei, Yijuan Lu, Dinei A. F. Florêncio, Cha Zhang, Wanxiang Che, Min Zhang 0005, Lidong Zhou |
ACL/IJCNLP (1) | 12 |
| 2021 | Forerunner: Constraint-based Speculative Transaction Execution for EthereumabstractEthereum is an emerging distributed computing platform that supports a decentralized replicated virtual machine at a large scale. Transactions in Ethereum are specified in smart contracts, disseminated through broadcast, accepted into the chain of blocks, and then executed on each node. In this new Dissemination-Consensus-Execution (DiCE) paradigm, the time interval between when a transaction is known (during the dissemination phase) to when the transaction is executed (after the consensus phase) offers a window of opportunity to accelerate transaction processing through speculative execution. However, the traditional speculative execution, which hinges on the ability to predict the future accurately, is inadequate because of DiCE's many-future nature. Zhongxin Guo, Runhuai Li, Shuo Chen 0001, Lidong Zhou, Yajin Zhou, Xian Zhang 0001 |
SOSP | 5 |
| 2021 | Geometric Partitioning: Explore the Boundary of Optimal Erasure Code RepairabstractErasure coding is widely used in building reliable distributed object storage systems despite its high repair cost. Regenerating codes are a special class of erasure codes, which are proposed to minimize the amount of data needed for repair. In this paper, we assess how optimal repair can help to improve object storage systems, and we find that regenerating codes present unique challenges: regenerating codes repair at the granularity of chunks instead of bytes, and the choice of chunk size leads to the tension between streamed degraded read time and repair throughput. Yingdi Shan, Kang Chen 0001, Tuoyu Gong, Lidong Zhou, Tai Zhou, Yongwei Wu 0001 |
SOSP | 4 |
| 2020 | TextNAS: A Neural Architecture Search Space Tailored for Text RepresentationabstractLearning text representation is crucial for text classification and other language related tasks. There are a diverse set of text representation networks in the literature, and how to find the optimal one is a non-trivial problem. Recently, the emerging Neural Architecture Search (NAS) techniques have demonstrated good potential to solve the problem. Nevertheless, most of the existing works of NAS focus on the search algorithms and pay little attention to the search space. In this paper, we argue that the search space is also an important human prior to the success of NAS in different applications. Thus, we propose a novel search space tailored for text representation. Through automatic search, the discovered network architecture outperforms state-of-the-art models on various public datasets on text classification and natural language inference tasks. Furthermore, some of the design principles found in the automatic network agree well with human intuition. Yujing Wang 0002, Yaming Yang 0001, Jing Bai 0010, Ce Zhang 0001, Guinan Su, Xiaoyu Kou, Yunhai Tong, Mao Yang 0004, Lidong Zhou |
AAAI | 10 |
| 2020 | Rammer: Enabling Holistic Deep Learning Compiler Optimizations with rTasks
Lingxiao Ma, Zhi Yang 0001, Jilong Xue, Youshan Miao, Wenxiang Hu, Fan Yang 0024, Lidong Zhou |
OSDI | 10 |
| 2020 | Retiarii: A Deep Learning Exploratory-Training Framework
Quanlu Zhang, Zhenhua Han, Fan Yang 0024, Yuge Zhang, Mao Yang 0004, Lidong Zhou |
OSDI | 7 |
| 2020 | Byzantine Ordered Consensus without Byzantine Oligarchy
Srinath Setty, Qi Chen 0009, Lidong Zhou, Lorenzo Alvisi |
OSDI | 4 |
| 2020 | HiveD: Sharing a GPU Cluster for Deep Learning with Guarantees
Zhenhua Han, Zhi Yang 0001, Quanlu Zhang, Fan Yang 0024, Lidong Zhou, Mao Yang 0004, Francis C. M. Lau 0001, Yifan Xiong 0001 |
OSDI | 6 |
| 2020 | AutoSys: The Design and Operation of Learning-Augmented Systems
Chieh-Jan Mike Liang, Hui Xue 0004, Mao Yang 0004, Lidong Zhou, Lifei Zhu, Zhao Lucis Li, Qi Chen 0009, Quanlu Zhang, Chuanjie Liu, Wenjun Dai |
USENIX ATC | 4 |
| 2020 | Distributed Graph Computation Meets Machine LearningabstractTuX2is a new distributed graph engine that bridges graph computation and distributed machine learning.TuX2inherits the benefits of elegant graph computation model, efficient graph layout, and balanced parallelism to scale to billion-edge graphs, while extended and optimized for distributed machine learning to support heterogeneity in data model, Stale Synchronous Parallel in scheduling, and a new Mini-batch, Exchange, GlobalSync, and Apply (MEGA) model for programming.TuX2further introduces a hybrid vertex-cut graph optimization and supports various consistency models in fault tolerance for machine learning. We have developed a set of representative distributed machine learning algorithms inTuX2, covering both supervised and unsupervised learning. Compared to the implementations on distributed machine learning platforms, writing those algorithms inTuX2takes only about 25 percent of the code: our graph computation model hides the detailed management of data layout, partitioning, and parallelism from developers. The extensive evaluation ofTuX2, using large datasets with up to 64 billion of edges, shows thatTuX2outperforms PowerGraph/PowerLyra, the state-of-the-art distributed graph engines, by an order of magnitude, while beating two state-of-the-art distributed machine learning systems by at least 60 percent. Wencong Xiao, Jilong Xue, Youshan Miao, Ming Wu 0007, Wei Li 0022, Lidong Zhou |
IEEE Trans. Parallel Distributed Syst. | 8 |
| 2019 | Astra: Exploiting Predictability to Optimize Deep LearningabstractWe present Astra, a compilation and execution framework that optimizes execution of a deep learning training job. Instead of treating the computation as a generic data flow graph, Astra exploits domain knowledge about deep learning to adopt a custom approach to compiler optimization. The key insight in Astra is to exploit the unique repetitiveness and predictability of a deep learning job, to perform online exploration of the optimization state space in a work-conserving manner while making progress on the training job. This dynamic state space exploration in Astra uses lightweight profiling and indexing of profile data, coupled with several techniques to prune the exploration state space. Effectively, the execution layer custom-wires the infrastructure end-to-end for each job and hardware, while keeping the compiler simple and maintainable. We have implemented Astra in two popular deep learning frameworks, PyTorch and Tensorflow. On state-of-the-art deep learning models, we show that Astra improves end-to-end performance of deep learning training by up to 3x, while approaching the performance of hand-optimized implementations such as cuDNN where available. Astra also significantly outperforms static compilation frameworks such as Tensorflow XLA both in performance and robustness. Muthian Sivathanu, Tapan Chugh, Sanjay Sri Vallabh Singapuram, Lidong Zhou |
ASPLOS | 4 |
| 2019 | Fast Distributed Deep Learning over RDMAabstractDeep learning emerges as an important new resource-intensive workload and has been successfully applied in computer vision, speech, natural language processing, and so on. Distributed deep learning is becoming a necessity to cope with growing data and model sizes. Its computation is typically characterized by a simple tensor data abstraction to model multi-dimensional matrices, a dataflow graph to model computation, and iterative executions with relatively frequent synchronizations, thereby making it substantially different from Map/Reduce style distributed big data computation. Jilong Xue, Youshan Miao, Ming Wu 0007, Lidong Zhou |
EuroSys | 6 |
| 2019 | NeuGraph: Parallel Deep Neural Network Computation on Large Graphs
Lingxiao Ma, Zhi Yang 0001, Youshan Miao, Jilong Xue, Ming Wu 0007, Lidong Zhou, Yafei Dai |
USENIX ATC | 6 |
| 2018 | TEE-KV: Secure Immutable Key-Value Store for Trusted Execution EnvironmentsabstractTrusted Execution Environments (TEEs) ensure strong data confidentiality for applications running in the TEEs even on untrusted servers. In particular, TEEs are expected to bring significant benefits to blockchain workloads for enterprise because it ensures confidentiality and correctness of transaction records without any heavy-weight data verification process such as proof-of-work. For example, Coco [3] improves both confidentiality and transaction throughput of existing blockchain protocols by utilizing TEE features. Atsushi Koshiba, Zhongxin Guo, Mitaro Namiki, Lidong Zhou |
SoCC | 5 |
| 2018 | Scheduling CPU for GPU-based Deep Learning JobsabstractDeep learning (DL) is popular in data-center as an important workload for artificial intelligence. With the recent breakthrough of using graphics accelerators and the popularity of DL framework, GPU server cluster dominates DL training in current practice. Cluster scheduler simply treats DL jobs as black-boxes and allocates GPUs as per job request specified by a user. However, other resources, e.g. CPU, are often allocated with workload-agnostic approaches. Kubeflow[1] performs heuristic static CPU resource assignment based on task types (e.g., worker, parameter-server), while [2] evenly divides CPUs of a server to each GPU. Despite the traditional impression that GPU is critical in DL, our observation suggests that the importance of CPU is undervalued. Identifying an appropriate CPU core number in a heterogeneous cluster is challenging yet performance critical to DL jobs. The diverse CPU usage characteristic is not well recognized in the following three aspects. Wencong Xiao, Zhenhua Han, Quanlu Zhang, Fan Yang 0024, Lidong Zhou |
SoCC | 7 |
| 2018 | Capturing and Enhancing In Situ System Observability for Failure Detection
Peng Huang 0005, Chuanxiong Guo, Jacob R. Lorch, Lidong Zhou, Yingnong Dang |
OSDI | 4 |
| 2018 | Gandiva: Introspective Cluster Scheduling for Deep Learning
Wencong Xiao, Romil Bhardwaj, Ramachandran Ramjee, Muthian Sivathanu, Nipun Kwatra, Zhenhua Han, Pratyush Patel, Quanlu Zhang, Fan Yang 0024, Lidong Zhou |
OSDI | 12 |
| 2018 | TerseCades: Efficient Data Compression in Stream Processing
Gennady Pekhimenko, Chuanxiong Guo, Myeongjae Jeon, Peng Huang 0005, Lidong Zhou |
USENIX ATC | 5 |
| 2017 | Gray Failure: The Achilles' Heel of Cloud-Scale SystemsabstractCloud scale provides the vast resources necessary to replace failed components, but this is useful only if those failures can be detected. For this reason, the major availability breakdowns and performance anomalies we see in cloud environments tend to be caused by subtle underlying faults, i.e., gray failure rather than fail-stop failure. In this paper, we discuss our experiences with gray failure in production cloud-scale systems to show its broad scope and consequences. We also argue that a key feature of gray failure is differential observability: that the system's failure detectors may not notice problems even when applications are afflicted by them. This realization leads us to believe that, to best deal with them, we should focus on bridging the gap between different components' perceptions of what constitutes failure. Peng Huang 0005, Chuanxiong Guo, Lidong Zhou, Jacob R. Lorch, Yingnong Dang, Murali Chintalapati, Randolph Yao |
HotOS | 3 |
| 2017 | Tux2: Distributed Graph Computation for Machine Learning
Wencong Xiao, Jilong Xue, Youshan Miao, Ming Wu 0007, Wei Li 0022, Lidong Zhou |
NSDI | 8 |
| 2016 | StreamScope: Continuous Reliable Distributed Processing of Big Data Streams
Wei Lin 0016, Haochuan Fan, Zhengping Qian, Jingren Zhou 0001, Lidong Zhou |
NSDI | 7 |
| 2016 | Realizing the Fault-Tolerance Promise of Cloud Storage Using Locks with Intent
Srinath Setty, Chunzhi Su, Jacob R. Lorch, Lidong Zhou, Hao Chen 0030, Parveen Patel, Jinglei Ren |
OSDI | 4 |
| 2015 | GraM: scaling graph computation to the trillionsabstractGraM is an efficient and scalable graph engine for a large class of widely used graph algorithms. It is designed to scale up to multicores on a single server, as well as scale out to multiple servers in a cluster, offering significant, often over an order-of-magnitude, improvement over existing distributed graph engines on evaluated graph algorithms. GraM is also capable of processing graphs that are significantly larger than previously reported. In particular, using 64 servers (1,024 physical cores), it performs a PageRank iteration in 140 seconds on a synthetic graph with over one trillion edges, setting a new milestone for graph engines. Ming Wu 0007, Fan Yang 0024, Jilong Xue, Wencong Xiao, Youshan Miao, Haoxiang Lin, Yafei Dai, Lidong Zhou |
SoCC | 9 |
| 2015 | ImmortalGraph: A System for Storage and Analysis of Temporal GraphsabstractTemporal graphs that capture graph changes over time are attracting increasing interest from research communities, for functions such as understanding temporal characteristics of social interactions on a time-evolving social graph. ImmortalGraph is a storage and execution engine designed and optimized specifically for temporal graphs. Locality is at the center of ImmortalGraph’s design: temporal graphs are carefully laid out in both persistent storage and memory, taking into account data locality in both time and graph-structure dimensions. ImmortalGraph introduces the notion of locality-aware batch scheduling in computation, so that common “bulk” operations on temporal graphs are scheduled to maximize the benefit of in-memory data locality. The design of ImmortalGraph explores an interesting interplay among locality, parallelism, and incremental computation in supporting common mining tasks on temporal graphs. The result is a high-performance temporal-graph system that is up to 5 times more efficient than existing database solutions for graph queries. The locality optimizations in ImmortalGraph offer up to an order of magnitude speedup for temporal iterative graph mining compared to a straightforward application of existing graph engines on a series of snapshots. Youshan Miao, Ming Wu 0007, Fan Yang 0024, Lidong Zhou, Vijayan Prabhakaran, Enhong Chen |
ACM Trans. Storage | 6 |
| 2015 | Spotting Code Optimizations in Data-Parallel Pipelines through PeriSCOPEabstractTo minimize the amount of data-shuffling I/O that occurs between the pipeline stages of a distributed data-parallel program, its procedural code must be optimized with full awareness of the pipeline that it executes in. Unfortunately, neither pipeline optimizers nor traditional compilers examine both the pipeline and procedural code of a data-parallel program so programmers must either hand-optimize their program across pipeline stages or live with poor performance. To resolve this tension between performance and programmability, this paper describes PeriSCOPE, which automatically optimizes a data-parallel program's procedural code in the context of data flow that is reconstructed from the program's pipeline topology. Such optimizations eliminate unnecessary code and data, perform early data filtering, and calculate small derived values (e.g., predicates) earlier in the pipeline, so that less data - sometimes much less data - is transferred between pipeline stages. PeriSCOPE further leverages symbolic execution to enlarge the scope of such optimizations by eliminating dead code. We describe how PeriSCOPE is implemented and evaluate its effectiveness on real production jobs. Xuepeng Fan, Hai Jin 0001, Xiaofei Liao, Hucheng Zhou, Sean McDirmid, Wei Lin 0016, Jingren Zhou 0001, Lidong Zhou |
IEEE Trans. Parallel Distributed Syst. | 10 |
| 2014 | Rex: replication at the speed of multi-coreabstractStandard state-machine replication involves consensus on a sequence of totally ordered requests through, for example, the Paxos protocol. Such a sequential execution model is becoming outdated on prevalent multi-core servers. Highly concurrent executions on multi-core architectures introduce non-determinism related to thread scheduling and lock contentions, and fundamentally break the assumption in state-machine replication. This tension between concurrency and consistency is not inherent because the total-ordering of requests is merely a simplifying convenience that is unnecessary for consistency. Concurrent executions of the application can be decoupled with a sequence of consensus decisions through consensus on partial-order traces, rather than on totally ordered requests, that capture the non-deterministic decisions in one replica execution and to be replayed with the same decisions on others. The result is a new multi-core friendly replicated state-machine framework that achieves strong consistency while preserving parallelism in multi-thread applications. On 12-core machines with hyper-threading, evaluations on typical applications show that we can scale with the number of cores, achieving up to 16 times the throughput of standard replicated state machines. Chuntao Hong, Mao Yang 0004, Dong Zhou 0006, Lidong Zhou, Li Zhuang |
EuroSys | 5 |
| 2014 | Chronos: a graph engine for temporal graph analysisabstractTemporal graphs capture changes in graphs over time and are becoming a subject that attracts increasing interest from the research communities, for example, to understand temporal characteristics of social interactions on a time-evolving social graph. Chronos is a storage and execution engine designed and optimized specifically for running in-memory iterative graph computation on temporal graphs. Locality is at the center of the Chronos design, where the in-memory layout of temporal graphs and the scheduling of the iterative computation on temporal graphs are carefully designed, so that common "bulk" operations on temporal graphs are scheduled to maximize the benefit of in-memory data locality. The design of Chronos further explores the interesting interplay among locality, parallelism, and incremental computation in supporting common mining tasks on temporal graphs. The result is a high-performance temporal-graph system that offers up to an order of magnitude speedup for temporal iterative graph mining compared to a straightforward application of existing graph engines on a series of snapshots. Youshan Miao, Ming Wu 0007, Fan Yang 0024, Lidong Zhou, Vijayan Prabhakaran, Enhong Chen |
EuroSys | 6 |
| 2014 | Cybertron: pushing the limit on I/O reduction in data-parallel programsabstractI/O reduction has been a major focus in optimizing data-parallel programs for big-data processing. While the current state-of-the-art techniques use static program analysis to reduce I/O, Cybertron proposes a new direction that incorporates runtime mechanisms to push the limit further on I/O reduction. In particular, Cybertron tracks how data is used in the computation accurately at runtime to filter unused data at finer granularity dynamically, beyond what current static-analysis based mechanisms are capable of, and to facilitate a new mechanism called constraint based encoding for more efficient encoding. Cybertron has been implemented and applied to production data-parallel programs; our extensive evaluations on real programs and real data have shown its effectiveness on I/O reduction over the existing mechanisms at reasonable CPU cost, and its improvement on end-to-end performance in various network environments. Tian Xiao, Hucheng Zhou, Xu Zhao 0004, Chencheng Ye 0001, Xi Wang 0005, Wei Lin 0016, Lidong Zhou |
OOPSLA | 10 |
| 2014 | Apollo: Scalable and Coordinated Scheduling for Cloud-Scale Computing
Eric Boutin, Jaliya Ekanayake, Wei Lin 0016, Jingren Zhou 0001, Zhengping Qian, Ming Wu 0007, Lidong Zhou |
OSDI | 8 |
| 2013 | TimeStream: reliable stream computation in the cloudabstractTimeStream is a distributed system designed specifically for low-latency continuous processing of big streaming data on a large cluster of commodity machines. The unique characteristics of this emerging application domain have led to a significantly different design from the popular MapReduce-style batch data processing. In particular, we advocate a powerful new abstraction called resilient substitution that caters to the specific needs in this new computation model to handle failure recovery and dynamic reconfiguration in response to load changes. Several real-world applications running on our prototype have been shown to scale robustly with low latency while at the same time maintaining the simple and concise declarative programming model. TimeStream handles an on-line advertising aggregation pipeline at a rate of 700,000 URLs per second with a 2-second delay, while performing sentiment analysis of Twitter data at a peak rate close to 10,000 tweets per second, with approximately 2-second delay. Zhengping Qian, Chunzhi Su, Zhuojie Wu, Hongyu Zhu 0003, Taizhi Zhang, Lidong Zhou, Zheng Zhang 0001 |
EuroSys | 7 |
| 2013 | Failure Recovery: When the Cure Is Worse Than the Disease
Sean McDirmid, Mao Yang 0004, Li Zhuang, Yingwei Luo, Tom Bergan, Madan Musuvathi, Zheng Zhang 0001, Lidong Zhou |
HotOS | 10 |
| 2013 | KuaFu: Closing the parallelism gap in database replicationabstractDatabase systems are nowadays increasingly deployed on multi-core commodity servers, with replication to guard against failures. Database engine is best designed to scale with the number of cores to offer a high degree of parallelism on a modern multi-core architecture. On the other hand, replication traditionally resorts to a certain form of serialization for data consistency among replicas. In the widely used primary/backup replication with log shipping, concurrent executions on the primary and the serialized log replay on a backup creates a serious parallelism gap. Our experiment on MySQL with a 16-core configuration shows that the serial replay of a backup can sustain only less than one third of the throughput achievable on the primary under an OLTP workload. This paper proposes KuaFu to close the parallelism gap on replicated database systems by enabling concurrent replay of transactions on a backup. KuaFu maintains write consistency on backups by tracking transaction dependencies. Concurrent replay on a backup does introduce read inconsistency between the primary and backups. KuaFu further leverages multi-version concurrency control to produce snapshots in order to restore the consistency semantics. We have implemented KuaFu on MySQL; our evaluations show that KuaFu allows a backup to keep up with the primary while preserving replication consistency. Chuntao Hong, Dong Zhou 0006, Mao Yang 0004, Carbo Kuo, Lidong Zhou |
ICDE | 6 |
| 2012 | Kineograph: taking the pulse of a fast-changing and connected worldabstractKineograph is a distributed system that takes a stream of incoming data to construct a continuously changing graph, which captures the relationships that exist in the data feed. As a computing platform, Kineograph further supports graph-mining algorithms to extract timely insights from the fast-changing graph structure. To accommodate graph-mining algorithms that assume a static underlying graph, Kineograph creates a series of consistent snapshots, using a novel and efficient epoch commit protocol. To keep up with continuous updates on the graph, Kineograph includes an incremental graph-computation engine. We have developed three applications on top of Kineograph to analyze Twitter data: user ranking, approximate shortest paths, and controversial topic detection. For these applications, Kineograph takes a live Twitter data feed and maintains a graph of edges between all users and hashtags. Our evaluation shows that with 40 machines processing 100K tweets per second, Kineograph is able to continuously compute global properties, such as user ranks, with less than 2.5-minute timeliness guarantees. This rate of traffic is more than 10 times the reported peak rate of Twitter as of October 2011. Raymond Cheng 0001, Aapo Kyrola, Youshan Miao, Xuetian Weng, Ming Wu 0007, Fan Yang 0024, Lidong Zhou, Feng Zhao 0001, Enhong Chen |
EuroSys | 8 |
| 2012 | Optimizing Data Shuffling in Data-Parallel Computation by Understanding User-Defined Functions
Hucheng Zhou, Rishan Chen, Xuepeng Fan, Haoxiang Lin, Jack Li 0001, Wei Lin 0016, Jingren Zhou 0001, Lidong Zhou |
NSDI | 10 |
| 2012 | Spotting Code Optimizations in Data-Parallel Pipelines through PeriSCOPE
Xuepeng Fan, Rishan Chen, Hucheng Zhou, Sean McDirmid, Chang Liu 0021, Wei Lin 0016, Jingren Zhou 0001, Lidong Zhou |
OSDI | 10 |
| 2012 | Managing Large Graphs on Multi-Cores with Graph Awareness
Vijayan Prabhakaran, Ming Wu 0007, Xuetian Weng, Frank McSherry, Lidong Zhou, Maya Haradasan |
USENIX ATC | 5 |
| 2011 | Practical software model checking via dynamic interface reductionabstractImplementation-level software model checking explores the state space of a system implementation directly to find potential software defects without requiring any specification or modeling. Despite early successes, the effectiveness of this approach remains severely constrained due to poor scalability caused by state-space explosion. DeMeter makes software model checking more practical with the following contributions: (i) proposing dynamic interface reduction, a new state-space reduction technique, (ii) introducing a framework that enables dynamic interface reduction in an existing model checker with a reasonable amount of effort, and (iii) providing the framework with a distributed runtime engine that supports parallel distributed model checking. Huayang Guo, Ming Wu 0007, Lidong Zhou |
SOSP | 3 |
| 2011 | G2: A Graph Processing System for Diagnosing Distributed Systems
Dong Zhou 0006, Haoxiang Lin, Mao Yang 0004, Fan Long, Chaoqiang Deng, Changshu Liu, Lidong Zhou |
USENIX ATC | 8 |
| 2010 | Comet: batched stream processing for data intensive distributed computingabstractBatched stream processing is a new distributed data processing paradigm that models recurring batch computations on incrementally bulk-appended data streams. The model is inspired by our empirical study on a trace from a large-scale production data-processing cluster; it allows a set of effective query optimizations that are not possible in a traditional batch processing model.We have developed a query processing system called Comet that embraces batched stream processing and integrates with DryadLINQ. We used two complementary methods to evaluate the effectiveness of optimizations that Comet enables. First, a prototype system deployed on a 40-node cluster shows an I/O reduction of over 40% using our benchmark. Second, when applied to a real production trace covering over 19 million machine-hours, our simulator shows an estimated I/O saving of over 50%. Bingsheng He, Mao Yang 0004, Rishan Chen, Wei Lin 0016, Lidong Zhou |
SoCC | 7 |
| 2010 | Language-based replay via data flow cutabstractA replay tool aiming to reproduce a program's execution interposes itself at an appropriate replay interface between the program and the environment. During recording, it logs all non-deterministic side effects passing through the interface from the environment and feeds them back during replay. The replay interface is critical for correctness and recording overhead of replay tools. Ming Wu 0007, Fan Long, Xi Wang 0005, Zhilei Xu, Haoxiang Lin, Xuezheng Liu, Huayang Guo, Lidong Zhou, Zheng Zhang 0001 |
SIGSOFT FSE | 9 |
| 2009 | Wave Computing in the Cloud
Bingsheng He, Mao Yang 0004, Rishan Chen, Wei Lin 0016, Lidong Zhou |
HotOS | 8 |
| 2009 | MODIST: Transparent Model Checking of Unmodified Distributed Systems
Tisheng Chen, Ming Wu 0007, Zhilei Xu, Xuezheng Liu, Haoxiang Lin, Mao Yang 0004, Fan Long, Lidong Zhou |
NSDI | 10 |
| 2009 | Vertical paxos and primary-backup replicationabstractNo abstract available. Leslie Lamport, Dahlia Malkhi, Lidong Zhou |
PODC | 3 |
| 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. | 4 |
| 2009 | Kinesis: A new approach to replica placement in distributed storage systemsabstractKinesis is a novel data placement model for distributed storage systems. It exemplifies three design principles: structure (division of servers into a few failure-isolated segments), freedom of choice (freedom to allocate the best servers to store and retrieve data based on current resource availability), and scattered distribution (independent, pseudo-random spread of replicas in the system). These design principles enable storage systems to achieve balanced utilization of storage and network resources in the presence of incremental system expansions, failures of single and shared components, and skewed distributions of data size and popularity. In turn, this ability leads to significantly reduced resource provisioning costs, good user-perceived response times, and fast, parallelized recovery from independent and correlated failures. This article validates Kinesis through theoretical analysis, simulations, and experiments on a prototype implementation. Evaluations driven by real-world traces show that Kinesis can significantly outperform the widely used Chain replica-placement strategy in terms of resource requirements, end-to-end delay, and failure recovery. John MacCormick, Nicholas Murphy, Venugopalan Ramasubramanian, Udi Wieder, Lidong Zhou |
ACM Trans. Storage | 6 |
| 2008 | Transactional Flash
Vijayan Prabhakaran, Thomas L. Rodeheffer, Lidong Zhou |
OSDI | 3 |
| 2008 | Vigilante: End-to-end containment of Internet worm epidemicsabstractWorm containment must be automatic because worms can spread too fast for humans to respond. Recent work proposed network-level techniques to automate worm containment; these techniques have limitations because there is no information about the vulnerabilities exploited by worms at the network level. We propose Vigilante, a new end-to-end architecture to contain worms automatically that addresses these limitations. In Vigilante, hosts detect worms by instrumenting vulnerable programs to analyze infection attempts. We introduce dynamic data-flow analysis : a broad-coverage host-based algorithm that can detect unknown worms by tracking the flow of data from network messages and disallowing unsafe uses of this data. We also show how to integrate other host-based detection mechanisms into the Vigilante architecture. Upon detection, hosts generate self-certifying alerts (SCAs), a new type of security alert that can be inexpensively verified by any vulnerable host. Using SCAs, hosts can cooperate to contain an outbreak, without having to trust each other. Vigilante broadcasts SCAs over an overlay network that propagates alerts rapidly and resiliently. Hosts receiving an SCA protect themselves by generating filters with vulnerability condition slicing : an algorithm that performs dynamic analysis of the vulnerable program to identify control-flow conditions that lead to successful attacks. These filters block the worm attack and all its polymorphic mutations that follow the execution path identified by the SCA. Our results show that Vigilante can contain fast-spreading worms that exploit unknown vulnerabilities, and that Vigilante's filters introduce a negligible performance overhead. Vigilante does not require any changes to hardware, compilers, operating systems, or the source code of vulnerable programs; therefore, it can be used to protect current software binaries. Manuel Costa, Jon Crowcroft, Miguel Castro 0001, Antony I. T. Rowstron, Lidong Zhou, Paul Barham 0001 |
ACM Trans. Comput. Syst. | 5 |
| 2008 | Niobe: A practical replication protocolabstractThe task of consistently and reliably replicating data is fundamental in distributed systems, and numerous existing protocols are able to achieve such replication efficiently. When called on to build a large-scale enterprise storage system with built-in replication, we were therefore surprised to discover that no existing protocols met our requirements. As a result, we designed and deployed a new replication protocol called Niobe . Niobe is in the primary-backup family of protocols, and shares many similarities with other protocols in this family. But we believe Niobe is significantly more practical for large-scale enterprise storage than previously published protocols. In particular, Niobe is simple, flexible, has rigorously proven yet simply stated consistency guarantees, and exhibits excellent performance. Niobe has been deployed as the backend for a commercial Internet service; its consistency properties have been proved formally from first principles, and further verified using the TLA + specification language. We describe the protocol itself, the system built to deploy it, and some of our experiences in doing so. John MacCormick, Chandramohan A. Thekkath, Marcus Jager, Kristof Roomp, Lidong Zhou, Ryan S. Peterson |
ACM Trans. Storage | 5 |
| 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 | 3 |
| 2007 | Graceful degradation via versions: specifications and implementationsabstractCorrectness of a fault-tolerant system hinges on the failure model, which typically constrains the number of concurrent failures in the system. These assumptions are sometimes violated in practice, inevitably leading to degraded system behavior that deviates from the system's specification and even causing complete unavailability of the system. Lidong Zhou, Vijayan Prabhakaran, Venugopalan Ramasubramanian, Roy Levin, Chandramohan A. Thekkath |
PODC | 1 |
| 2007 | Bouncer: securing software by blocking bad inputabstractAttackers exploit software vulnerabilities to control or crash programs. Bouncer uses existing software instrumentation techniques to detect attacks and it generates filters automatically to block exploits of the target vulnerabilities. The filters are deployed automatically by instrumenting system calls to drop exploit messages. These filters introduce low overhead and they allow programs to keep running correctly under attack. Previous work computes filters using symbolic execution along the path taken by a sample exploit, but attackers can bypass these filters by generating exploits that follow a different execution path. Bouncer introduces three techniques to generalize filters so that they are harder to bypass: a new form of program slicing that uses a combination of static and dynamic analysis to remove unnecessary conditions from the filter; symbolic summaries for common library functions that characterize their behavior succinctly as a set of conditions on the input; and generation of alternative exploits guided by symbolic execution. Bouncer filters have low overhead, they do not have false positives by design, and our results show that Bouncer can generate filters that block all exploits of some real-world vulnerabilities. Manuel Costa, Miguel Castro 0001, Lidong Zhou, Marcus Peinado |
SOSP | 3 |
| 2006 | Brief Announcement: Chasing the Weakest System Model for Implementing Omega and Consensus
Martin Hutle, Dahlia Malkhi, Ulrich Schmid 0001, Lidong Zhou |
SSS | 4 |
| 2005 | Distributed Blinding for Distributed ElGamal Re-EncryptionabstractA protocol is given to take an ElGamal ciphertext encrypted under the key of one distributed service and produce the corresponding ciphertext encrypted under the key of another distributed service, but without the plaintext ever becoming available. Each distributed service comprises a set of servers and employs threshold cryptography to maintain its service private key. Unlike prior work, the protocol requires no assumptions about execution speeds or message delivery delays. The protocol also imposes fewer constraints on where and when various steps are performed, which can bring improvements in end-to-end performance for some applications (e.g., a trusted publish/subscribe infrastructure.) Two new building blocks employed — a distributed blinding protocol and verifiable dual encryption proofs — could have uses beyond re-encryption protocols. Lidong Zhou, Michael A. Marsh, Fred B. Schneider, Anna Redz |
ICDCS | 1 |
| 2005 | Troubleshooting multihop wireless networksabstractEffective network troubleshooting is critical for maintaining efficient and reliable network operation. Troubleshooting is especially challenging in multihop wireless networks because the behavior of such networks depends on complicated interactions between many unpredictable factors such as RF noise, signal propagation, node interference, and traffic flows. In this paper we propose a new direction for research on fault diagnosis in wireless networks. Specifically, we present a diagnostic system that employs trace-driven simulations to detect faults and perform root cause analysis. We apply this approach to diagnose performance problems caused by packet dropping, link congestion, external noise, and MAC misbehavior. In a 25 node multihop wireless network, we are able to diagnose over 10 simultaneous faults of multiple types with more than 80% coverage. Our framework is general enough for a wide variety of wireless and wired networks. Lili Qiu, Paramvir Bahl, Ananth Rao, Lidong Zhou |
SIGMETRICS | 4 |
| 2005 | Vigilante: end-to-end containment of internet wormsabstractWorm containment must be automatic because worms can spread too fast for humans to respond. Recent work has proposed network-level techniques to automate worm containment; these techniques have limitations because there is no information about the vulnerabilities exploited by worms at the network level. We propose Vigilante, a new end-to-end approach to contain worms automatically that addresses these limitations. Vigilante relies on collaborative worm detection at end hosts, but does not require hosts to trust each other. Hosts run instrumented software to detect worms and broadcast self-certifying alerts (SCAs) upon worm detection. SCAs are proofs of vulnerability that can be inexpensively verified by any vulnerable host. When hosts receive an SCA, they generate filters that block infection by analysing the SCA-guided execution of the vulnerable software. We show that Vigilante can automatically contain fast-spreading worms that exploit unknown vulnerabilities without blocking innocuous traffic. Manuel Costa, Jon Crowcroft, Miguel Castro 0001, Antony I. T. Rowstron, Lidong Zhou, Paul Barham 0001 |
SOSP | 5 |
| 2005 | Omega Meets Paxos: Leader Election and Stability Without Eventual Timely Links
Dahlia Malkhi, Florian Oprea, Lidong Zhou |
DISC | 3 |
| 2005 | APSS: proactive secret sharing in asynchronous systemsabstractAPSS, a proactive secret sharing (PSS) protocol for asynchronous systems, is explained and proved correct. The protocol enables a set of secret shares to be periodically refreshed with a new, independent set, thereby thwarting mobile-adversary attacks. Protocols for asynchronous systems are inherently less vulnerable to denial-of-service attacks, which slow processor execution or delay message delivery. So APSS tolerates certain attacks that PSS protocols for synchronous systems cannot. Lidong Zhou, Fred B. Schneider, Robbert van Renesse |
ACM Trans. Inf. Syst. Secur. | 1 |
| 2004 | A Multi-Radio Unification Protocol for IEEE 802.11 Wireless NetworksabstractWe present a link layer protocol called the multi-radio unification protocol or MUP. On a single node, MUP coordinates the operation of multiple wireless network cards tuned to non-overlapping frequency channels. The goal of MUP is to optimize local spectrum usage via intelligent channel selection in a multihop wireless network. MUP works with standard-compliant IEEE 802.11 hardware, does not require changes to applications or higher-level protocols, and can be deployed incrementally. The primary usage scenario for MUP is a multihop community wireless mesh network, where cost of the radios and battery consumption are not limiting factors. We describe the design and implementation of MUP, and analyze its performance using both simulations and measurements based on our implementation. Our results show that under dynamic traffic patterns with realistic topologies, MUP significantly improves both TCP throughput and user perceived latency for realistic workloads. Atul Adya, Paramvir Bahl, Jitendra Padhye, Alec Wolman, Lidong Zhou |
BROADNETS | 5 |
| 2004 | Boxwood: Abstractions as the Foundation for Storage Infrastructure
John MacCormick, Nick Murphy, Marc Najork, Chandramohan A. Thekkath, Lidong Zhou |
OSDI | 5 |
| 2002 | Implementing IPv6 as a Peer-to-Peer Overlay NetworkabstractThis paper proposes to implement an IPv6 routing infrastructure as a self-organizing overlay network on top of the current IPv4 infrastructure. The overlay network builds upon a distributed IPv6 edge router with a master/slave architecture. We show how different slaves can be constructed to tunnel through NATs and firewalls, as well as to improve robustness of the routing infrastructure and to provide efficient and resilient implementations for features such as multicast, anycast, and mobile IP using currently available peer-to-peer (P2P) protocols. The resulting IPv6 overlay network would restore the end-to-end property of the original Internet, support evolution and dynamic updating of the protocols running on the overlay network, make available IPv6 and the associated features to network applications immediately, and provide an ideal underlying infrastructure for P2P applications, without changing networking hardware and software in the core Internet. Lidong Zhou, Robbert van Renesse, Michael A. Marsh |
SRDS | 1 |
| 2002 | COCA: A secure distributed online certification authorityabstractCOCA is a fault-tolerant and secure online certification authority that has been built and deployed both in a local area network and in the Internet. Extremely weak assumptions characterize environments in which COCA's protocols execute correctly: no assumption is made about execution speed and message delivery delays; channels are expected to exhibit only intermittent reliability; and with 3t+ 1 COCA servers up totmay be faulty or compromised. COCA is the first system to integrate a Byzantine quorum system (used to achieve availability) with proactive recovery (used to defend against mobile adversaries which attack, compromise, and control one replica for a limited period of time before moving on to another). In addition to tackling problems associated with combining fault-tolerance and security, new proactive recovery protocols had to be developed. Experimental results give a quantitative evaluation for the cost and effectiveness of the protocols. Lidong Zhou, Fred B. Schneider, Robbert van Renesse |
ACM Trans. Comput. Syst. | 1 |