Home / Software y Cloud / Is it there or not? The punch of the Bloom filter

Is it there or not? The punch of the Bloom filter

Ilustración de un filtro de Bloom

In a distributed system there is no “single queue” letting nodes find out what is happening. Each server only knows a handful of neighbors, and no central authority tells them who has just fallen over. Even so, a cluster of hundreds of machines can detect failures and spread information within seconds. The engine behind that daily miracle is the gossip protocol, also called an epidemic protocol: the way servers tell things to each other, neighbor to neighbor, until everyone knows everything.

The corridor metaphor

The name is no accident: it mimics how a rumor spreads in an office. One person tells two or three colleagues, those tell a few more, and within a few rounds the whole building knows, even though nobody called a general meeting. What fascinates is the math: if in each round every node contacts k neighbors and the network has n nodes, the information covers practically the whole network in O(log n) rounds. Propagation grows exponentially: 1, 3, 9, 27… until saturation.

Push, pull and exchanging rumors

There are two basic contagion mechanisms. In push mode, the node with news actively sends it to its neighbors (“I tell you what I know”). In pull mode, the inquiring node pulls news from others (“tell me what you know”). Real systems combine both in push-pull mode: each exchange is bidirectional, so both parties end up with the information the other had. When a node receives a message it already knows, the conversation loses interest and propagation dies out on its own, a property called anti-entropy.

Gossip versus quorum: two philosophies

Compare it with the alternative we have seen in protocols such as Raft. Raft is sequential and quorum-based: it demands a majority, orders the replicas and, before a failure, waits for the leader to decide. It is correct and deterministic, but expensive: latency and message count grow with cluster size. Gossip, by contrast, trades strict ordering guarantees for horizontal scalability: each node talks only with a handful of neighbors, so the per-node cost stays roughly constant even with thousands of machines. In exchange, information converges with probability, not with instant certainty: it is an eventually consistent system.

SWIM: failure detection at scale

The best-known use case is failure detection: knowing which nodes are still alive. The SWIM algorithm (Scalable Weakly-consistent Infection-style process group Membership) solves it in a distributed way: each node randomly picks a member and sends it a ping; if there is no answer, it does not immediately conclude the node is dead. It asks k witnesses to send an indirect ping, and only if all of them fail does it declare the node down. This avoids the false positives caused by a mere network hiccup and removes the bottleneck of a central coordinator.

Even better: each node keeps a suspicion clock with exponential decay, a variant called the phi-accrual failure detector (the one used by Cassandra). Instead of a yes/no, it computes a continuous probability that the node is down, based on the distribution of inter-heartbeat times. The stranger a silence is, the more suspicious it becomes; in this way the system can tell a micro-latency from a real outage.

Membership is contagious too

Gossip does more than detect dead nodes: it also disseminates which members are alive, their addresses and their metadata, a set called the membership list. Cassandra and Consul use it so that each node learns the topology of the ring without a central server; the open-source memberlist library from HashiCorp implements SWIM with arbitrary gossip payloads, letting services replicate resources, publish node health and distribute credentials over a single channel. It is also the backbone of reachability detection in Kubernetes and of metadata dissemination in Prometheus.

Engineering concerns: port, neighbors and version

In practice, a well-designed gossip protocol manages three things. First, a dedicated port (e.g. 7946 in Consul) separate from the data-traffic port. Second, an anti-partitioning mechanism: if the network splits in two, each half keeps gossiping internally and both eventually converge separately, making it clear where every node is. Third, a version vector or timestamp per entry: when two nodes exchange the membership list, they keep the entry with the highest version and drop conflicts, ensuring the freshest rumor is not lost under a stale one.

Gossip for shared cryptography

Here is a concrete case with security tension: several nodes must share a secret (for example, an encryption key) without it ever passing through a central channel. Through gossip, each node passes fragments of the material with vector protection; eventually they all converge to the same key, even though each pair exchanged different pieces. It is a scheme of eventual cryptographic replication, used for instance for cluster secrets in Consul, and it shows that gossip can be as good as authoritative synchronization when it is well designed.

The price: complexity of reasoning

Why gossip is not used for everything: its probabilistic behavior makes it hard to reason about the state at a given instant, and the emphasis on eventual delivery clashes with operations that need strong or linearizable consistency. A financial ledger must know right away which transaction won; there, quorum rules. The usual wisdom is hybrid: use a quorum protocol (Raft, Paxos) for critical state and gossip for cluster health, metadata dissemination and low-criticality event propagation. Each sheep, as the rumors say, with its own partner.