Files
nexus/sreweekly/articles/106/05-building-a-distributed-log-from-scratch-part-4-trade-offs-and-lessons.html
2026-09-12 17:23:01 +08:00

25 lines
37 KiB
HTML
Raw Blame History

This file contains invisible Unicode characters
This file contains invisible Unicode characters that are indistinguishable to humans but may be processed differently by a computer. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
<!doctype html><html lang=en dir=auto data-theme=auto><head><meta charset=utf-8><meta http-equiv=X-UA-Compatible content="IE=edge"><meta name=viewport content="width=device-width,initial-scale=1,shrink-to-fit=no"><meta name=robots content="index, follow"><title>Building a Distributed Log from Scratch, Part 4: Trade-Offs and Lessons Learned | Brave New Geek</title><meta name=keywords content="architecture,building a distributed log from scratch,distributed log,distributed systems,kafka,message queues,message-oriented middleware,messaging,nats,nats streaming,performance,raft,scalability,stream processing,systemantics,systems theory"><meta name=description content="In part three of this series we talked about scaling message delivery in a distributed log. In part four, we’ll look at some key trade-offs involved with such systems and discuss a few lessons learned while building NATS Streaming.
Competing Goals
There are a number of competing goals when building a distributed log (these goals also extend to many other types of systems). Recall from part one that our key priorities for this type of system are performance, high availability, and scalability. The preceding parts of this series described at various levels how we can accomplish these three goals, but astute readers likely noticed that some of these things conflict with one another."><meta name=author content><link rel=canonical href=https://bravenewgeek.com/building-a-distributed-log-from-scratch-part-4-trade-offs-and-lessons-learned/><link crossorigin=anonymous href=/assets/css/stylesheet.4861a452a4c13a9a1fbf2085400b74a7de96b1beeb94dee57b4273e5dffcf337.css integrity="sha256-SGGkUqTBOpofvyCFQAt0p96Wsb7rlN7le0Jz5d/88zc=" rel="preload stylesheet" as=style><link rel=icon href=https://bravenewgeek.com/favicon.ico><link rel=icon type=image/png sizes=16x16 href=https://bravenewgeek.com/favicon.ico><link rel=icon type=image/png sizes=32x32 href=https://bravenewgeek.com/favicon.ico><link rel=apple-touch-icon href=https://bravenewgeek.com/favicon.ico><link rel=mask-icon href=https://bravenewgeek.com/favicon.ico><meta name=theme-color content="#2e2e33"><meta name=msapplication-TileColor content="#2e2e33"><link rel=alternate hreflang=en href=https://bravenewgeek.com/building-a-distributed-log-from-scratch-part-4-trade-offs-and-lessons-learned/><noscript><style>#theme-toggle,.top-link{display:none}</style><style>@media(prefers-color-scheme:dark){:root{--theme:rgb(29, 30, 32);--entry:rgb(46, 46, 51);--primary:rgb(218, 218, 219);--secondary:rgb(155, 156, 157);--tertiary:rgb(65, 66, 68);--content:rgb(196, 196, 197);--code-block-bg:rgb(46, 46, 51);--code-bg:rgb(55, 56, 62);--border:rgb(51, 51, 51);color-scheme:dark}.list{background:var(--theme)}.toc{background:var(--entry)}}</style></noscript><script>localStorage.getItem("pref-theme")==="dark"?document.querySelector("html").dataset.theme="dark":localStorage.getItem("pref-theme")==="light"?document.querySelector("html").dataset.theme="light":window.matchMedia("(prefers-color-scheme: dark)").matches?document.querySelector("html").dataset.theme="dark":document.querySelector("html").dataset.theme="light"</script><link rel=preconnect href=https://fonts.googleapis.com><link rel=preconnect href=https://fonts.gstatic.com crossorigin><link rel=stylesheet href="https://fonts.googleapis.com/css2?family=JetBrains+Mono:wght@400;500&family=Source+Serif+4:ital,opsz,wght@0,8..60,400;0,8..60,600;1,8..60,400&family=Space+Grotesk:wght@500;600;700&display=swap"><meta property="og:url" content="https://bravenewgeek.com/building-a-distributed-log-from-scratch-part-4-trade-offs-and-lessons-learned/"><meta property="og:site_name" content="Brave New Geek"><meta property="og:title" content="Building a Distributed Log from Scratch, Part 4: Trade-Offs and Lessons Learned"><meta property="og:description" content="In part three of this series we talked about scaling message delivery in a distributed log. In part four, we’ll look at some key trade-offs involved with such systems and discuss a few lessons learned while building NATS Streaming.
Competing Goals There are a number of competing goals when building a distributed log (these goals also extend to many other types of systems). Recall from part one that our key priorities for this type of system are performance, high availability, and scalability. The preceding parts of this series described at various levels how we can accomplish these three goals, but astute readers likely noticed that some of these things conflict with one another."><meta property="og:locale" content="en_us"><meta property="og:type" content="article"><meta property="article:section" content="posts"><meta property="article:published_time" content="2018-01-18T16:01:13-06:00"><meta property="article:modified_time" content="2018-02-23T16:10:13-06:00"><meta property="article:tag" content="Architecture"><meta property="article:tag" content="Building a Distributed Log From Scratch"><meta property="article:tag" content="Distributed Log"><meta property="article:tag" content="Distributed Systems"><meta property="article:tag" content="Kafka"><meta property="article:tag" content="Message Queues"><meta name=twitter:card content="summary"><meta name=twitter:title content="Building a Distributed Log from Scratch, Part 4: Trade-Offs and Lessons Learned"><meta name=twitter:description content="In part three of this series we talked about scaling message delivery in a distributed log. In part four, we’ll look at some key trade-offs involved with such systems and discuss a few lessons learned while building NATS Streaming.
Competing Goals There are a number of competing goals when building a distributed log (these goals also extend to many other types of systems). Recall from part one that our key priorities for this type of system are performance, high availability, and scalability. The preceding parts of this series described at various levels how we can accomplish these three goals, but astute readers likely noticed that some of these things conflict with one another."><script type=application/ld+json>{"@context":"https://schema.org","@type":"BreadcrumbList","itemListElement":[{"@type":"ListItem","position":1,"name":"Posts","item":"https://bravenewgeek.com/posts/"},{"@type":"ListItem","position":2,"name":"Building a Distributed Log from Scratch, Part 4: Trade-Offs and Lessons Learned","item":"https://bravenewgeek.com/building-a-distributed-log-from-scratch-part-4-trade-offs-and-lessons-learned/"}]}</script><script type=application/ld+json>{"@context":"https://schema.org","@type":"BlogPosting","headline":"Building a Distributed Log from Scratch, Part 4: Trade-Offs and Lessons Learned","name":"Building a Distributed Log from Scratch, Part 4: Trade-Offs and Lessons Learned","description":"In part three of this series we talked about scaling message delivery in a distributed log. In part four, we’ll look at some key trade-offs involved with such systems and discuss a few lessons learned while building NATS Streaming.\nCompeting Goals There are a number of competing goals when building a distributed log (these goals also extend to many other types of systems). Recall from part one that our key priorities for this type of system are performance, high availability, and scalability. The preceding parts of this series described at various levels how we can accomplish these three goals, but astute readers likely noticed that some of these things conflict with one another.\n","keywords":["architecture","building a distributed log from scratch","distributed log","distributed systems","kafka","message queues","message-oriented middleware","messaging","nats","nats streaming","performance","raft","scalability","stream processing","systemantics","systems theory"],"articleBody":"In part three of this series we talked about scaling message delivery in a distributed log. In part four, we’ll look at some key trade-offs involved with such systems and discuss a few lessons learned while building NATS Streaming.\nCompeting Goals There are a number of competing goals when building a distributed log (these goals also extend to many other types of systems). Recall from part one that our key priorities for this type of system are performance, high availability, and scalability. The preceding parts of this series described at various levels how we can accomplish these three goals, but astute readers likely noticed that some of these things conflict with one another.\nIt’s easy to make something fast if it’s not fault-tolerant or scalable. If our log runs on a single server, our only constraints are how fast we can send data over the network and how fast the disk I/O is. And this is how a lot of systems, including many databases, tend to work—not only because it performs well, but because it’s simple. We can make these types of systems fault-tolerant by introducing a standby server and allowing clients to failover, but there are a couple issues worth mentioning with this.\nWith data systems, such as a log, high availability does not just pertain to continuity of service, but also availability of data. If I write data to the system and the system acknowledges that, that data should not be lost in the event of a failure. So with a standby server, we need to ensure data is replicated to avoid data loss (otherwise, in the context of a message log, we must relax our requirement of guaranteed delivery).\nNATS Streaming initially shipped as a single-node system, which raised immediate concerns about production-readiness due to a single point of failure. The first step at trying to address some of these concerns was to introduce a fault-tolerance mode whereby a group of servers would run and only one would run as the active server. The active server would obtain an exclusive lock and process requests. Upon detecting a failure, standby servers would attempt to obtain the lock and become the active server.\nAside from the usual issues with distributed locks, this design requires a shared storage layer. With NATS Streaming, this meant either a shared volume, such as Gluster or EFS, or a shared MySQL database. This poses a performance challenge and isn’t particularly “cloud-native” friendly. Another issue is data is not replicated unless done so out-of-band by the storage layer. When we add in data replication, performance is hamstrung even further. But this was a quick and easy solution that offered some solace with respect to a SPOF (disclosure: I was not involved with NATS or NATS Streaming at this time). The longer term solution was to provide first-class clustering and data-replication support, but sometimes it’s more cost effective to provide fast recovery of a single-node system.\nAnother challenge with the single-node design is scalability. There is only so much capacity that one node can handle. At a certain point, scaling out becomes a requirement, and so we start partitioning. This is a common technique for relational databases where we basically just run multiple databases and divide up the data by some key. NATS Streaming is no different as it offers a partitioning story for dividing up channels between servers. The trouble with partitioning is it complicates things as it typically requires cooperation from the application. To make matters worse, NATS Streaming does not currently offer partitioning at the channel level, which means if a single topic has a lot of load, the solution is to manually partition it into multiple channels at the application level. This is why Kafka chose to partition its topics by default.\nSo performance is at odds with fault-tolerance and scalability, but another factor is what I call simplicity of mechanism. That is, the simplicity of the design plays an important role in the performance of a system. This plays out at multiple levels. We saw that, at an architectural level, using a simple, single-node design performs best but falls short as a robust solution. In part one, we saw that using a simple file structure for our log allowed us to take advantage of the hardware and operating system in terms of sequential disk access, page caching, and zero-copy reads. In part two, we made the observation that we can treat the log itself as a replicated WAL to solve the problem of data replication in an efficient way. And in part three, we discussed how a simple pull-based model can reduce complexity around flow control and batching.\nAt the same time, simplicity of “UX” makes performance harder. When I say UX, I mean the ergonomics of the system and how easy it is to use, operate, etc. NATS Streaming initially optimized for UX, which is why it fills an interesting space. Simplicity is a core part of the NATS philosophy, so it caught a small mindshare with developers frustrated or overwhelmed by Kafka. There is appetite for a “Kafka lite,” something which serves a similar purpose to Kafka but without all the bells and whistles and probably not targeted at large enterprises—a classic Innovator’s Dilemma to be sure.\nNATS Streaming tracks consumer positions automatically, provides simple APIs, and uses a simple push-based protocol. This also means building a client library is a much less daunting task. The downside is the server needs to do more work. With a single node, as NATS Streaming was initially designed, this isn’t much of a problem. Where it starts to rear its head is when we need to replicate that state across a cluster of nodes. This has important implications with respect to performance and scale. Smart middleware has a natural tendency to become more complex, more fragile, and slower. The end-to-end principle attests to this. Amusingly, NATS Streaming was originally named STAN because it’s the opposite of NATS, a fast and simple messaging system with minimal guarantees.\nSimplicity of mechanism tends to simply push complexity around in the system. For example, NATS Streaming provides an ergonomic API to clients by shifting the complexity to the server. Kafka scales and performs exceptionally well by shifting the complexity to other parts of the system, namely the client and ZooKeeper.\nScalability and fault-tolerance are equally at odds with simplicity for reasons mostly described above. The important point here is that these cannot be an afterthought. As I learned while implementing clustering in NATS Streaming, you can’t cleanly and effectively bolt on fault-tolerance onto an existing complex system. One of the laws of Systemantics comes to mind here: “A complex system designed from scratch never works and cannot be patched up to make it work. You have to start over, beginning with a working simple system.” Scalability and fault-tolerance need to be designed from day one.\nLastly, availability is inherently at odds with consistency. This is simply the CAP theorem. Guaranteeing strong consistency requires a quorum when replicating data, which hinders availability and performance. The key here is minimize what you need to replicate or relax your requirements.\nLessons Learned The section above already contains several lessons learned in the process of working on NATS Streaming and implementing clustering, but I’ll capture a few important ones here.\nFirst, distributed systems are complex enough. Simple is usually better—and faster. Again, we go back to the laws of systems here: “A complex system that works is invariably found to have evolved from a simple system that works.”\nSecond, lean on existing work. A critical part to delivering clustering rapidly was sticking with Raft and an existing Go implementation for leader election and data replication. There was considerable time spent designing a proprietary solution before I joined which still had edge cases not fully thought through. Not only is Raft off the shelf, it’s provably correct (implementation bugs notwithstanding). And following from the first lesson learned, start with a solution that works before worrying about optimization. It’s far easier to make a correct solution fast than it is to make a fast solution correct. Don’t roll your own coordination protocol if you don’t need to (and chances are you don’t need to).\nThere are probably edge cases for which you haven’t written tests. There are many failures modes, and you can only write so many tests. Formal methods and property-based testing can help a lot here. Similarly, chaos and fault-injection testing such as Kyle Kingsbury’s Jepsen help too.\nLastly, be honest with your users. Don’t try to be everything to everyone. Instead, be explicit about design decisions, trade-offs, guarantees, defaults, etc. If there’s one takeaway from Kyle’s Jepsen series it’s that many vendors are dishonest in their documentation and marketing. MongoDB became infamous for having unsafe defaults and implementation issues early on, most likely because they make benchmarks look much more impressive.\nIn part five of this series, we’ll conclude by outlining the design for a new log-based system that draws from ideas in the previous entries in the series.\n","wordCount":"1510","inLanguage":"en","datePublished":"2018-01-18T16:01:13-06:00","dateModified":"2018-02-23T16:10:13-06:00","mainEntityOfPage":{"@type":"WebPage","@id":"https://bravenewgeek.com/building-a-distributed-log-from-scratch-part-4-trade-offs-and-lessons-learned/"},"publisher":{"@type":"Organization","name":"Brave New Geek","logo":{"@type":"ImageObject","url":"https://bravenewgeek.com/favicon.ico"}}}</script></head><body id=top><header class=site-header><div class="wrap site-header-inner"><a class=wordmark href=https://bravenewgeek.com/ accesskey=h title="Brave New Geek (Alt + H)"><span class=wordmark-name>Brave New Geek</span>
<span class=wordmark-tag>Introspections of a software engineer</span></a><nav class=site-nav aria-label=Primary><a href=https://bravenewgeek.com/archive/>Archive</a>
<a href=https://bravenewgeek.com/tags/>Tags</a>
<a href=https://bravenewgeek.com/about-me/>About</a>
<button id=theme-toggle class=theme-toggle accesskey=t title="Toggle theme (Alt + T)" aria-label="Toggle light/dark theme">
<svg class="moon" width="18" height="18" viewBox="0 0 24 24" fill="none" stroke="currentColor" stroke-width="2" stroke-linecap="round" stroke-linejoin="round"><path d="M21 12.79A9 9 0 1111.21 3 7 7 0 0021 12.79z"/></svg>
<svg class="sun" width="18" height="18" viewBox="0 0 24 24" fill="none" stroke="currentColor" stroke-width="2" stroke-linecap="round" stroke-linejoin="round"><circle cx="12" cy="12" r="5"/><line x1="12" y1="1" x2="12" y2="3"/><line x1="12" y1="21" x2="12" y2="23"/><line x1="4.22" y1="4.22" x2="5.64" y2="5.64"/><line x1="18.36" y1="18.36" x2="19.78" y2="19.78"/><line x1="1" y1="12" x2="3" y2="12"/><line x1="21" y1="12" x2="23" y2="12"/><line x1="4.22" y1="19.78" x2="5.64" y2="18.36"/><line x1="18.36" y1="5.64" x2="19.78" y2="4.22"/></svg></button></nav></div></header><main class=main><article class="post wrap"><div class=post-return><a href=https://bravenewgeek.com/><span class=pager-arrow>&larr;</span> the log</a></div><header class=post-header><div class=post-meta><span class=post-offset>#72</span><time datetime=2018-01-18>2018-01-18</time><span>8 min read</span>
<span class=post-cats><a href=https://bravenewgeek.com/category/distributed-systems-2/>Distributed Systems</a><a href=https://bravenewgeek.com/category/messaging/>Messaging</a><a href=https://bravenewgeek.com/category/software-architecture/>Software Architecture</a><a href=https://bravenewgeek.com/category/software-engineering/>Software Engineering</a><a href=https://bravenewgeek.com/category/systems-theory/>Systems Theory</a></span></div><h1 class=post-title>Building a Distributed Log from Scratch, Part 4: Trade-Offs and Lessons Learned</h1></header><div class="post-content md-content"><p>In <a href=https://bravenewgeek.com/building-a-distributed-log-from-scratch-part-3-scaling-message-delivery/>part three</a> of this series we talked about scaling message delivery in a distributed log. In part four, we’ll look at some key trade-offs involved with such systems and discuss a few lessons learned while building NATS Streaming.</p><h3 id=competing-goals>Competing Goals<a hidden class=anchor aria-hidden=true href=#competing-goals>#</a></h3><p>There are a number of competing goals when building a distributed log (these goals also extend to many other types of systems). Recall from <a href=https://bravenewgeek.com/building-a-distributed-log-from-scratch-part-1-storage-mechanics/>part one</a> that our key priorities for this type of system are performance, high availability, and scalability. The preceding parts of this series described at various levels how we can accomplish these three goals, but astute readers likely noticed that some of these things conflict with one another.</p><p>It’s easy to make something fast if it’s not fault-tolerant or scalable. If our log runs on a single server, our only constraints are how fast we can send data over the network and how fast the disk I/O is. And this is how a lot of systems, including many databases, tend to work—not only because it performs well, but because it’s simple. We can make these types of systems fault-tolerant by introducing a standby server and allowing clients to failover, but there are a couple issues worth mentioning with this.</p><p>With data systems, such as a log, high availability does not just pertain to continuity of service, but also availability of data. If I write data to the system and the system acknowledges that, that data should not be lost in the event of a failure. So with a standby server, we need to ensure data is replicated to avoid data loss (otherwise, in the context of a message log, we must relax our requirement of guaranteed delivery).</p><p>NATS Streaming initially shipped as a single-node system, which raised immediate concerns about production-readiness due to a single point of failure. The first step at trying to address some of these concerns was to introduce a fault-tolerance mode whereby a group of servers would run and only one would run as the active server. The active server would obtain an exclusive lock and process requests. Upon detecting a failure, standby servers would attempt to obtain the lock and become the active server.</p><p>Aside from the <a href=https://martin.kleppmann.com/2016/02/08/how-to-do-distributed-locking.html>usual issues with distributed locks</a>, this design requires a shared storage layer. With NATS Streaming, this meant either a shared volume, such as Gluster or EFS, or a shared MySQL database. This poses a performance challenge and isn’t particularly “cloud-native” friendly. Another issue is data is not replicated unless done so out-of-band by the storage layer. When we add in data replication, performance is hamstrung even further. But this was a quick and easy solution that offered some solace with respect to a SPOF (disclosure: I was not involved with NATS or NATS Streaming at this time). The longer term solution was to provide first-class clustering and data-replication support, but sometimes it’s more cost effective to provide fast recovery of a single-node system.</p><p>Another challenge with the single-node design is scalability. There is only so much capacity that one node can handle. At a certain point, scaling out becomes a requirement, and so we start partitioning. This is a common technique for relational databases where we basically just run multiple databases and divide up the data by some key. NATS Streaming is no different as it offers a <a href=https://github.com/nats-io/nats-streaming-server#partitioning>partitioning story</a> for dividing up channels between servers. The trouble with partitioning is it complicates things as it typically requires cooperation from the application. To make matters worse, NATS Streaming does not currently offer partitioning at the channel level, which means if a single topic has a lot of load, the solution is to manually partition it into multiple channels at the application level. This is why Kafka chose to partition its topics by default.</p><p>So performance is at odds with fault-tolerance and scalability, but another factor is what I call <em>simplicity of mechanism</em>. That is, the simplicity of the design plays an important role in the performance of a system. This plays out at multiple levels. We saw that, at an architectural level, using a simple, single-node design performs best but falls short as a robust solution. In <a href=https://bravenewgeek.com/building-a-distributed-log-from-scratch-part-1-storage-mechanics/>part one</a>, we saw that using a simple file structure for our log allowed us to take advantage of the hardware and operating system in terms of sequential disk access, page caching, and zero-copy reads. In <a href=https://bravenewgeek.com/building-a-distributed-log-from-scratch-part-2-data-replication/>part two</a>, we made the observation that we can treat the log itself as a replicated WAL to solve the problem of data replication in an efficient way. And in <a href=https://bravenewgeek.com/building-a-distributed-log-from-scratch-part-3-scaling-message-delivery/>part three</a>, we discussed how a simple pull-based model can reduce complexity around flow control and batching.</p><p>At the same time, simplicity of “UX” makes performance harder. When I say UX, I mean the ergonomics of the system and how easy it is to use, operate, etc. NATS Streaming initially optimized for UX, which is why it fills an interesting space. Simplicity is a core part of the NATS philosophy, so it caught a small mindshare with developers frustrated or overwhelmed by Kafka. There is appetite for a “Kafka lite,” something which serves a similar purpose to Kafka but without all the bells and whistles and probably not targeted at large enterprises—a classic Innovator’s Dilemma to be sure.</p><p>NATS Streaming tracks consumer positions automatically, provides simple APIs, and uses a simple push-based protocol. This also means building a client library is a much less daunting task. The downside is the server needs to do more work. With a single node, as NATS Streaming was initially designed, this isn’t much of a problem. Where it starts to rear its head is when we need to replicate that state across a cluster of nodes. This has important implications with respect to performance and scale. Smart middleware has a natural tendency to become <a href="https://www.youtube.com/watch?v=JHQlA_tB10c">more complex, more fragile, and slower</a>. The <a href=https://bravenewgeek.com/from-the-ground-up-reasoning-about-distributed-systems-in-the-real-world/>end-to-end principle</a> attests to this. Amusingly, NATS Streaming was originally named STAN because it’s the opposite of NATS, a fast and simple messaging system with minimal guarantees.</p><p>Simplicity of mechanism tends to simply push complexity around in the system. For example, NATS Streaming provides an ergonomic API to clients by shifting the complexity to the server. Kafka scales and performs exceptionally well by shifting the complexity to other parts of the system, namely the client and ZooKeeper.</p><p>Scalability and fault-tolerance are equally at odds with simplicity for reasons mostly described above. The important point here is that these <em>cannot</em> be an afterthought. As I learned while implementing clustering in NATS Streaming, you can’t cleanly and effectively <em>bolt on</em> fault-tolerance onto an existing complex system. One of the laws of <a href=https://en.wikipedia.org/wiki/Systemantics>Systemantics</a> comes to mind here: “A complex system designed from scratch never works and cannot be patched up to make it work. You have to start over, beginning with a working simple system.” Scalability and fault-tolerance need to be designed from day one.</p><p>Lastly, availability is inherently at odds with consistency. This is simply the <a href=https://bravenewgeek.com/cap-and-the-illusion-of-choice/>CAP theorem</a>. Guaranteeing strong consistency requires a quorum when replicating data, which hinders availability and performance. The key here is minimize what you need to replicate or relax your requirements.</p><h3 id=lessons-learned>Lessons Learned<a hidden class=anchor aria-hidden=true href=#lessons-learned>#</a></h3><p>The section above already contains several lessons learned in the process of working on NATS Streaming and implementing clustering, but I’ll capture a few important ones here.</p><p>First, distributed systems are complex enough. Simple is usually better—and faster. Again, we go back to the laws of systems here: “A complex system that works is invariably found to have evolved from a simple system that works.”</p><p>Second, lean on existing work. A critical part to delivering clustering rapidly was sticking with Raft and an existing Go implementation for leader election and data replication. There was considerable time spent designing a proprietary solution before I joined which still had edge cases not fully thought through. Not only is Raft off the shelf, it’s <em>provably</em> correct (implementation bugs notwithstanding). And following from the first lesson learned, start with a solution that <em>works</em> before worrying about optimization. It’s far easier to make a <em>correct</em> solution <em>fast</em> than it is to make a <em>fast</em> solution <em>correct</em>. Don’t roll your own coordination protocol if you don’t need to (and chances are you don’t need to).</p><p>There are probably edge cases for which you haven’t written tests. There are <em>many</em> failures modes, and you can only write so many tests. Formal methods and property-based testing can help a lot here. Similarly, chaos and fault-injection testing such as Kyle Kingsbury’s <a href=https://github.com/jepsen-io/jepsen>Jepsen</a> help too.</p><p>Lastly, be honest with your users. Don’t try to be everything to everyone. Instead, be explicit about design decisions, trade-offs, guarantees, defaults, etc. If there’s one takeaway from Kyle’s <a href=https://jepsen.io/analyses>Jepsen series</a> it’s that many vendors are dishonest in their documentation and marketing. MongoDB became infamous for having <a href=https://aphyr.com/posts/322-call-me-maybe-mongodb-stale-reads>unsafe defaults</a> and <a href=https://aphyr.com/posts/284-call-me-maybe-mongodb>implementation</a> <a href=https://jepsen.io/analyses/mongodb-3-4-0-rc3>issues</a> early on, most likely because they make benchmarks look much more impressive.</p><p>In <a href=https://bravenewgeek.com/building-a-distributed-log-from-scratch-part-5-sketching-a-new-system/>part five</a> of this series, we’ll conclude by outlining the design for a new log-based system that draws from ideas in the previous entries in the series.</p></div><footer class=post-footer><ul class=post-tags><li><a href=https://bravenewgeek.com/tag/architecture/>Architecture</a></li><li><a href=https://bravenewgeek.com/tag/building-a-distributed-log-from-scratch/>Building a Distributed Log From Scratch</a></li><li><a href=https://bravenewgeek.com/tag/distributed-log/>Distributed Log</a></li><li><a href=https://bravenewgeek.com/tag/distributed-systems/>Distributed Systems</a></li><li><a href=https://bravenewgeek.com/tag/kafka/>Kafka</a></li><li><a href=https://bravenewgeek.com/tag/message-queues/>Message Queues</a></li><li><a href=https://bravenewgeek.com/tag/message-oriented-middleware/>Message-Oriented Middleware</a></li><li><a href=https://bravenewgeek.com/tag/messaging/>Messaging</a></li><li><a href=https://bravenewgeek.com/tag/nats/>Nats</a></li><li><a href=https://bravenewgeek.com/tag/nats-streaming/>Nats Streaming</a></li><li><a href=https://bravenewgeek.com/tag/performance/>Performance</a></li><li><a href=https://bravenewgeek.com/tag/raft/>Raft</a></li><li><a href=https://bravenewgeek.com/tag/scalability/>Scalability</a></li><li><a href=https://bravenewgeek.com/tag/stream-processing/>Stream Processing</a></li><li><a href=https://bravenewgeek.com/tag/systemantics/>Systemantics</a></li><li><a href=https://bravenewgeek.com/tag/systems-theory/>Systems Theory</a></li></ul><nav class=post-nav aria-label="Adjacent posts"><a class=post-nav-link href=https://bravenewgeek.com/building-a-distributed-log-from-scratch-part-5-sketching-a-new-system/><span class=post-nav-dir><span class=pager-arrow>&larr;</span> newer</span>
<span class=post-nav-title>Building a Distributed Log from Scratch, Part 5: Sketching a New System</span>
</a><a class="post-nav-link post-nav-right" href=https://bravenewgeek.com/building-a-distributed-log-from-scratch-part-3-scaling-message-delivery/><span class=post-nav-dir>older <span class=pager-arrow>&rarr;</span></span>
<span class=post-nav-title>Building a Distributed Log from Scratch, Part 3: Scaling Message Delivery</span></a></nav></footer><section class=wp-comments><h2>Comments</h2><p class=wp-comments-notice>Comments are from this blog's WordPress era and are preserved read-only.</p><article class=wp-comment><header><span class=wp-comment-author>Kyle</span>
<time class=wp-comment-date>March 21, 2018</time></header><div class=wp-comment-body><p>for &#8220;Guaranteeing strong consistency requires a quorum when replicating data, which hinders availability and performance.&#8221; here. As I understand, strong consistency means all copy are consistent, not only quorum. Quorum consistency is a trade of between C and A. Strong consistency lives in traditional DB system, not catch up current tendency.</p></div><div class=wp-comment-replies><article class=wp-comment><header><span class=wp-comment-author>Tyler Treat</span>
<time class=wp-comment-date>March 21, 2018</time></header><div class=wp-comment-body><p>Quorum is sufficient as long as reads are served by the leader.</p></div></article></div></article><article class=wp-comment><header><span class=wp-comment-author>Kyle</span>
<time class=wp-comment-date>March 27, 2018</time></header><div class=wp-comment-body><p>To make clear, as in CAP theorem, there are two Quorum protocol, one is partial quorum protocol that is popular used (partition tolerant), which is strong consistency, should not have any divergent; another is full strict quorum protocol, such as 2PC (used in MySQL cluster, not partition tolerant), which is week consistency, can have different copies.</p></div></article></section></article></main><footer class=site-footer><div class="wrap site-footer-inner"><div class=footer-meta><span class=footer-copy>&copy; 2026 Tyler Treat</span>
<span class="footer-sep footer-dot">·</span>
<span class=footer-links><a href=/feed/>rss</a>
<span class=footer-sep>·</span>
<a href=https://github.com/tylertreat target=_blank rel="noopener noreferrer me">github</a>
<span class=footer-sep>·</span>
<a href=https://www.linkedin.com/in/ttreat/ target=_blank rel="noopener noreferrer me">linkedin</a></span></div></div></footer><a href=#top id=top-link class="top-link hidden" aria-label="go to top" title="Go to Top (Alt + G)" accesskey=g><svg viewBox="0 0 24 24" fill="none" stroke="currentColor" stroke-width="2" stroke-linecap="round" stroke-linejoin="round" class="feather feather-chevrons-up"><polyline points="17 11 12 6 7 11"/><polyline points="17 18 12 13 7 18"/></svg>
</a><script>let menu=document.getElementById("menu");if(menu){const e=localStorage.getItem("menu-scroll-position");e&&(menu.scrollLeft=parseInt(e,10)),menu.onscroll=function(){localStorage.setItem("menu-scroll-position",menu.scrollLeft)}}document.querySelectorAll('a[href^="#"]').forEach(e=>{e.addEventListener("click",function(e){e.preventDefault();var t=this.getAttribute("href").substr(1);window.matchMedia("(prefers-reduced-motion: reduce)").matches?document.querySelector(`[id='${decodeURIComponent(t)}']`).scrollIntoView():document.querySelector(`[id='${decodeURIComponent(t)}']`).scrollIntoView({behavior:"smooth"}),t==="top"?history.replaceState(null,null," "):history.pushState(null,null,`#${t}`)})})</script><script>var toplink=document.getElementById("top-link");window.onscroll=function(){const e=window.innerHeight;document.body.scrollTop>e||document.documentElement.scrollTop>e?toplink.classList.remove("hidden"):toplink.classList.add("hidden")}</script><script>document.getElementById("theme-toggle").addEventListener("click",()=>{const e=document.querySelector("html");e.dataset.theme==="dark"?(e.dataset.theme="light",localStorage.setItem("pref-theme","light")):(e.dataset.theme="dark",localStorage.setItem("pref-theme","dark"))})</script><script>document.querySelectorAll("pre > code").forEach(e=>{const n=e.parentNode.parentNode,t=document.createElement("button");t.classList.add("copy-code"),t.innerHTML="copy";function s(){t.innerHTML="copied!",setTimeout(()=>{t.innerHTML="copy"},2e3)}t.addEventListener("click",t=>{if("clipboard"in navigator){navigator.clipboard.writeText(e.textContent),s();return}const n=document.createRange();n.selectNodeContents(e);const o=window.getSelection();o.removeAllRanges(),o.addRange(n);try{document.execCommand("copy"),s()}catch{}o.removeRange(n)}),n.classList.contains("highlight")?n.appendChild(t):n.parentNode.firstChild==n||(e.parentNode.parentNode.parentNode.parentNode.parentNode.nodeName=="TABLE"?e.parentNode.parentNode.parentNode.parentNode.parentNode.appendChild(t):e.parentNode.appendChild(t))})</script></body></html>