Before 2017 every consumer of published content at The New York Times integrated with each producing system's own API and bootstrapped its history by calling those APIs. The 2017 publishing pipeline replaced that with one Kafka log. What changed about the cost of adding the eleventh consumer, and about who pays for a backfill?
Show the full answer Hide the answer
The situation they were in
Producers of content exposed APIs to read it and feeds to be told it had changed, and every consumer integrated with every producer it cared about. With M producers and N consumers the integration count grows as M times N, and each of those integrations carries its own authentication, pagination, error semantics and schema drift.
The sharper problem is bootstrapping. A new consumer had to walk every producer's API to build its initial state, which made onboarding a load event on systems that had no way to schedule or price it. The producers absorbed a cost created by somebody else's roadmap, which is the structural reason this kind of integration estate ossifies: adding a consumer is a negotiation, so people stop adding consumers and start copying databases instead.
What they chose
A single ordered log as the source of truth for published content, which they call the Monolog, described by Boerge Svingen on the Confluent engineering blog in September 2017. Every producing system appends a published asset to it; every consumer builds its own store by reading the log. Retention is effectively forever, and the entire archive back to 1851 was under 100 GB, which is why a single-partition topic on a single disk was a design rather than a compromise.
Why it fit their constraints
Integration count collapses from M times N to M plus N. A new consumer learns one schema.
The part worth taking away is where the backfill cost went. Bootstrapping stops being an unscheduled load event on the producers and becomes a sequential read from a broker, which is the one operation a log is unambiguously good at. The producers are not involved and cannot be hurt by a consumer's bad week. Replay being cheap for the party that pays for it is what makes "just rebuild it from the log" a sentence anyone can say out loud, and it is a property of the data size as much as of the technology.
What it cost them
The log's schema is a permanent public contract. A consumer replaying from offset zero must still interpret a message written in 2017, so every field ever published is supported forever and the schema registry becomes a governance function rather than a convenience.
Producers also gave up control of audience. A topic is readable by whoever can read the topic, so "who sees this field" stops being a per-consumer API decision and becomes a broker authorisation decision that someone has to actually administer.
And the ordering that makes replay deterministic caps throughput at what one partition and one consumer can sustain.
When not to copy this
Three conditions have to hold, and most systems break at least one.
- The corpus must be small enough that a full replay is an operation rather than a project. Under 100 GB, not under 100 TB. A clickstream or a telemetry feed arrives at a rate that forces partitioning, and partitioned means per-key ordering only, which weakens "rebuild anything from the log" to "rebuild per-key state from the log".
- The data must not be per-person. An immutable forever-log and an erasure obligation are in direct conflict, and that conflict cannot be engineered away once the log exists. Published articles are a corpus with no subject-access problem; orders and messages are not.
- M times N has to actually be large. With three producers and four consumers, twelve integrations is less total work than operating a log platform, a schema registry and the discipline that keeps both honest. If there is one consumer, the log is a queue and the database is still the truth.