Migrating to Etcd's Raft Implementation

October 3, 2026

In my attempt to learn and explore the Raft algorithm, I used a library called dragonboat.

But just like what I mentioned in my previous article, I encountered performance issue: high latency and small throughput. This is what I wrote:

5 RPS already blew the p90 latency to 4.6s.

To be really honest, I was testing it based on it's Key-Value (KV) database example, without any fine tuning; and I was testing it in a cluster of 3$ VMs running in only 1 core CPU and 1 Gb memory.

But still, I expected more: at least, at 5 RPS, the p90 should still be sub-second.

This led to my journey to find another Golang's Raft implementation. There were two other candidates ready in my mind:

Both have nice APIs, good documentation, and more importantly, just like dragonboat, have an KV database example in which I can break down and appropriate it in to my own library.

I choose the etcd-io arbitrarily based on subjective feeling: somewhat the API feels nicer, especially based on its usage in their own KV database example (all three libraries provide KV database example).

Unlike dragonboat, Etcd raft is more flexible—it exposes its Write Ahead Log (WAL), storage, and network layer so the user can extend. Of course, they also provides the ready to use default implementation.

Hashicorp raft is also the same: they exposes their WAL, storage, and network layer for extension. But the way they handles commit & apply operation is more opinionated—you need to provide your Finite State Machine (FSM) implementation that has the Apply function to "apply"; while Etcd raft is more flexible (or, lower level) and allow separations between the commit and apply—you can "apply", but if you don't write it in the WAL on your own, it will not be "applied"; and if the node restart, it will consider that as not yet applied.

Besides this API, maybe I'm also afraid in case Hashicorp changed its library's licensing policy in the future. ahem.

Again, most of my decision are based on my subjective feeling after taking a quick look on their KV database example that allowed me to see the usage of their API, without any complicated fine tunings.

Nevertheless, I'm grateful to the great engineers implementing these libraries. It allowed me to learn by examples.

But then again, I don't want to only learn. I want it to be useful for my use cases:

Create a cheap, replicated storage; or a stateful service/app with embedded local DB, that allowed omitting external DB; that supports moderate load for creating service/app Proof of Concepts (PoC) but can scale to production load easily.

(Load) Testing it early and doing some descriptive statistics on the latency (not even an advanced analysis) allowed me to reason about the performance and it's scalability.

Load Testing the Library

Heads up:

Like said previously, when choosing raft library:

I FIRST based it on subjective feeling when trying the APIs, ease of running the examples, etc. without any complicated fine tuning. SECOND only finalize it after performance (load) testing.

This article is not meant to test all of the library and do a fair comparison. My capacity is limited, and I short-circuited to whoever fits above criteria first. Feel free to do your own benchmarking to reproduce this load tests, or to test the other raft library.

The test started with an ideal imagination: what if I have 3 nodes raft cluster that spans across cities (not only availability zone [AZ]). If one city is down (wars, ahem), I still have my raft app/service running.

I did it also to understand the latency of a (larger) region, knowing the test I did here puts a latency cap for smaller region: eg. within the same cities but different AZ, or inside the same data center (DC), or inside within my LAN.

So I spawned it in these three cities in South East Asia (SEA): Jakarta, Singapore, and Bangkok. Each has 2 core CPU / 4Gb memory (around 10 USD each in Tencent Cloud)

Each installed Tailscale to allow private and secure communication for my raft application to bind and communicate with each other over TCP/IP (the default transport implementation for etcd is HTTP).

Then, I installed my app called storage in each node. It is a simple "CRUD" app (albeit it uses "upsert") that is backed by an embedded KV database. Every delete and upsert operation—or mutation, or "command" in Event Sourcing lingo—is routed to the etcd raft library and replicated.

For the actual load testing, I spawned another VM in Singapore as the load-generator that runs a k6 script here. The load testing script outputs a csv file that is analyzed using the script here.

The load is doing just a single upsert operation to the app and waited for the raft "apply" result synchronously before returning.

Finally, I used Claude/ChatGPT to visualize the load testing result across various runs.

Result

The performance is fit for my use cases.

For a total of 30 USD a month, I already get replicated storage (state machine even) that supports 1,6k write RPS within p9999=159.99 ms.

Result The result ("nice") chart

There were assumptions:

Outliers were removed to make the data easier to understand and to remove factor such as a short network issue. I used Inter-Quartile Range with k=12 value chosen arbitrarily.

There were some limitations:

Above 2k RPS, the latency graph started to plateau.

I expected exponential graph from theory, but I suspected the load tester VM cannot keep up with the load (see reduced "samples kept" count in the result table); or, network bandwidth limitations makes the plateau. I haven't investigated deeply in this part.

Plateau The not-so-nice chart showing the latency "plateau". I suspect: (1) because it only represent "successful request" (see number of samples kept), (2) load-testing VM limited resource, (3) network bandwidth. But haven't investigated anything. How about reads?

Reads benefit are given:

Eventually-consistent is good enough for me, albeit linearizable-read seemed nice. Dragonboat provides linearizable read out-of-the-box, but for Etcd, I still need to figure out on how to implement it. But using optimistic-lock seems good enough for this case to avoid write attempt of stale read data. Since I can use cloudflare tunnel, request is routed to the nearby node from the user, with better latency. If you live in Bangkok, you are likely to access the Bangkok node. Additionally, I can use cloudflare load balancer for various load balancing scenario.

What's next?

I will emphasis my motivation with Raft:

  • Learning raft allows me to think of (stateful) app/service reliability from the get go, not as an after-thought.
  • Learning raft allows me to knowledge-transfer it to other consensus algorithm. Making learning and integrating them later easier.
  • Implementing raft allows me to have cheaper replicated storage for day-to-day use.
  • Implementing raft allows me to create an app/service without external database.

In exploring all raft related functionality, the learning path is immense, but not steep. I can do it slowly while iterating, one reliable application/service at a time. For example, how to allow snapshotting easily, scaling-up, scaling-down, enable follower node, linearizable-reads, managing downstream connection (including load balancing), and many more.

Since some of my previously created application/service are tightly coupled with the dragonboat implementation (intentional design decision). You'll find me busy refactoring those in foreseeable future.