Most social architectures are explained through the timeline: how a post fans out to followers. LinkedIn's architecture is better explained by two other ideas. The first is network distance. Almost every member-facing surface, from search results to the feed to recommendations of people to connect with, shows or uses whether someone is a first, second or third degree connection. The second is derived data. A member's profile is written rarely and read in dozens of forms: as a search document, a feed card, an analytics row, a recommendation feature. LinkedIn built its platform around turning one write into all of those reads.

This article reconstructs that architecture from LinkedIn's engineering publications, explains why each piece exists, and works through the graph-distance problem with code. It is a system design study, not an inventory of current internals; companies retire systems, and some components named here have successors. The push fan-out alternative is covered in Twitter timeline generation.

Advertisement

What makes the problem different

LinkedIn reports over a billion registered members. The scale is large but the shape is what matters. Connections are mutual and deliberately curated, so the graph is undirected, and most members have hundreds rather than millions of connections. That makes the graph far denser in the middle and flatter at the top than a follower graph. A celebrity on a follow-based network creates a fan-out problem on write; on LinkedIn, the expensive operation is reading the neighbourhood: who is within two or three hops of me, and how is this person connected to me?

The workload is also extremely read-heavy and heterogeneous. One profile edit should change what appears in search, in the feed, in typeahead, in recruiter tools and in analytics. Serving all of those from one database would mean one schema optimised for nothing. LinkedIn's answer, described repeatedly by its engineers, is to separate the source of truth from many purpose-built derived stores and connect them with a log.

The big picture

Source of truth, change log, derived serving storesWeb and mobileclientsAPI frontendsRest.li resourcesMid-tier servicesprofile, feed, searchSource-of-truth storesEspresso (doc store on MySQL)writesChange captureDatabus, later BrooklinbinlogKafkachange events, activity eventsStream processingSamza jobsBatchHadoop / SparkVenicederived key-valuePinotreal-time OLAPGalenesearch indexesLIquidin-memory graphbatch pushMid-tier services read derived stores for almost every page; only edits touch the source of truth
Figure 1. LinkedIn's data architecture as described in its engineering publications: writes land in a source-of-truth store, change capture publishes them to Kafka, and stream and batch jobs build derived serving stores.

Requests come in through frontend API services. LinkedIn built and open-sourced Rest.li, a framework for typed REST resources with schema-defined data models, so that hundreds of services can call each other with consistent conventions. Mid-tier services own domains such as profiles, connections, feed and search.

Writes go to a source-of-truth store. For much of LinkedIn's member data that is Espresso, a document store LinkedIn built on top of MySQL storage nodes, partitioned by key, with routing, secondary indexing and replication layered above. Every committed change is captured from the database's replication log by a change-capture service. LinkedIn first built Databus for this and later Brooklin, a more general streaming system, and the change events land in Kafka.

From Kafka, stream processors (LinkedIn created Apache Samza for this) and batch jobs on Hadoop and Spark compute derived views. Those views are served by specialised stores: Venice for derived key-value data such as features and recommendations, Apache Pinot for real-time analytics such as who viewed your profile, the Galene search platform for search indexes, and LIquid, an in-memory graph database, for connection queries. Mid-tier services read mostly from these derived stores.

Advertisement

The log at the centre

Apache Kafka began at LinkedIn as the answer to an integration mess: every system needed data from every other system, and point-to-point pipelines multiplied. Jay Kreps' essay on the log, written while at LinkedIn, sets out the principle. If every change is appended to an ordered, durable, replayable log, any number of consumers can build their own view at their own pace, and a new view can be built by replaying history.

This gives three properties that the rest of the architecture depends on. Consumers are decoupled: search indexing falling behind does not slow profile writes. Ordering is per partition, so keying change events by member ID guarantees each member's updates are applied in order by every consumer. And views are rebuildable: if a derived store is corrupted or its schema changes, you reprocess the log or the batch snapshot rather than migrate in place. The partitioning mechanics that make this scale are covered in Kafka partition architecture.

Change capture from the database's replication log, rather than having services publish events themselves, matters for correctness. If a service writes to the database and then publishes to Kafka, a crash between the two loses the event. Reading the commit log means an event exists exactly when the write committed. Where change capture is not available, the outbox pattern gives the same guarantee.

Derived data: one write, many reads

A derived store is a materialised view computed from the log. Its consumer must be idempotent and ordered per key, because after a failure it will reprocess events it already applied. The sketch below builds a per-member search document from profile change events, storing the source offset with each record so replays are harmless.

def apply_profile_change(event, store):
    """event: {member_id, offset, fields: {...}, deleted: bool} from a member-keyed partition."""
    key = event["member_id"]
    current = store.get(key)              # derived record, or None
    if current and current["src_offset"] >= event["offset"]:
        return                            # already applied: replay after a crash is a no-op
    if event["deleted"]:
        store.put(key, {"tombstone": True, "src_offset": event["offset"]})
        return
    doc = dict(current or {})
    doc.update(project_for_search(event["fields"]))  # only fields search needs
    doc["src_offset"] = event["offset"]
    store.put(key, doc)

def project_for_search(fields):
    out = {}
    for f in ("headline", "current_title", "current_company", "skills", "region"):
        if f in fields:
            out[f] = fields[f]
    return out

Venice generalises this. LinkedIn described it as the successor to Voldemort's read-only mode for serving derived data: a batch job computes a full dataset offline and pushes it as a new version, which the store swaps in atomically, while a hybrid mode also applies a stream of near-real-time updates on top of the latest batch version. That is the lambda architecture implemented inside a storage system, so application teams do not have to merge batch and stream results themselves. The same separation of write model and read models appears in CQRS.

The graph: computing network distance

Network distance is the signature LinkedIn query. Showing a 2nd badge next to a search result means answering, for the viewer and that candidate, whether they share at least one connection. Naively this is a graph search per result, and a results page may hold dozens of candidates.

The arithmetic explains the design. If a typical member has a few hundred connections, their second-degree network is the union of their connections' connections, which can reach tens or hundreds of thousands of members after overlap. The third degree can reach millions. Precomputing and storing every member's second-degree set is too large and changes with every new connection, so the practical approach is to keep the first-degree adjacency lists in memory, sharded across machines, and compute distances at query time with set intersections.

def label_distances(viewer, candidates, adj, max_degree=3):
    """adj: member -> set of connections (first-degree adjacency, held in memory).
    Returns {candidate: 1, 2, 3 or None}. Intersections only; never materialise degree 3."""
    f1 = adj.get(viewer, set())
    labels = {}
    for c in candidates:
        if c == viewer:
            continue
        if c in f1:
            labels[c] = 1
            continue
        c1 = adj.get(c, set())
        # degree 2: viewer and candidate share a connection
        small, large = (c1, f1) if len(c1) < len(f1) else (f1, c1)
        if any(m in large for m in small):
            labels[c] = 2
            continue
        # degree 3: some connection of the candidate is connected to some connection of the viewer
        if max_degree >= 3 and any(adj.get(m, set()) & f1 for m in c1):
            labels[c] = 3
        else:
            labels[c] = None
    return labels

Walk through one request. The viewer has 400 connections, and search returns 20 candidates. For each candidate the degree-1 check is a hash lookup. The degree-2 check intersects two sets of a few hundred entries, iterating the smaller one: microseconds. The degree-3 check is the expensive part, because it touches the adjacency lists of every connection of the candidate, perhaps 400 lookups of 400-entry sets per candidate, and those lists live on other shards. Real systems cap the work per request, batch the cross-shard fetches, cache the viewer's second-degree set for the session, and accept that a third-degree label may occasionally be omitted.

LinkedIn described LIquid as a general in-memory graph database built for this class of query, with a declarative, Datalog-like query language, replacing an earlier purpose-built graph service. The design lesson holds regardless of product: keep the graph in memory, partition it so a member's adjacency list lives in one place, and push the intersection to the data rather than shipping adjacency lists to the caller.

The feed: pull, then rank

LinkedIn's engineers have described the feed as largely pull-based. Rather than writing every post into every connection's inbox, activity is stored per actor, and at read time the feed service gathers recent activity from the viewer's network and followed entities, then ranks it. A dense, mutual graph with modest degree makes this tractable, and it lets ranking change without rewriting stored timelines.

Ranking is multi-stage. A first pass retrieves and lightly scores a large candidate set from network activity, followed companies and topics, and recommended content. A second pass applies heavier machine-learned models whose features are read from derived stores such as Venice, then business rules enforce diversity and freshness. Features are precomputed by batch and stream jobs, so the request path reads them rather than computing them.

Failure modes and operations

  • Derived data lag and read-your-writes. A member edits their headline, reloads, and sees the old one because the search document has not caught up. Fix: route the editing member's own reads to the source of truth for a short window, or render the change optimistically on the client.
  • Consumer lag cascades. A slow Samza job falls behind, and every view it feeds goes stale together. Monitor lag per consumer group in time, not just messages, and alert on the age of the oldest unprocessed event.
  • Reprocessing storms. Rebuilding a large view by replaying history can saturate brokers and the target store. Throttle replays, and use the batch push path for full rebuilds rather than the stream.
  • Schema evolution. A producer drops a field a consumer relies on. Use schema-defined events with compatibility checks in a registry, and enforce backward compatibility in CI.
  • Graph hot spots. Members with very large networks make degree queries expensive and their shards hot. Cap per-request work and replicate hot adjacency lists.
  • Tail latency from fan-out. A page that calls twenty services is only as fast as the slowest. Use timeouts with partial rendering, and hedge requests to replicated stores. Caching in front of derived stores, covered in caching strategies, absorbs the hottest keys.

Trade-offs worth arguing about

The derived-data architecture buys independent scaling and fit-for-purpose stores at the cost of eventual consistency everywhere except the source of truth, plus many specialised systems to operate. LinkedIn could justify building Espresso, Venice, Pinot and a graph database; most companies should use existing equivalents such as a managed document store, a CDC tool feeding Kafka, a key-value store loaded by batch jobs, and an OLAP engine. Pull-based feeds make ranking flexible and storage cheap but put the work on the read path, which is only affordable when network degree is bounded. Query-time graph distance avoids huge precomputed sets but needs the graph in memory and careful sharding.

What to do next

  1. Identify your source-of-truth stores and list every derived view each one feeds, including caches and search indexes.
  2. Capture changes from the database commit log with a CDC tool, or adopt the outbox pattern, instead of dual writes.
  3. Key change events by entity ID so per-entity order is preserved, and make every consumer idempotent by storing the source offset.
  4. Measure consumer lag as event age and alert on it per derived view.
  5. Decide a read-your-writes policy for user edits and implement it explicitly.
  6. If you show relationship distance, prototype the intersection algorithm above on a sample of your graph and measure degree-3 cost before promising it.
  7. Choose pull or push for your feed from your degree distribution, not from what another company did.
Key takeaway: LinkedIn's architecture is organised around network distance and derived data. Writes land in a source-of-truth store, change capture turns each commit into an ordered event in Kafka, and stream and batch jobs build purpose-built serving stores for key-value features, analytics, search and the graph. Distance labels are computed at query time from in-memory adjacency lists using intersections, and the feed is gathered and ranked on read. Copy the pattern, not the bespoke systems: CDC, idempotent consumers, explicit read-your-writes and bounded graph queries.