Contents

Contents

LinkedIn Built Kafka. What Did It Change When Kafka Was No Longer Enough?

LinkedIn Built Kafka. What Did It Change When Kafka Was No Longer Enough? featured image

Kafka started at LinkedIn around 2010, giving services a way to publish data, consume it independently and replay it when necessary. By 2025, LinkedIn reported processing more than 32 trillion records and 17 PB per day, across 400,000 topics, over 10,000 machines and roughly 150 clusters.

Operating that infrastructure required an expanding set of services for balancing, recovery and cluster management. LinkedIn's response was Northguard, a new log storage system, and Xinfra, a virtualization layer that lets applications use both Kafka and Northguard.

Northguard keeps producers, consumers, ordered logs, replication and replay. The changes happen underneath the surface: smaller replication units, a sharded control plane and topics that can move between storage systems.

A note on the evidence: This article examines LinkedIn's June 2025 design and deployment report. Northguard and Xinfra were internal, closed-source systems at the time; LinkedIn told InfoQ it would consider open sourcing them later. I could not find a public release for this article. The performance and reliability claims therefore come from LinkedIn, without a public implementation we can inspect or benchmark independently. Results on other workloads, hardware and Kafka configurations remain unverified.

Kafka partitions do a lot of jobs

A Kafka partition is an ordered log, but it also determines replication, broker placement and how consumers divide the work. Several responsibilities share the same boundary:

1.kafka-partition

With twelve partitions, a topic has twelve ordered logs. Each has a replica assignment, a leader accepting writes and followers copying the data. In a conventional consumer group, partitions also determine the available processing parallelism.

LinkedIn's write-up identifies this coupling as a problem at its scale: a partition is a heavyweight replication unit, so moving or repairing one can affect a substantial amount of data.

Adding a broker means dealing with history

Suppose a cluster has three brokers, A, B and C, and is approaching its storage or throughput limit. We add Broker D. It starts empty; using its capacity for existing topics usually requires moving some partition replicas:

2.broker%20a%20-%20broker%20d

Those transfers compete with production traffic for network and disk bandwidth. LinkedIn built Cruise Control to automate partition reassignment, cluster balancing and recovery after failures. The transfer still has to happen, often under a throttle to protect the live workload.

Northguard reduces the amount of data bound to each replica assignment.

Northguard separates the log from its replication units

Northguard's hierarchy looks roughly like this:

3.%20three%20segments

A range represents a log covering a contiguous portion of the keyspace. It contains segments, each with its own sequence of records. An active segment accepts writes; a sealed segment is immutable. LinkedIn describes sealing a segment when it reaches 1 GB, has been active for over an hour, or loses a replica.

The segment is the unit of replication. Kafka also divides partition logs into segment files, but those files normally follow the partition's replica assignment. Northguard can give successive segments in the same range different replica sets:

4.%20log%20striping

The range preserves logical ordering while its segments are distributed across brokers. LinkedIn calls this log striping.

Don't move the old data

Now add Broker E. Northguard can include it in the replica sets of new segments:

5.broker-e-joins

Existing segments can stay where they are. As new segments are created, the allocator distributes new writes across the expanded cluster. This gives a new broker work without first copying historical logs onto it.

It also gives the allocator frequent opportunities to correct uneven placement. The effect develops as segments turn over; adding a broker does not instantly redistribute retained data. For me, this is the most useful change: capacity can start serving new traffic without a large reassignment job first.

Failure recovery changes as well

Suppose active Segment 42 has replicas on A, B and C, and Broker C fails. Northguard can seal that segment and create Segment 43 with a new replica set:

6.missing%20replica

Writes continue on Segment 43 while sealed-segment replication restores the missing copy of Segment 42. Repairing the old data is separate from placing new writes.

LinkedIn credits this smaller replication unit with preserving produce availability without the consistency compromises it made in its Kafka configuration. That comparison concerns LinkedIn's deployment choices; it does not mean Kafka inherently requires sacrificing consistency whenever a broker fails.

Ranges replace indexed partitions

Kafka's partition count also affects key placement. With a mapping such as hash(key) % number_of_partitions, adding partitions can send subsequent records for a key to a different partition. That matters for ordering, stateful processing and keyed joins.

Northguard instead divides the keyspace into ranges that can split:

7.split

The range being split provides the synchronization point. Producers using unrelated ranges can continue. If R1 splits into R2 and R3, records in R1 precede records in either child. Merging ranges creates a corresponding ordering relationship between the parents and the merged range.

This also helps stream processing. Consider joining orders and customers by customer ID. If their key partitioning is aligned, matching records can be processed together. With ten partitions in one topic and sixteen in the other, a framework may need to reshuffle data first.

Northguard uses buddy-style splits and merges, giving ranges aligned keyspace boundaries across topics. LinkedIn chose this partly so processing frameworks could avoid such shuffles when joining streams on the same key.

The control plane had to change too

Kafka's move from ZooKeeper to KRaft changed metadata storage, but the cluster still has one logical metadata state machine and an active controller. LinkedIn describes that centralized control plane as a bottleneck at millions of partition replicas.

Northguard distributes metadata across vnodes. Each vnode is an independent Raft-backed state machine, led by a coordinator. Together, the vnodes cover a hash ring:

8.mmetadata%20hash%20ring

Topic metadata is assigned by topic name; range and segment metadata by range ID. LinkedIn calls this a Dynamically-Sharded Replicated State Machine, or DS-RSM, and describes deployments with 128 or more coordinators and state machines.

Each coordinator manages the metadata assigned to its vnode, including range splits and merges, segment state, replica assignments and topic deletion. These operations no longer all pass through one metadata leader.

Not every piece of state goes through Raft

Brokers also need connection addresses, broker attributes and enough information about the vnode ring to route requests. Northguard distributes this minimal global state through gossip, using SWIM for membership and failure detection.

Authoritative topic, range and segment metadata remains in the Raft-backed vnodes:

9.raft%20and%20swim

Gossip helps a broker find the coordinator responsible for a request. The coordinator's replicated state machine decides the metadata change. Membership dissemination and authoritative updates therefore have separate paths.

Northguard also revisits the disk path

Kafka relies on sequential I/O, the operating system page cache and zero-copy transfer where possible. Northguard's default segment store, fps store, uses a WAL, a file per segment, Direct I/O, a RocksDB sparse index and an application-managed cache.

The store batches appends, writes the WAL, appends records to segment files, calls fsync and updates the index. Direct I/O bypasses normal page-cache buffering; Northguard populates its own cache using knowledge of active consume streams.

This gives LinkedIn more explicit control over caching. It also keeps local disks in the write path, unlike the object-storage approach discussed in my diskless Kafka article. Segment storage is pluggable, but the published default uses small replication units on local disks to meet LinkedIn's latency requirements.

The durability comparison needs some care

LinkedIn reports the following configuration parameters for its two deployments:

LinkedIn's deploymentDisk synchronizationReported sync/batch thresholds
KafkaLazy sync10 seconds or 20,000 records
Northguardfsync on all replicas before producer acknowledgement10 ms, 20,000 records or 10 MB

LinkedIn says Northguard meets its Kafka performance SLOs with those stronger disk-persistence guarantees. The announcement does not provide a reproducible benchmark setup or latency percentiles for that comparison.

Kafka commonly relies on replication for acknowledged-write durability rather than requiring an fsync on every replica before each acknowledgement. The result reported here is that LinkedIn could afford the stronger sync requirement after redesigning its storage path. It does not establish an inherent durability limit in Kafka or a general performance advantage for Northguard.

Produce is a stream, not a sequence of independent RPCs

Metadata operations use ordinary request-response calls. Produce, consume and replication use sessionized streaming protocols.

A producer handshakes with the active segment leader, then pipelines appends within a window limiting data in flight. The broker can acknowledge multiple appends together, but only after their records are committed. Consume reverses the data flow, with the client controlling the window; replication uses the same streaming model.

Changing the protocol also creates a migration problem. Thousands of applications already use Kafka clients and know which physical cluster holds their topics. They need another abstraction before storage can change transparently.

Xinfra virtualizes Pub/Sub

Xinfra sits between applications and the underlying log storage:

10.kafka%20with%20xinfra

Applications use Xinfra clients with a unified API. A topic can move between physical clusters or between Kafka and Northguard without changing its logical name. This depends on adopting the Xinfra client; existing Kafka clients do not gain that indirection automatically.

By the June 2025 announcement, LinkedIn reported that more than 90% of its applications were running Xinfra clients.

Topics have epochs

An Xinfra topic records changes to its physical storage through epochs:

11.epoch

The application can keep using payments while one epoch lives in Kafka and a later one lives in Northguard. Each epoch contains the shards representing that version of the topic. Clients obtain the physical mapping from Xinfra metadata instead of hardcoding it into application configuration.

Kafka to Northguard without stopping the world

Xinfra creates a new epoch in the destination. Producers migrate first and temporarily write to both systems:

12.kafka-northguard

Dual writes allow rollback if the migration fails. Consumers move later; once the transition completes, dual writes stop. Consumers can still read previous epochs until retention removes their data. LinkedIn states that applications continue operating and ordering guarantees are maintained throughout the transition.

By June 2025, LinkedIn reported migrating thousands of topics carrying trillions of records per day. Kafka and Northguard were still running side by side. These figures describe the rollout at publication time, rather than a complete replacement of Kafka.

Xinfra also hides the cluster boundary

Xinfra can group topics from different physical clusters into one virtual cluster:

13.virtual%20cluster

A consumer can subscribe to those topics through the same logical environment even when they are stored in different systems. This lets LinkedIn support workloads larger than one physical cluster and change topic placement without making every application track the topology.

The complexity did not disappear

The application API becomes simpler, but Xinfra adds infrastructure that the platform team must operate:

  • MySQL stores virtual and physical topic and cluster mappings.
  • ZooKeeper handles membership, leadership and consumer-group allocation.
  • Vitess stores consumer checkpoints, with Couchbase providing a low-latency cache.

For LinkedIn's roughly 150 clusters, centralizing this work can be worth the cost. With two Kafka clusters, maintaining the virtualization layer might cost more than the problem it solves.

Perhaps the partition was the real thing LinkedIn outgrew

Northguard and Xinfra separate responsibilities at three levels. Ranges define the logical keyspace and ordering; segments define replication and placement. Independent Raft groups distribute metadata management. Xinfra separates the topic applications use from the physical system storing it:

14.data%20plane

Reading this as LinkedIn outgrowing the partition is my interpretation. Its write-up identifies several concrete problems: coarse replication units, indexed partitions, balancing and centralized metadata. The common thread is how many operations depend on the same placement and ordering boundaries.

Kafka had already carried LinkedIn to an extraordinary scale. For a deployment with forty topics and six brokers, building Northguard and Xinfra would be a very expensive response to routine operational problems. The useful question is which costs in your own system come from the workload, and which come from coupling ordering, replication, placement and parallelism to one partition.

Reviewed by Krzysztof Atłasik

Blog Comments powered by Disqus.