Every distributed system eventually has to answer the same question: how does news travel?
The answer determines almost everything else about the system — its latency profile, its failure modes, its operational cost, and the kinds of guarantees it can make to the applications running on top of it.
There are three mainstream answers, and they are often confused with each other:
- Gossip protocols — the Cassandra/SWIM approach. Nodes periodically exchange state with a small number of random peers, and updates propagate exponentially through the network.
- Event streaming — the Kafka approach. Producers append records to a partitioned, replicated log; consumers read from that log at their own pace.
- Remote procedure calls — the gRPC approach. A client sends a typed request to a server, and the server responds.
They are not competitors. They are tools with different jobs. But if you pick the wrong one for a given workload, you pay for it in latency, complexity, or correctness — and sometimes all three.
This post explains what we use at Beelogik, why, and where gossip belongs in the picture.
The case for Kafka: durable, ordered, replayable
Apache Kafka is a distributed, partitioned, replicated commit log. Producers write records to topics; each topic is split into partitions; each partition is a strictly ordered sequence of records; and each partition is replicated across brokers for fault tolerance.
The properties that fall out of this model are exactly what we need for the workloads we run:
- Durability. Records are persisted to disk and replicated before they are acknowledged. A broker can die, and no committed record is lost.
- Ordering. Within a partition, records are ordered. This is what makes Kafka usable for state-change events, where the order of operations matters.
- Replay. Because the log is durable and consumers track their position (the “offset”), any consumer can rewind and replay history. This is impossible with a pure request/response system.
- Backpressure. If a consumer is slow, the log simply grows. No messages are dropped, no producer is blocked. The system degrades gracefully instead of falling over.
- Decoupling. Producers do not know who consumes their records. Adding a new consumer requires no change on the producer side.
- Exactly-once semantics. Kafka’s transactional API lets a producer write to multiple partitions atomically, and lets a consumer process and commit offsets in a single transaction.
For Beelogik, this is the foundation. Every state change in a Beelogik deployment is an event on a Kafka topic. Every downstream system — analytics, billing, observability, reconciliation — reads from those topics at its own pace.
The case for gRPC: typed contracts, low overhead
Kafka is excellent at moving events between services asynchronously. It is not the right tool for a client asking a server a question and waiting for an answer.
For synchronous request/response — the classic “the client needs an answer now” pattern — we use gRPC.
gRPC is a remote procedure call framework built on HTTP/2 and Protocol Buffers. It gives us:
- Typed contracts. The interface between client and server is defined in a
.protofile. Both sides generate code from it, so the contract cannot drift. - HTTP/2 multiplexing. Many requests flow over a single connection, with no head-of-line blocking between streams.
- Bidirectional streaming. Clients and servers can push streams of messages to each other in real time — useful for telemetry, subscriptions, and long-running operations.
- Low overhead. Protocol Buffers are compact and fast to parse. The wire format is smaller than JSON by a factor of 3–10 for typical payloads.
- Deadlines and cancellation. Every call carries a deadline. If the deadline is missed, both sides receive a cancellation signal and release resources.
For Beelogik, gRPC is the edge. Every external client — the console, the CLI, partner integrations — speaks gRPC to our gateways. Those gateways translate incoming calls into Kafka events, which the rest of the system consumes.
The case against gossip (for us)
Gossip protocols are elegant. They solve a real problem — keeping state eventually consistent across a large number of nodes without a coordinator — and they solve it well. Cassandra, Riak, Consul, and Serf all lean on gossip, and each of them is a solid piece of engineering.
But gossip is not free. The properties that make it attractive also make it unsuitable for the workloads Beelogik is optimized for.
Eventual consistency is not enough for payments
Gossip protocols converge eventually. For configuration data, membership, or caching, “eventually” is measured in hundreds of milliseconds, and nobody notices. For a payment settlement, “eventually” is measured in whether the customer’s money moved or not, and everybody notices.
Beep, our cross-border payments platform, settles BRL → CNY transactions. The state of each settlement must be unambiguous at every moment. A shopper’s Pix transfer either happened or it didn’t. There is no “it will happen in a few gossip rounds.”
Kafka’s log model gives us linearizable ordering within a partition, and its transactional API gives us atomic writes across partitions. That is what a payments workload requires.
Bandwidth grows with N²
Gossip protocols scale their propagation depth logarithmically, but their total bandwidth still grows with the square of the number of nodes, because every node eventually gossips with every other node.
In a 128-node cluster, that is fine. In a 1,000-node cluster, it becomes the dominant cost. Kafka’s model avoids this entirely — brokers exchange messages with each other only for replication, and consumers pull from brokers on their own schedule. The network cost scales linearly with the number of consumers, not quadratically with the number of nodes.
Debugging a gossip system is hard
When something goes wrong in a gossip-based system, there is no single log to read. The state is distributed by design, and reconstructing the causal chain of an incident requires instrumenting every node and reasoning about partial orderings.
Kafka inverts this. Because every event is a durable record in a log, every incident has a first-class artifact: the sequence of records that led to it. We can replay that sequence, inspect it, and reproduce the failure. That is the single biggest operational advantage of the event-driven model.
Where gossip does belong
We are not saying gossip is wrong. We are saying it is the wrong default for what we build.
Gossip remains the right choice when:
- You need to propagate low-stakes data (health, membership, approximate metrics) to a very large, loosely-coupled cluster.
- You can tolerate eventual consistency, and you prefer availability over linearizability.
- You want a system with no coordinator at all, and you are willing to pay the bandwidth cost for it.
That is Cassandra’s territory, and it is a good place to be. It is just not our territory.
At Beelogik, the state we care about — orders, payments, settlements, user identity — needs ordering, durability, and replayability. Those are the properties Kafka delivers.
How Kafka and gRPC fit together
Here is the canonical shape of a Beelogik request, from a client’s point of view:
- A client calls a gRPC gateway. The call is typed, deadline-bound, and can be streamed.
- The gateway validates the call and publishes an event to Kafka. The event is durable and replicated before the gateway responds to the client.
- Downstream services consume the event. Each service does its work — settlement, compliance, notification — and emits its own events.
- A status endpoint (also gRPC) reflects progress. If the client wants to wait for completion, it subscribes to a status stream and receives updates as events progress.
Every state transition is a Kafka record. Every synchronous interaction is a gRPC call. Nothing relies on gossip or on eventual convergence of critical state.
That is what “event-driven” actually means in our stack. It is not a buzzword — it is a specific set of guarantees that only Kafka’s log model provides.
The trade-offs we accept
Kafka is not free either. We pay for it in three ways:
- Operational complexity. A Kafka cluster needs careful capacity planning, monitoring, and tuning. We have automated most of this, but it is not a zero-config system.
- Latency floor. A round trip through Kafka adds a few milliseconds versus direct request/response. For most workloads this is invisible; for ultra-low-latency query paths, we bypass Kafka and use gRPC directly.
- Topic design. Poorly-designed topics (too many, too few, or badly-keyed partitions) can turn a Kafka cluster into a footgun. We invest heavily in the schema and topic registry to prevent this.
These are real costs. We pay them because the alternative — an eventually-consistent system that cannot cleanly represent a payment — is not acceptable for what we build.
Why this matters for your business
You do not need to know how Kafka’s replication protocol works to benefit from it. What you get, at the product level, is:
- Ordered, durable state. Every event is a record. Every record is persisted before it is acknowledged.
- Replayable history. Any incident can be reproduced from the event log.
- Decoupled services. Adding a new consumer requires no change to producers.
- Strict consistency for critical operations. When a write returns, it is committed — period.
The beehive does not gossip about where the nectar is. It carries the information back, records it, and every bee that needs it can read it. That is the model we build on.