Instagram Feed Generation System
A few weeks ago I sat down with a public system-design doc for an Instagram-style feed and asked myself a plain question: could I actually build this, not sketch it? Not a slide with boxes, arrows, and a bullet point that says “shard by user_id at scale” — an actually-running system where that sharding is real, where a follow edge, a like counter, and a celebrity’s post all have to survive contact with four independent databases at once without silently disagreeing with each other.
feed-generation-service is what came out of that: four independent Postgres shards, self-routing IDs, a real graph database for the social graph, and a script that proves two independently-implemented shard routers agree, byte-for-byte, on where a brand-new entity lands. This is the walkthrough of what I built and why — requirements, invariants, API, schema, diagrams, and the actual code, past the one sentence a system design interview usually stops at.
Six services, six languages, one design: each service is written in the language its role would realistically use in production (Go for the write path, Java for a Kafka consumer, Python for embeddings, Node for I/O-bound event handling, Rust for a stateless scorer). Every code sample below is pulled directly from the repo, not paraphrased.
Functional requirements
- Create a post (image/video/carousel) with a caption.
- Follow / unfollow another user.
- Like / unlike a post; comment on a post.
- Get
@mentionnotifications when named in a caption. - Generate a ranked, paginated feed combining posts from people you follow, celebrities you follow, and semantically similar content from people you don’t.
- Recall content a user hasn’t already been shown.
Non-functional requirements
Targets inherited from the source design doc, at 500M DAU:
| Requirement | Target |
|---|---|
| Feed read throughput | ~80,000 QPS |
| Feed read latency | p99 < 200ms |
| Write availability | A post commit must succeed even if every downstream consumer (fan-out, embedding, notifications) is degraded or down |
| Consistency | Eventual consistency for fan-out (a new post reaching a follower’s feed can lag by a bounded window); strong, read-your-own-writes consistency for “did I just like this” |
| Data locality | The common read for any entity (a user’s posts, a post’s comments, a user’s own feed) must resolve without a cross-shard query |
System invariants
Properties that must hold regardless of load, failure, or which shard is involved — the things a correctness review actually checks:
- An ID’s shard bits never change after minting. A post’s shard is fixed the instant it’s created; nothing ever moves a row to a different shard.
- A post always lives on its author’s shard. Never re-hashed independently, so “get this user’s posts” is always single-shard.
- Exactly one physical shard owns any given shard-ID range. No two shards can claim the same ID space.
- Go and Java always agree on shard placement for the same input. Verified, not assumed — see scripts/verify_shard_parity.sh.
- A follow edge and its derived facts (
followerCount,isCelebrity) commit atomically. Neo4j updates both in one transaction — never observably out of sync. - A like or comment is idempotent per (post, user). Retrying the same toggle never double-counts.
- Fan-out is at-least-once, never zero-out. A follower is never silently skipped; they land in the hot tier, the cold tier, or (for a celebrity’s post) an outbox — never nowhere.
API
The client-facing surface, all through Envoy at :8080:
| Method & path | Purpose |
|---|---|
POST /v1/users | Create a user; mints a self-routing ID |
POST /v1/posts | Create a post; {userId, mediaUrl, mediaType, caption} |
POST /v1/follow / POST /v1/unfollow | Follow graph mutation (Neo4j) |
GET /v1/users/{id}/following | Who this user follows |
POST /v1/posts/{id}/like / POST /v1/posts/{id}/unlike | Idempotent, rate-limited (20/user/post/min) |
POST /v1/posts/{id}/comments / GET /v1/posts/{id}/comments | Comment write (moderated) / read (stampede-protected) |
GET /v1/feed?userId=&cursor= | The ranked, paginated feed |
POST /v1/feed/reset-seen | Demo-only: clears seen-state history |
POST /v1/posts’s request and response shapes, straight from the handler:
type createPostRequest struct {
UserID string `json:"userId"`
MediaURL string `json:"mediaUrl"`
MediaType int16 `json:"mediaType"` // 1=Image, 2=Video, 3=Carousel
Caption string `json:"caption"`
}
// writeJSON(w, http.StatusCreated, postResponse(post))
And what a feed response item actually carries — enough for the client to render a post and know why it’s there:
type FeedItem struct {
PostID string
AuthorID string
MediaURL string
Caption string
LikeCount int64
CommentCount int64
Source string // in-network | celebrity | vector
Score float64
IsAd bool
CreatedAt time.Time
LikedByMe bool
}
Data modeling and schema
Postgres — identical schema on all 4 shards (db/shard-schema.sql), every primary key an application-minted Snowflake ID, never a database SERIAL:
CREATE TABLE IF NOT EXISTS users (
user_id BIGINT PRIMARY KEY,
username VARCHAR(64) UNIQUE NOT NULL,
last_active_at TIMESTAMP WITH TIME ZONE DEFAULT NOW()
);
CREATE TABLE IF NOT EXISTS posts (
post_id BIGINT PRIMARY KEY,
user_id BIGINT NOT NULL,
media_url TEXT NOT NULL,
media_type SMALLINT NOT NULL, -- 1: Image, 2: Video, 3: Carousel
caption TEXT,
created_at TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT NOW()
);
CREATE INDEX IF NOT EXISTS idx_user_posts ON posts (user_id, created_at DESC);
CREATE TABLE IF NOT EXISTS likes (
post_id BIGINT NOT NULL,
user_id BIGINT NOT NULL,
created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(),
PRIMARY KEY (post_id, user_id) -- naturally idempotent
);
CREATE TABLE IF NOT EXISTS comments (
comment_id BIGINT PRIMARY KEY,
post_id BIGINT NOT NULL,
user_id BIGINT NOT NULL,
body TEXT NOT NULL,
created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW()
);
CREATE INDEX IF NOT EXISTS idx_comments_post ON comments (post_id, created_at);
Notably absent: a follows table, and a like_count/comment_count column. Both are deliberate — covered below. The social graph is a separate model entirely, in Neo4j:
(:User {userId, username, followerCount, isCelebrity})
-[:FOLLOWS {createdAt}]->
(:User)
High-level design
flowchart TB
Client([Client]) --> GW[Envoy Gateway]
GW -->|writes| PI["post-ingestion-service (Go)"]
GW -->|reads| FA["feed-aggregation-service (Go)"]
GW --> UI[web-ui]
PI --> Sharded[("Postgres x4 · Neo4j\ngraph + counters")]
PI --> KF[["Kafka post-created"]]
KF --> FW["fanout-worker (Java)"]
KF --> VP["vector-pipeline (Python)"]
KF --> NS["notification-service (Node/TS)"]
FW -->|writes| Tiers[("Redis hot inbox\nBadgerDB cold tier")]
VP --> QD[("Qdrant")]
NS --> Sharded
FA --> Tiers
FA --> QD
FA -.gRPC.-> RK["ranking-service (Rust)"]
In plain terms, here’s what each box actually does:
- Envoy Gateway — the one door every client request walks through. It looks at the path and sends writes to
post-ingestion-service, feed reads tofeed-aggregation-service, and everything else to the static web UI. - post-ingestion-service (Go) — handles everything a user writes: new posts, follow/unfollow, likes, comments. One synchronous database commit per request, nothing else.
- fanout-worker (Java) — listens for new posts and decides who needs to see them next. This is where “push the post into the right inboxes” actually happens.
- vector-pipeline (Python) — reads the caption of every new post, turns it into a vector, and stores it so the system can later find “posts like this one.”
- notification-service (Node/TS) — watches new posts for
@mentionsand writes a notification for whoever got mentioned. - feed-aggregation-service (Go) — builds a user’s feed when they open the app: pulls candidates from several places at once, asks
ranking-serviceto sort them, and hands back a page. - ranking-service (Rust) — a pure function: given a pile of candidate posts and their engagement counts, returns them in the order they should be shown. No database, no memory of past calls.
And the data stores, also in plain terms:
- Postgres (4 shards) — the system of record. Real rows for users, posts, comments. Split four ways so no single database has to hold everyone.
- Neo4j — who-follows-whom, stored as a graph instead of a table, because “who follows this person” is a graph question, not a relational one.
- Redis — fast, disposable state: inboxes, engagement counters, rate limits, “has this user already seen this post.”
- BadgerDB — a cheap backup shelf for inbox entries belonging to people who haven’t opened the app recently. Slower to read than Redis, far cheaper to keep around.
- Kafka — the mail room. One post comes in, three departments (
fanout-worker,vector-pipeline,notification-service) each get their own copy and read it on their own schedule. - Qdrant — a similarity index over post captions, used to recall posts a user might like even if they don’t follow the author.
Write path: a user creates a post
sequenceDiagram
participant C as Client
participant PI as post-ingestion-service
participant PG as Postgres (author's shard)
participant K as Kafka
participant FW as fanout-worker
participant N4 as Neo4j
participant R as Redis (hot inbox)
participant BG as BadgerDB (cold tier)
C->>PI: POST /v1/posts
PI->>PI: mint ID, inheriting author's shard
PI->>PG: INSERT post
PG-->>PI: OK
PI->>K: publish post-created (FlatBuffers)
PI-->>C: 201 Created
K->>FW: post-created event
FW->>N4: is author a celebrity? who follows them?
N4-->>FW: follower list / isCelebrity flag
alt author is a celebrity
FW->>R: ZADD celebrity:outbox -- one write, any follower count
else normal user
FW->>FW: split followers active vs. dormant, per shard
FW->>R: ZADD feed:user:<id> for each active follower
FW->>BG: gRPC Append for each dormant follower
end
The client gets 201 Created the moment the single Postgres write lands — everything downstream of Kafka is fire-and-forget from the client’s point of view. The one decision worth pausing on is what fanout-worker does next, and this is the real branch, from FanoutService.java:
public void handle(PostCreatedEvent event) {
long authorId = event.userId;
long postId = event.postId;
long score = event.createdAt;
// Celebrity status lives on the Neo4j User node -- one graph query,
// no Postgres shard lookup needed for this decision at all.
if (socialGraph.isCelebrity(authorId)) {
hotInbox.pushToCelebrityOutbox(authorId, postId, score);
return;
}
List<Long> followers = socialGraph.getFollowers(authorId);
List<Long> activeFollowers = activity.filterActive(followers, activeWithinDays);
List<Long> dormantFollowers = new ArrayList<>(followers);
dormantFollowers.removeAll(activeFollowers);
for (Long followerId : activeFollowers) {
hotInbox.pushToInbox(followerId, postId, score, feedInboxMaxItems);
}
// Dormant followers used to simply be dropped here. Now they go to
// the cold tier instead -- see the caching section below.
for (Long followerId : dormantFollowers) {
coldTier.append(followerId, postId, score); // gRPC to feed-aggregation-service
}
}
Without the celebrity branch, one post from a million-follower account would mean a million individual writes on the critical path of fan-out. With it, a celebrity’s post costs exactly one write regardless of follower count, and their followers pull the post from a shared outbox when they read their own feed instead.
Why not just push, or just pull
This is really three different answers to the same question — “how does a post get into a feed” — chosen per follower based on how likely that write is to ever get read:
- Pure push (write a copy into every follower’s inbox, unconditionally) makes reads cheap — a feed open is just “read my precomputed inbox” — but wastes real work on followers who rarely open the app: bandwidth and storage spent on a write that may sit unread for months, or never be read at all.
- Pure pull (compute the feed from scratch on every read, for everyone) wastes nothing on inactive users — no work happens until they actually ask — but makes every read expensive: scatter-gathering across everyone a user follows, on every single feed open, including for the heavy, frequent-checking majority who’d benefit most from a precomputed answer.
Neither extreme is right for every follower of the same post, so the split above picks per-follower, not globally:
- Active followers get pushed, because they’re likely to read soon — the precomputation pays for itself.
- Dormant followers aren’t pushed to the expensive tier, but aren’t skipped either. They get one cheap append to BadgerDB (disk, roughly two orders of magnitude cheaper per GB than the Redis write an active follower gets). That’s not free, but it’s cheap enough that paying it now beats the alternative: if it were skipped entirely (pure pull for anyone dormant), a user who follows hundreds of accounts and comes back after months would trigger a full cross-follow-graph scatter-gather the moment they opened the app — exactly the expensive-read problem pure pull has, just deferred and concentrated into one bad moment instead of amortized.
- Celebrities are pure pull, deliberately — the one case where even the cheap per-follower write doesn’t scale, since a million followers times any per-follower cost, however small, still adds up on the critical path of a single post.
So the answer to “isn’t pushing to inactive users wasteful” is: yes, which is exactly why they don’t get the expensive push — they get the cheapest write that still avoids forcing an expensive read later.
Read path: a user opens their feed
sequenceDiagram
participant C as Client
participant FA as feed-aggregation-service
participant R as Redis
participant BG as BadgerDB
participant N4 as Neo4j
participant Q as Qdrant
participant PG as Postgres shards
participant RK as ranking-service
C->>FA: GET /v1/feed?userId=
par fan-in, three sources concurrently
FA->>R: hot inbox (ZREVRANGE)
FA->>BG: cold-tier fallback if hot inbox is empty
FA->>N4: which celebrities does this user follow?
FA->>Q: ANN search against a "taste vector"
end
FA->>FA: merge, dedupe, drop anything seen in the last 48h
FA->>PG: hydrate post metadata, grouped by shard
FA->>RK: gRPC Rank(candidates)
RK-->>FA: ranked list with scores
FA->>FA: cap 2 posts/author, insert ads
FA-->>C: feed page + next cursor
Three candidate sources run in parallel because none of them depend on each other: the user’s own hot inbox, the outboxes of celebrities they follow, and a vector search against the average embedding of their recent activity. Ranking itself is a single composite score, computed by ranking-service from recency decay and normalized engagement counts — the real formula, from main.rs:
fn score_candidate(weights: &Weights, c: &Candidate) -> RankedItem {
let (p_like, p_comment, p_share, p_dwell, p_hide) = estimate_probabilities(c);
let score = weights.like * p_like
+ weights.comment * p_comment
+ weights.share * p_share
- weights.hide * p_hide
+ weights.dwell * p_dwell;
// ...
}
ranking-service never sees a user_id’s or post_id’s shard — it’s a pure function over pre-hydrated inputs, which is exactly why it’s stateless and trivially language-agnostic.
Out-of-network recall: the taste vector, ANN, and HNSW
The in-network and celebrity sources above only ever return posts from people a user already follows. On its own, that’s a discovery dead end: a new user who follows three accounts gets a feed of three accounts’ worth of posts, forever, until they manually go find more people — and even a long-time user never sees anything from an account they haven’t already found. Every large feed product solves this with some form of out-of-network recall: content from accounts you don’t follow, surfaced because it’s similar to what you already engage with. That’s the business reason this source exists at all — it’s not a ranking nicety, it’s what makes a feed keep giving people a reason to open the app instead of running dry.
Here’s how that similarity is actually computed. Every post’s caption gets turned into a 384-number vector (all-MiniLM-L6-v2, run in vector-pipeline) — captions that mean similar things end up as vectors that sit close together in that 384-dimensional space. To recommend for a specific user, feed-aggregation-service builds a taste vector: it looks up that user’s own recent posts plus their followees’ recent posts, and averages those posts’ embeddings together into one vector that represents “the kind of thing this person posts or follows”:
seedPosts, _ := s.postMeta.RecentByAuthors(ctx, seedAuthors, s.cfg.TasteSeedPostLimit)
tasteVector, ok := s.vector.DeriveTasteVector(ctx, seedIDs)
cands, _ := s.vector.Search(ctx, tasteVector, s.cfg.VectorCandidateLimit)
That search is where ANN (Approximate Nearest Neighbor) comes in. The exact way to find “which posts are most similar to this taste vector” is to compare it against every single post vector in the system and sort by distance — correct, but far too slow once there are millions of posts. ANN trades a small, usually-unnoticeable amount of accuracy for a huge amount of speed: instead of the mathematically exact closest matches, it finds matches that are almost certainly among the closest, in a fraction of the time. HNSW (Hierarchical Navigable Small World), the specific algorithm Qdrant uses here, does this by building several layers of connections between vectors, like a road network: a sparse top layer of long “highway” links lets a search jump straight into roughly the right neighborhood in a few hops, and each layer below it has denser, shorter links that narrow the search down street by street, until it lands among the true nearest neighbors — without ever touching most of the posts that were never going to be close matches in the first place.
One business rule worth being explicit about, because it’s easy to assume otherwise: there’s no reserved quota of “N out of the 20 items on a page must be out-of-network.” VECTOR_CANDIDATE_LIMIT (250) just sets how many similarity candidates are eligible to compete for a spot — every one of them still goes through the same composite ranking score as in-network and celebrity candidates, and only the highest-scoring ones survive into the final page. That’s a deliberate choice: it means discovery content never crowds out a genuinely great in-network feed, but it also means a user whose follows post constantly could see little to no out-of-network content on a given page, purely because nothing beat it on score. A real production feed commonly hard-reserves a handful of slots for exploration regardless of score, specifically so discovery can’t be starved out entirely — a real trade-off this repo makes visibly rather than hiding it.
How scrolling actually works
Infinite scroll looks like “give me page 2, then page 3” from the client’s side, but this system deliberately never implements it as OFFSET/LIMIT against a live, constantly-changing candidate list. If it did, a new post landing in your inbox between two scroll events would shift everyone after it by one position, and you’d either see a duplicate or skip a post entirely — pagination drift, not a bug in the traditional sense, just OFFSET doing exactly what it’s documented to do against a list that isn’t stable.
Instead, each response carries an opaque, encrypted cursor, and scrolling further is just the client echoing it back:
end := start + s.cfg.PageSize - 2 // reserve 2 organic slots for ad insertion
page := diversified[start:end]
// ... build the response ...
nextCursor, _ := s.cursor.Encode(domain.FeedCursor{Offset: end})
return FeedResult{Items: finalItems, NextCursor: nextCursor}, nil
That Offset is sealed with AES-256-GCM before it ever reaches the client:
func (c *CursorCodec) Encode(cursor domain.FeedCursor) (string, error) {
// ... AES-256-GCM seal ...
ciphertext := gcm.Seal(nonce, nonce, plaintext, nil)
return base64.RawURLEncoding.EncodeToString(ciphertext), nil
}
So a scroll event is just: client calls GET /v1/feed?userId=&cursor=<opaque>, the server decrypts it back into an offset, re-runs the same fan-in → rank → diversify pipeline, and slices out the next window starting where the last one left off. Two things make this actually correct, not just opaque:
- The offset is into a freshly re-ranked list, not a frozen one. Every scroll event re-derives candidates and re-ranks them — the cursor only remembers how far in you were, not what was there.
- Every post actually returned gets marked seen (
seen:user:<id>, 48h TTL) before the response goes out. That’s what stops a pull-to-refresh or a re-opened app from showing you the same post twice, independent of whatever the cursor offset says.
The client never needs to know or guess an offset — it just holds onto whatever string the last response gave it and sends it back unchanged, which also means the offset’s meaning can change shape later without breaking any client that’s already deployed.
Why the social graph isn’t in the sharded database
A follow relationship connects two users who can each be on any of the four Postgres shards. The first version of this system sharded follows relationally anyway — storing each edge twice, once on each endpoint’s shard, so both “who do I follow” and “who follows me” stayed single-shard. It worked, and it was still the wrong tool: every follow or unfollow became a hand-rolled distributed transaction across two independent databases, with no shared transaction and a real, permanent possibility of the two copies drifting apart.
Moving the graph to Neo4j mirrors what Meta’s own TAO does for the same reason: a purpose-built graph store, not a cleverly-sharded relational table, because relational sharding and graph traversal want opposite things from a key.
The hash ring and self-routing IDs
Every entity ID is a 64-bit value that carries its own shard number, minted by pkg/sharding:
// 41 bits timestamp | 8 bits shard ID | 14 bits sequence.
const (
sequenceBits = 14
shardIDBits = 8
shardIDShift = sequenceBits
timestampShift = sequenceBits + shardIDBits
maxSequence = (1 << sequenceBits) - 1
MaxShardID = (1 << shardIDBits) - 1
)
func (g *IDGenerator) Next() int64 {
// ...
return (ts << timestampShift) | (int64(g.shardID) << shardIDShift) | g.sequence
}
// ExtractShardID recovers the shard that minted id -- a pure bit-shift,
// never a network call or a lookup table.
func ExtractShardID(id int64) int {
return int(id>>shardIDShift) & MaxShardID
}
That covers existing entities: their shard is baked into their own ID forever. The one thing that’s actually mutable is where a brand-new signup lands, and that’s the hash ring’s only job:
func NewRing(cfg Config) *Ring {
r := &Ring{numShards: len(cfg.Shards)}
for _, shard := range cfg.Shards {
for v := 0; v < cfg.VirtualNodesPerShard; v++ {
// Vnode index BEFORE shard ID matters: CRC32 is a linear code,
// and "shard-<id>-vnode-<v>" (a constant prefix, only the
// low-order counter varying) measured ~29% max deviation
// across shards. "vnode-<v>-shard-<id>" breaks that
// correlation and brings it under 2%.
key := fmt.Sprintf("vnode-%d-shard-%d", v, shard.ID)
r.nodes = append(r.nodes, vnode{
hash: crc32.ChecksumIEEE([]byte(key)),
shardID: shard.ID,
})
}
}
sort.Slice(r.nodes, func(i, j int) bool { return r.nodes[i].hash < r.nodes[j].hash })
return r
}
func (r *Ring) ShardForNewEntity(placementKey string) int {
h := crc32.ChecksumIEEE([]byte(placementKey))
idx := sort.Search(len(r.nodes), func(i int) bool { return r.nodes[i].hash >= h })
if idx == len(r.nodes) {
idx = 0 // wrap around the ring
}
return r.nodes[idx].shardID
}
150 virtual nodes per shard exist so that with only 4 real shards, new-signup placement still spreads roughly evenly instead of landing on 4 arbitrary hash buckets. This exact algorithm is reimplemented independently in Java for fanout-worker; a script feeds 10,000 sample IDs through both and diffs the results, because if the two ever silently disagreed about placement, one service could create a user on shard 2 while another writes related data to shard 3 — silent corruption, with no error at the moment it happens.
Hot/cold tier caching
Every active user’s feed inbox lives in Redis — fast, and roughly two orders of magnitude more expensive per GB than disk. At real scale, keeping every user’s inbox in RAM regardless of whether they opened the app yesterday or a year ago means paying RAM prices for data that, for the dormant majority, is read almost never. The fan-out code above already writes dormant followers to a cheaper tier instead of dropping them; the read side is what decides when to pull them back:
func (s *FeedService) inNetworkCandidates(ctx context.Context, userID string) []domain.Candidate {
exists, err := s.hotInbox.Exists(ctx, userID)
if exists {
cands, _ := s.hotInbox.GetInbox(ctx, userID, s.cfg.InNetworkCandidateLimit)
return cands
}
// Hot tier miss: this user has been dormant. Fall back to the cold
// tier and, if we find anything, promote it back to the hot tier --
// a read means they're active again now.
cold, err := s.coldInbox.Get(ctx, userID)
if err != nil || len(cold) == 0 {
return nil
}
if err := s.hotInbox.Promote(ctx, userID, cold); err != nil {
slog.Warn("failed to promote cold-tier inbox to hot tier", "userID", userID, "error", err)
}
return cold
}
A returning-after-months user’s first request pays one extra, disk-backed (but still local) read — instead of either losing their accumulated fan-out history entirely, or the system having paid RAM cost to keep their empty inbox warm the whole time nobody was looking at it. The cold tier itself is BadgerDB, a pure-Go LSM-tree store standing in for the reference design’s RocksDB — the same substitution CockroachDB made for its own Pebble engine, and for the same reason: no CGO, no native toolchain dependency.
What’s simplified relative to real production scale
Four shards instead of hundreds, one Neo4j instance and one Kafka broker instead of replicated clusters, heuristic ranking instead of a trained model — every one of these is a named, reasoned trade-off, not an oversight. The full breakdown, written for system-design-interview prep, is in the repo’s trade-offs doc, including a prioritized “what I’d build first toward production” list. The complete design write-up this post is based on is in doc/DESIGN.md.