Byzantine Fault Tolerance Explained
Byzantine fault tolerance lets a distributed system keep working correctly even when some nodes fail arbitrarily or actively lie, not just crash cleanly.
Byzantine fault tolerance, or BFT, is a property of a distributed system that lets it continue operating correctly even when some of its nodes fail in arbitrary, unpredictable ways — including sending conflicting or actively misleading information to different parts of the system — rather than simply crashing. It’s a much harder standard than tolerating clean crashes, and it’s the assumption behind consensus mechanisms that need to work even when some participants can’t be trusted to fail honestly.
Where the name comes from
The term comes from the “Byzantine Generals Problem,” a thought experiment about several generals surrounding a city, needing to agree on a single coordinated plan — attack or retreat — while communicating only through messengers who might be delayed, lost, or, in the worst case, carrying messages from a traitorous general who’s actively trying to sabotage the coordination. The generals need a protocol that reaches agreement despite some fraction of participants behaving maliciously rather than just going silent.
Translated into distributed systems terms: nodes in a network need to agree on a value or an ordering of operations, even when some nodes might be compromised, buggy, or actively adversarial, and every node needs a way to reach the correct decision without necessarily being able to identify which other nodes are the faulty ones.
Crash faults vs Byzantine faults
Most distributed-systems fault tolerance is designed around crash faults: a node stops responding entirely, and the rest of the system just needs to detect the absence and route around it. This is the assumption behind Raft-style consensus, described in more detail in how the Raft algorithm works — a node either participates correctly or is silent, never both.
A Byzantine fault is strictly worse: a faulty node can send different, contradictory information to different peers, can appear to work correctly to some observers while misbehaving toward others, or can actively try to violate the protocol in whatever way benefits an attacker most. The system has to reach correct agreement without any hard guarantee about which nodes are being honest.
| Crash fault tolerance | Byzantine fault tolerance | |
|---|---|---|
| Failure mode assumed | Node stops responding | Node can lie, send conflicting messages, or behave adversarially |
| Typical minimum nodes for consensus | 2f + 1 to tolerate f failures | 3f + 1 to tolerate f failures |
| Common use case | Internal clusters you control (databases, schedulers) | Systems with untrusted or unverified participants (blockchains, some safety-critical systems) |
| Overhead | Lower — fewer messages, simpler protocols | Higher — more message rounds, more redundancy |
Why it takes more nodes
Tolerating f Byzantine faults requires at least 3f + 1 total nodes, compared to 2f + 1 for crash-only fault tolerance — see the general comparison in Raft vs Paxos for why crash-tolerant systems can get away with a smaller quorum. The extra margin exists because honest nodes need enough of a majority that even if f nodes actively lie in the most damaging way possible, the honest nodes can still out-vote the false information and converge on the same correct answer. With fewer nodes, a well-placed set of liars could split the honest nodes’ view of reality in two, and there’d be no way to tell which half was correct.
This is also why quorum-based consensus in a Byzantine setting requires larger overlapping majorities than crash-tolerant quorums: the overlap has to be large enough that it’s mathematically impossible for two conflicting decisions to both gather enough honest votes.
Where Byzantine fault tolerance actually shows up
Most internal infrastructure — a company’s own database cluster, its own job scheduler, its own service mesh — doesn’t need Byzantine fault tolerance, because the nodes are all operated by the same trusted party. A crashed node is a hardware failure or a bug, not sabotage, and crash-fault-tolerant protocols like Raft are the appropriate, cheaper tool.
BFT matters when a system’s participants aren’t all trusted equally, or aren’t controlled by a single operator at all:
- Blockchains and public distributed ledgers, where anyone can run a node and some fraction of participants are assumed to be actively adversarial rather than merely unreliable — this is the foundational assumption most blockchain consensus mechanisms are built to survive.
- Safety-critical and aerospace systems, where a sensor or controller might fail in a way that produces plausible-looking but wrong data, rather than failing silently, and the system needs to keep functioning correctly regardless.
- Multi-organization systems where each participant runs their own infrastructure and no single party can be fully trusted to report honestly, such as some interbank settlement or federated verification systems.
The practical cost
Byzantine fault tolerance isn’t free, and that cost is exactly why it isn’t the default. It requires more nodes for the same fault tolerance, more message rounds to reach agreement (since nodes typically need to cross-check what other nodes claim other nodes said), and correspondingly higher latency and lower throughput than a crash-fault-tolerant protocol running the same workload. Systems adopt it because the trust model genuinely requires it, not because it’s a strictly better default — for a cluster you fully control, the added overhead buys protection against a threat model that mostly doesn’t apply.
The takeaway
Byzantine fault tolerance is what lets a distributed system reach correct agreement even when some participants actively lie or behave adversarially, not just fail silently — a strictly harder problem than crash fault tolerance, requiring more nodes (3f + 1 instead of 2f + 1) and more communication rounds to solve. It’s the right tool when a system’s participants can’t all be trusted equally, such as public blockchains or multi-organization systems, and unnecessary overhead for infrastructure you fully control, where crash-fault-tolerant consensus is the cheaper, sufficient choice.
Keep reading
Chisato · · 5 min read Raft vs Paxos: Consensus Algorithms Compared
Raft and Paxos both let a distributed cluster agree on a value despite failures — Raft trades some flexibility for a design built to be understood.
The Lycoris Team · · 4 min read Read-Your-Writes Consistency Explained
Read-your-writes consistency guarantees a client sees its own writes immediately, even when other clients might not yet. How it's implemented.
The Lycoris Team · · 5 min read Two-Phase Commit vs. Saga Pattern Explained
Two-phase commit locks resources until every node agrees to a transaction; the saga pattern trades that guarantee for availability using compensating steps.