Track C · Design interviews

System Design: AI Products on Real Infrastructure

Many "AI system design" rounds at product companies are really distributed-systems rounds with a model somewhere in the middle. The interviewer wants to know whether you can keep a vector index consistent with a source database, stop a retry storm from burning a provider quota, pause a workflow for three days waiting for a human, or meter tokens accurately enough to bill on. This page covers a framework for hybrid designs, a refresher on the recurring building blocks, ten case studies with diagrams and back-of-envelope numbers, cross-cutting deep dives and a question bank. The AI-centric designs (support agent, enterprise RAG, coding agent, deep research, voice, memory assistant) are in B6 and aren't repeated here.

TL;DR: the 8–12 things to be able to say out loud

  • Treat the LLM as a slow, expensive, rate-limited, non-deterministic remote dependency. Timeouts, retries with jitter, circuit breakers, queues, caches and idempotency all apply, but a retry can return a different answer.
  • Place each model call deliberately: synchronous, async just after the response, or offline batch. Most AI work (extraction, enrichment, evals, summaries) belongs off the hot path.
  • Little's law sizes everything. At 8 s per call, 500 req/s means 4,000 calls in flight, and token quota, not CPU, is usually the binding limit.
  • Kafka gives per-key ordering, replay and fan-out. Key by the entity whose order matters and make consumers idempotent; delivery is at-least-once.
  • Never dual-write. Use CDC or a transactional outbox so the database and the stream can't disagree.
  • Exactly-once is an end-to-end property you build from at-least-once delivery plus idempotent writes (deterministic IDs, version checks, unique constraints).
  • Store model outputs; don't recompute them, tagged with the model and prompt versions that produced them.
  • Long-running AI jobs want a durable workflow engine: model and tool calls as retried activities, human waits as timers and signals.
  • Re-embedding is a blue/green migration: new index, dual-write, backfill, evaluate, flip an alias, keep the old one for rollback.
  • Multi-tenancy means fairness at every queue: per-tenant token buckets, fair scheduling, tenant-led keys, noisy-neighbour isolation.
  • Quality is a non-functional requirement with its own pipeline: traces to a stream, sampled spans to async judges, scores to alerts.
  • Do the cost arithmetic out loud. Human review minutes often dominate, so straight-through rate is the real lever.

Part 0. A framework for hybrid AI + infrastructure design

The classic high-level design template is still the backbone: requirements, estimates, API, data model, high-level design, deep dives, bottlenecks. Kleppmann, DDIA What changes is that one or more boxes in the diagram behave very unlike a database or a microservice. An LLM call takes seconds instead of milliseconds. It costs cents instead of fractions of a cent. It's metered by tokens per minute by a third party, it can return a different answer for the same input, and "200 OK" tells you nothing about whether the answer was right. The table shows how each step of the template shifts.

Template stepClassic versionWhat changes when AI components are involved
Functional requirementsFeatures and user flows.Add what a correct output looks like and who judges it. Agree whether results can stream.
Non-functional requirementsLatency, availability, durability, consistency.Add quality (an eval target), cost per request, freshness of derived data, data governance (what may reach a third-party model, residency) and degradation behaviour when the model is slow or down.
EstimatesQPS, storage, bandwidth.Add input and output tokens per request, tokens per minute against quotas, concurrency via Little's law, GPU-hours and dollars per day.
APIRequest-response.Streaming (SSE) for generation; async job APIs (202 plus polling or webhooks) for slow work; idempotency keys on every mutation, because clients retry slow requests.
Data modelEntities and access patterns.Model outputs as versioned data (model, prompt version, input hash), embeddings tagged with their model, trace and eval records, lineage for deletion.
High-level designServices, stores, queues.Model calls behind a gateway; sync vs async placement; a tracing and eval path; quota-bound work behind queues so bursts become backlog, not errors.
Deep dives and bottlenecksSharding, caching, hot spots.Idempotency of AI side effects, rate-limit-aware scheduling, reprocessing on model change, cascades, review queues. The bottleneck is usually token quota or GPUs, then human review, then index memory.

Where to put the model call

The most consequential decision in a hybrid design is where each model call sits relative to the user's request. There are three placements, and a good design uses all three for different jobs.

1. Synchronous (user waits) Client API service LLM gatewaystream, timeout chat answers, agent turns, query rewrite budget: seconds; must degrade gracefully 2. Async after the response (user doesn't wait) API service DB + outboxcommit first Kafka / queue AI workersrate-limited pool Derived storesindex, memory, scores embedding, memory extraction, moderation of published posts, online evals, summaries; budget: seconds to minutes 3. Offline batch (scheduled) Lake (S3/Parquet) Spark / Batch API Bulk load backfills, re-embeds, nightly user profiles, eval sweeps; ~50% cheaper via batch APIs
Three placements for model calls. Push work down the list whenever the product allows it: each step down buys cost, throughput and failure isolation at the price of freshness.

Rule of thumb: a model call goes on the synchronous path only if the user is waiting for its output, and then it gets a latency budget, a timeout, a fallback and streaming. Everything else goes behind a durable queue, where a burst becomes backlog instead of a wave of 429s. For bulk work, OpenAI and Anthropic both offer batch endpoints at about half price; OpenAI's has a 24-hour completion window and a separate rate-limit pool. OpenAI Batch Anthropic Batches

The LLM as a distributed-systems dependency

Intuition

Size AI systems the way you'd size a call centre, not a web server. A web request occupies a thread for milliseconds. An LLM call occupies a stream, a slot of provider concurrency and a chunk of token quota for seconds. Little's law (L = λW) turns that into a number: 300 req/s × 6 s average duration = 1,800 concurrent streams. Little's law Async I/O makes holding 1,800 streams cheap for your servers; the provider's quota is what you run out of.

Quality, evals and cost as design elements

Quality needs the same treatment as availability: a target, a measurement path and an alert. Every design here includes four hooks: an async trace record per model call (inputs, outputs, versions, tokens, latency, cost); a sampling tap to evaluators (judges, heuristics, humans); a feedback join attaching user signals to the trace by request ID; and versioned prompts and models, so a regression maps to a change. Case 5 is the pipeline behind these hooks; eval methods are in B5. Cost per request belongs in the estimates step, said out loud ("3k input tokens, 60% cached, 300 output, so about X cents; at 2M requests a day, $Y; input dominates, so caching is the first lever").

Interview angle

A common probe is "your provider's rate limit is 10M input tokens per minute; is that enough?" Strong answers convert units on the spot: 10M per minute is about 167k tokens per second, which at 3k input tokens per request is roughly 55 requests per second, uncached. Then they reach for levers: prompt caching (cached tokens may not count against the limit), moving work to batch, routing simple traffic to a smaller model with its own quota, multiple accounts or regions, provisioned capacity, or self-hosting part of the load.

Latency and cost numbers every AI system designer should know

Treat every row as an order of magnitude for back-of-envelope work, not a benchmark. Hardware, models, prompt sizes and providers vary a lot, and LLM speed numbers change every few months. Check a live comparison such as Artificial Analysis for current provider figures.

OperationRough figureDesign implication
Hosted LLM time-to-first-token, short prompt, non-reasoning model~0.2–1 sStream it. Reasoning models can take far longer.
TTFT with a very long prompt (~100k tokens), uncachedseveral secondsPrompt caching cuts it substantially.
Decode speed per stream~30–200 tokens/s (small models and specialised hardware faster)A 500-token answer takes ~3–15 s.
Embedding API call, batch of ~100 short texts~0.1–0.5 sBatch; quota, not latency, bounds throughput.
Self-hosted small embedding model, one modern GPU~1k–10k chunks/sWins for big re-embeds.
Cross-encoder rerank of ~50 candidates on GPU~10–50 msFits a 300 ms budget.
Small transformer classifier per item (batched, GPU)~1–10 msFirst stage of cascades.
In-memory HNSW query, millions of vectors, recall ~0.95~1–10 msMemory is the cost: 1B × 768-d float32 ≈ 3 TB. ann-benchmarks
Object-storage-backed vector search (warm cache)tens of ms; cold queries slowerMuch cheaper per vector; Notion reported 50–70 ms. Notion 2026
Redis GET/SET in the same AZsub-ms; ~100k+ ops/s per nodeFine for quota checks; hot keys still bottleneck.
Postgres indexed point read (hot data)~1 ms; one primary handles thousands to tens of thousands of writes/sShard before it saturates.
Kafka produce→consume, same clustera few ms to tens of ms; tens to hundreds of MB/s per brokerConsumers, not Kafka, bottleneck.
S3 request rate per prefixat least 3,500 writes/s and 5,500 reads/s; small-object latency ~100–200 msKeep S3 off tight sync paths. AWS S3 docs
DynamoDB per-partition throughput3,000 read units/s and 1,000 write units/sKey design matters. AWS DynamoDB docs
ClickHouse batch insert / aggregate scaninsert in batches of thousands+ rows; scans of hundreds of millions of rows/s per serverNever insert row by row.
Cross-region round trip~60–150 msResidency and failover add to TTFT.
Human review of one item~30 s–3 minOften the largest cost line.
Go deeper

Part 1. Building blocks refresher

A compressed refresher on the infrastructure that recurs in the case studies, focused on what interviewers probe: ordering, delivery guarantees, failure behaviour and scaling limits. The DDIA chapters on replication, partitioning, transactions and stream processing are the best single source for more depth. DDIA 2e

Kafka: the log as the backbone

A Kafka topic is split into partitions: append-only, replicated, ordered logs. Ordering is guaranteed only within a partition, and the default partitioner hashes the record key, so all records with one key land in one partition in order. Kafka design docs Key choice is therefore a design decision: document_id if edits to a document must apply in order, conversation_id for chat. Adding partitions later remaps keys and breaks in-flight ordering, so over-provision up front.

May be out of date

Kafka features relevant to AI workloads are moving: tiered storage to object stores (KIP-405) and "share groups" (KIP-932), which add queue-style per-record acknowledgement and delivery counts. KIP-405 KIP-932 I'm not certain which release made share groups production-ready, so check current release notes.

CDC with Debezium, and the transactional outbox

Writing to Postgres and then publishing to Kafka (a dual write) leaves them inconsistent after a crash between the two, and concurrent writers can publish out of commit order. Both standard fixes derive events from the database's commit log:

Gotchas: a replication slot retains WAL until the connector confirms it, so a stalled connector can fill the primary's disk; alert on slot lag. Large rows don't belong in events: publish a pointer to object storage (the "claim check" pattern), as Shopify did for records over Kafka's 1 MB limit. Shopify 2021

Serviceone txn Postgres documents outbox WAL (commit order) Debeziumreplication slot topic doc-changes P0: doc 7, doc 12… P1: doc 3, doc 9… P2: doc 5, doc 8… P3: … group: embedder4 consumers group: analyticsown offsets key = doc_id → same partition → per-document order preserved
Outbox plus CDC. The database commit is the single source of truth; Kafka is a faithful, ordered, replayable projection of it. Each consumer group tracks its own offsets.

Stream processing vs batch

Flink and Kafka Streams run continuous, stateful computations. The concepts that come up are event time vs processing time, watermarks (the processor's estimate that no older events are still coming, which lets windows close while tolerating some lateness), windows (tumbling, sliding, session), and checkpointed keyed state. Flink: time Flink: state The Dataflow paper is the clearest treatment of correctness vs latency vs cost for out-of-order data. Akidau+ 2015 Kafka Streams runs the same ideas as a library inside your service. Kafka Streams Batch (Spark, batch APIs) stays right for backfills and re-embeds. Spark AI systems usually run both into one sink with idempotent, versioned writes. Notion described exactly this for vector search: Spark batch jobs for onboarding workspaces and Kafka consumers for live edits. Notion 2026

Queues vs logs

Queue (SQS, RabbitMQ, Redis job queues)Log (Kafka, Kinesis, Redpanda)
ConsumptionCompeting consumers; a message is deleted when acknowledged.Offsets per group; records kept for the retention period.
OrderingNone (SQS standard) or per message group (FIFO).Per partition, by key.
Per-message retryNative: visibility timeout, redelivery, DLQ. SQS DLQDIY: a failed record blocks its partition unless moved to retry topics or a DLQ. Uber 2018
Replay / parallelismNo replay; scale workers freely.Replay by rewinding; parallelism capped by partitions.
Best AI fitIndependent slow tasks: OCR a page, one extraction, judge one trace.Change streams, event sourcing, fan-out to derived views, metering.

SQS suits long LLM tasks because a worker can extend a message's visibility timeout while it works, up to 12 hours from first receipt. SQS visibility Real systems often combine both: after Redis memory exhaustion caused an outage, Slack put Kafka in front of its Redis job queue as a durable buffer, with a rate-limited relay feeding Redis. Slack 2017

Workflow engines for long-running AI jobs

For jobs with many steps, durations of minutes to days, flaky dependencies and a need to survive deploys, a durable workflow engine beats hand-rolled state machines. In Temporal, a workflow is ordinary code whose progress is recorded as an event history and replayed after a crash, which is why workflow code must be deterministic. Side effects (LLM calls, tools, writes) run as activities with retry policies. Temporal workflows Activities Retry policies Durable timers sleep for days without holding a worker; signals and updates deliver external input such as approvals; continue-as-new bounds history length. Timers Message passing Continue-As-New See B4 for durable execution in agent frameworks; case 4 is the platform view.

Caches

Cache-aside or write-through, with jittered TTLs. The classic failure is a stampede: a hot key expires and thousands of requests hit the backend at once, which hurts when the backend is a multi-second LLM call. Cache stampede Defences: request coalescing (Discord's data services coalesce identical reads Discord 2023), leases (Facebook's memcache Nishtala+ 2013), and probabilistic early recomputation Vattani+ 2015. In Redis, Lua scripts run atomically (the basis of correct rate limiters), sorted sets make delayed-job schedulers, and Streams give a light log with consumer groups. Redis scripting Redis Streams AI-specific caches: exact-match response caches keyed by (model, prompt version, input hash, parameters); provider prompt caches, which reward stable prefixes; embedding caches keyed by content hash; and semantic caches, which carry correctness risk and should be tenant-scoped and opt-in. GPTCache

Databases and storage

Sharding, partitioning and multi-tenancy

Hash partitioning spreads load; range partitioning keeps ranges together but risks hot spots; consistent hashing limits data movement when nodes change. DeCandia+ 2007 A practical trick is many logical shards on fewer physical databases: Notion sharded Postgres by workspace ID into 480 logical shards over 32 databases. Notion 2021 In AI systems the natural shard key is usually the tenant, because retrieval, memory and ACLs are tenant-scoped. Tenancy models are silo (dedicated resources), pool (shared, with tenant IDs everywhere) and bridge (pool for small tenants, silo for large or regulated ones). Noisy neighbours are contained by per-tenant quotas, fair queuing and shuffle sharding. AWS: shuffle sharding AWS: fairness

The reliability toolkit

Consistency models and read-your-writes

From strongest to weakest: linearizable, sequential, causal, session guarantees (read-your-writes, monotonic reads), eventual. Jepsen The one AI users notice is read-your-writes: after I send a message or state a fact, the next turn must see it. Jepsen: RYW Implement it with leader or strongly consistent reads after a write, a session log position replicas must reach, a session cache of recent writes, or sticky routing. Derived views (vector index, memory, search) are eventually consistent by construction, so the design must say what the next turn does when they lag. Usually it also reads the raw recent data.

Which building block for which job

NeedTypical choiceGuaranteeWatch out for
Propagate DB changes to an index or memoryDebezium CDC → Kafka, or outboxCommit-ordered per key, replayableSlot growth, schema changes, big rows
Slow independent AI tasksSQS / RabbitMQ / share groupsAt-least-once, per-message retry, DLQVisibility timeout shorter than the task
Multi-step jobs with waitsTemporal or similarDurable state, retries, timers, signalsDeterminism, history size, big payloads
Windowed aggregates (metering, alerts)Flink / Kafka StreamsExactly-once state via checkpointsLate events, state size
Backfills, re-embedsSpark, batch APIsThroughput, lower priceCut-over with the live stream
Quotas, dedup keys, hot tailsRedisAtomic scripts, sub-msHot keys, fail-open policy
Chat historyDynamoDB / Cassandra / ScyllaDBPartition-local ordered readsHot partitions
Vectorspgvector → vector DB → object-storage enginesANN recall/latency trade-offFiltered recall, churn, memory
Traces, usage, scoresClickHouse + object storageFast aggregates, cheap retentionSmall inserts, updates
Interview angle

Interviewers test whether you know what Kafka doesn't do: per-message retries, delayed delivery, ordering across partitions, exactly-once for external side effects. Pair each with its fix: retry topics and DLQs, a timer service, key design, idempotent sinks.

Go deeper

Part 2. Case studies

Each case follows the same shape: clarifying questions, requirements, estimates, API, data model, diagram, data flow, two or three deep dives, failure modes and trade-offs. All numbers are illustrative assumptions you'd agree with the interviewer. The point is to show the arithmetic, not to defend the inputs. In an interview you'd cover perhaps half of each case in 45 minutes and go deep where the interviewer pulls you.

Case 1: Real-time RAG ingestion pipeline

Prompt: "Our assistant answers questions over documents in several source systems (a Postgres-backed product, a wiki, a ticketing tool) that change constantly. Design the pipeline that keeps the vector index fresh." B6 case 2 covers the query side of enterprise RAG; this case is the write path.

Clarifying questions

Requirements

Functional: ingest creates, updates and deletes; chunk, embed, upsert; propagate ACL changes; support backfills and model migrations. Non-functional: p95 commit-to-searchable under 60 s; deleted documents unservable within seconds; idempotent and replayable; cost proportional to change volume; streaming keeps working during a full re-embed.

Estimates

Assumptions: 50M documents × 8 chunks × ~400 tokens = 400M chunks and 160B tokens; 2% of documents change per day.

API and data model

Internal APIs: POST /sources/{id}/backfill (start a snapshot job), POST /indexes (create an index with an embedding model and chunker version), POST /indexes/{alias}:swap (cut over). The core data lives in three places:

doc_registry (Postgres, sharded by tenant) vector index (per alias: index_v1, index_v2 …) tenant_id, doc_id PK id = doc_id:chunk_hash source_version (LSN / row version, monotonic) vector content_uri (s3://… claim check) doc_id, tenant_id, acl_groups[], source_version chunk_hashes[] embed_model, chunker_version acl_version, acl_groups[] deleted_at (tombstone kept N days) embed_model, chunker_version

The doc registry is the key to correctness. It records, per document, the latest version the pipeline has applied, so stale or reordered events can be detected and dropped, and deletes leave a tombstone that later out-of-order updates can't resurrect.

Product Postgres Wiki / SaaS Data lake snapshot Debezium Webhook +poll connectors Backfill job Kafka: livekey = doc_id Kafka: backfilllow priority Chunk + embedworkersversion check, diff Doc registry (PG) Redis: quota +embedding cache Embedding APIor GPU pool (batched) vector store alias "prod" → v1live v2building DLQ + replay tool Delete / ACLfast path (no embed) Query serviceACL filter + deny-list Live lane gets reserved embedding quota; backfill uses what's left. Workers are idempotent: deterministic chunk IDs plus a version check in the registry mean replays, duplicates and reordering are harmless.
Write path for a continuously fresh vector index. The registry and deterministic chunk IDs give an effectively exactly-once result from at-least-once delivery.

Data flow

  1. A row commits; Debezium emits a change event keyed by doc_id with the commit's log position as source_version. SaaS sources arrive by webhook (verified, deduplicated) plus periodic polling for missed ones. Large bodies go to S3; the event carries only the URI.
  2. A worker loads the registry row and drops the event unless its version is newer than the applied one and newer than any tombstone.
  3. It chunks and hashes the content. Only chunks whose hash is new (and not in the content-hash embedding cache) get embedded, in batches of hundreds, with a permit from a shared, token-aware Redis rate limiter.
  4. It upserts vectors with deterministic IDs (doc_id:chunk_hash), deletes IDs of vanished chunks, commits the new version to the registry with a compare-and-set, and only then commits the Kafka offset.
  5. Delete and ACL-only events take a fast path with no embedding: registry tombstone, deny-list entry, and a filtered update or delete of vectors.

Deep dive 1: ordering, idempotency and deletes

At-least-once delivery duplicates events, retries and rebalances reorder them, and the backfill lane can race the live lane. Three mechanisms make the result correct anyway: per-document ordering from keying by doc_id; monotonic versions that make the apply step a conditional write ("apply version 42 only if the current version is below 42"); and deterministic IDs that make upserts repeatable. Together they make transport-level exactly-once unnecessary; a duplicate costs a few embedding calls, mostly absorbed by the cache.

Deletes are a security issue, not a freshness issue. Debezium emits a delete event plus a tombstone. Debezium The pipeline writes a versioned registry tombstone (so a late update can't resurrect the doc), pushes the ID onto a Redis deny-list with a TTL longer than worst-case lag, and deletes vectors by filter. The query service drops deny-listed hits, closing the gap between commit and physical deletion, and for sensitive corpora re-checks top results against the source of truth. Response and semantic caches need short TTLs or invalidation too. Debezium now ships an embeddings transform and a Milvus sink that propagate deletes in a simple pipeline. Debezium 2025

Deep dive 2: re-embedding with zero downtime (blue/green indexes)

A new embedding model invalidates every vector, because vectors from different models aren't comparable. Create index_v2 with the model and chunker version in its config; dual-write the live stream to both indexes; backfill v2 from a consistent snapshot (lake, compacted topic or incremental snapshot) through the batch lane, where version checks make overlap with live writes safe; evaluate recall, answer quality and latency on a golden set and shadow traffic; then flip an alias atomically, which Qdrant and Elasticsearch both support for this purpose. Qdrant aliases Elastic aliases The query service must switch query-embedding models at the same moment, so it reads the model name from the index config. Keep v1 for a rollback window. Budget double storage and write load during the migration. Notion describes switching embedding models during a vector-store migration and later moving to self-hosted open-source embedding models. Notion 2026

Deep dive 3: rate limits, lanes and ACLs

The provider quota is shared across all workers, so limit centrally: a Redis token bucket charged in tokens, with the live lane guaranteed a share (say 30%, burstable to 100%) and backfill taking what's left. On 429, back off with jitter and lower the refill rate (additive increase, multiplicative decrease). Batch by token count. Shopify found bursty input left its streaming embedder with bundles of one element, defeating GPU batching, so a short micro-batching window can pay off. Shopify 2024 Figma went the other way and debounced index updates to every few hours, cutting the data processed to about 12%. Figma 2024 Freshness is a dial. Permissions never trigger re-embedding: store group or principal IDs (not user lists) as filterable metadata, update them with metadata-only writes, and resolve a user's groups at query time.

Failure modes and alternatives

Alternatives: pgvector in the source database removes the pipeline for Postgres-resident data and makes deletes and ACLs transactional, a strong answer at small to medium scale. Nightly batch rebuilds are fine when a day of staleness is acceptable. Embedding inside Kafka Connect is quick to stand up but gives less control over batching, limits and versioning.

Interview angle

The probe is almost always deletes: "A user deletes a document; 30 seconds later someone asks about it." Walk through the tombstone, the query-time deny-list, filtered vector deletion, cache invalidation and the source re-check, then state the guarantee in one sentence: once the delete is processed, the document can't be served, and physical removal follows within pipeline lag.

Go deeper

Case 2: Conversation memory and chat history at scale

Prompt: "Design storage for chat history and long-term memory for a consumer assistant with tens of millions of users." Memory techniques are in B3 and the assistant's product design is B6 case 6. This case is the storage and pipeline engineering underneath.

Clarifying questions and requirements

Functional: append messages (including streamed replies), list and load conversations, extract and edit memories, search past chats, delete at every granularity. Non-functional: read-your-writes within a conversation; p99 history load under 50 ms; durable once acknowledged; memory extraction within minutes of a session ending; account deletion everywhere within a published SLA (say 30 days, including backups ageing out); per-region storage.

Estimates

Assumptions: 50M monthly and 10M daily users; 12 user messages per daily user, each answered.

API and data model

POST /conversations/{cid}/messages {client_msg_id, content} → 200 + SSE stream of the reply GET /conversations/{cid}/messages?before=<msg_id>&limit=50 GET /users/me/conversations?cursor=… GET/PATCH/DELETE /users/me/memories/{mid} DELETE /conversations/{cid} DELETE /users/me (async; returns a deletion job id) messages PK = conversation_id SK = message_id (ULID: time-ordered, unique) role, content | content_uri, model, prompt_version, token_counts, status (streaming|complete|error) conversations_by_user PK = user_id SK = last_activity#conversation_id (title, summary_id) summaries PK = conversation_id SK = covers_until_message_id memories PK = user_id SK = memory_id fact, source_msg_ids[], confidence, valid_from/valid_to memory_vectors namespace = user_id (vector store; payload: memory_id)

client_msg_id makes sends idempotent, so a retried slow send doesn't create two (paid) replies. Partitioning by conversation_id keeps a conversation together and ordered, like Discord's channel-plus-time-bucket scheme. Discord 2023 One conversation never comes near DynamoDB's per-partition limit of 1,000 writes per second. AWS Very long agentic conversations can add a time bucket to the key.

Client Chat servicebuilds context,streams reply Redis: conv tail Messages storePK conv_id, SK ULID Memories + vectors LLM gateway Kafkakey = conversation_id Memory extractor Summarizer Chat search indexer Archiver → S3 Deletionorchestrator(Temporal) CDC / streams Hot path: synchronous write of the user message, then context build. Everything after the reply is asynchronous and derived from the event stream.
Chat history and memory. The message store is the system of record; memories, summaries, the search index and the archive are derived views rebuilt from the stream.

Data flow for one turn

  1. The service writes the user message (conditional put on client_msg_id) to the store and the Redis tail. It's durable before any model call.
  2. It builds context from the tail (falling back to a strongly consistent read), the latest summary, and top-k memories from the user's vector namespace.
  3. It streams the reply, checkpointing the partial assistant message every few seconds (status streaming) so a crash leaves a recoverable record, then finalizes with token counts.
  4. A change stream (DynamoDB Streams, Debezium or an outbox) emits events to Kafka keyed by conversation. Consumers derive views: memory extraction (debounced until the conversation is idle, say 10 minutes), summarization (past a token threshold), chat search and archival.

Deep dive 1: read-your-writes for the next turn

It fails if history is read from a lagging replica, or if the previous reply isn't persisted because the user sent a new message mid-stream. Write the user message synchronously, read history with strongly consistent reads or from the write-through tail, and serialize turns per conversation with a short lease. Derived views stay off this path deliberately: "I'm vegetarian" followed by a recipe request is answered from the tail, not memory. Memory only matters across sessions, where minutes of lag are acceptable; if not, extract at session close with high priority.

Deep dive 2: hot vs cold storage and summarization

Keep 30–90 days hot, then archive conversations to object storage grouped per user (cheap deletion and export) rather than per day (cheap analytics, expensive deletion). Opening an old conversation rehydrates it with a TTL. A DynamoDB TTL attribute can expire hot items automatically. DynamoDB TTL Summaries are separate records with covers_until_message_id, so context is "summary plus later messages", and summaries can be regenerated with a better model without touching messages.

Deep dive 3: GDPR deletion fan-out

Account deletion must reach every store derived from the user's messages: messages, caches, the archive, memories, the vector namespace, chat search, analytics, traces, eval datasets and backups. GDPR Art. 17 Run it as a durable workflow: record the request, tombstone the account to block reads and writes, fan out one idempotent, retried activity per store from a data inventory, and audit each. Two techniques make it tractable: lineage (every derived record carries user_id, so deletion is a partition drop or filtered delete) and crypto-shredding (per-user keys for archived data, destroyed on deletion, which makes copies in backups unreadable). Kafka data expires with retention or via tombstones in compacted topics.

Failure modes and trade-offs

Interview angle

"How do you prove you deleted everything?" Deletion is a workflow with a per-store checklist generated from a data inventory, each step is audited, and an automated test plants a canary user's data everywhere and asserts it's gone. Be honest about backups: they age out within the retention window or are crypto-shredded.

Go deeper

Case 3: Multi-tenant LLM gateway / AI platform API

Prompt: "Every team (or every customer) calls LLMs through us. Design the gateway: routing across providers, quotas, streaming, fallbacks, metering for billing, PII and audit."

Clarifying questions

Requirements and estimates

Functional: an OpenAI-compatible API with streaming, routing by alias and policy, per-tenant quotas and priorities, retries and fallbacks, caching, PII redaction, metering and audit. Non-functional: under ~20 ms added p99 before the first upstream byte; higher availability than any single provider; token-accurate metering with no double counting; no tenant can starve another.

Assumptions: 1,500 tenants; peak 2,000 requests/s; average 2,500 input tokens (half cacheable) and 350 output tokens; average request duration 7 s.

Tenantapps / SDKs gateway pods (stateless, async I/O) AuthN/Ztenant, key PII redactplaceholders AdmissionRPM+TPM bucket Fair queuepriority classes Cacheexact / semantic Routerpolicy, health, cost SSE proxycount, cancel circuit breakers · retry budget · per-provider bulkheads · request_id (idempotency) Provider A (region 1, 2) Provider B Self-hosted (vLLM) Kafka: usagekey = tenant_id Flink meteringdedup, windows Billing ledger (PG) ClickHouse usage Audit logS3, object lock Redisbuckets, leases
Multi-tenant LLM gateway. The hot path stays stateless apart from quota checks; everything about money (metering, billing, audit) is derived asynchronously from an idempotent usage stream.

Deep dive 1: distributed token buckets for token-denominated quotas

Each tenant has request and token buckets per model class; a Redis Lua script refills by elapsed time × rate, checks both, and debits atomically. Redis scripting Output tokens aren't known up front, so charge an estimate at admission (tokenized input plus the tenant's p90 output for the route, capped by max_tokens) and reconcile when the stream ends. Reserving the full max_tokens is safe but starves tenants who set it high. Huge tenants create hot keys: shard their bucket across N keys at 1/N rate each, or have pods lease token blocks and spend locally, trading slight over-admission for far fewer Redis calls. If Redis is down, fail open to conservative per-pod limits; Stripe applies the same "fail safely" principle. Stripe 2017 A second set of buckets mirrors each upstream account's limits, fed by rate-limit response headers, so the router avoids sending traffic that's certain to get a 429.

Deep dive 2: priority and fairness

When upstream capacity is short, requests shouldn't fail at random. Admit into priority classes (interactive, standard, batch) with weighted fair queuing across tenants within each class, so a burst queues behind its own tenant's share. Interactive requests get a short maximum queue time (say 2 s), then a 429 with retry-after; batch can wait or be redirected to a provider batch API. It's the criticality-based shedding from Google SRE applied to tokens. Google SRE

Deep dive 3: streaming proxy, retries and fallbacks

The SSE proxy forwards chunks unbuffered at every hop, sends heartbeat comments through load balancers, captures the final usage block, and cancels upstream when the client disconnects. SSE spec Retries are transparent only before the first byte: retry on another deployment of the same model, then fall back along a pre-approved chain of models that passed the tenant's evals. After tokens have been sent, emit an error event and record partial usage. Wrap each provider-model-region in a circuit breaker and cap retries with a budget. LiteLLM's proxy exposes fallbacks, cooldowns and retries as configuration. LiteLLM

Caching, PII, metering and audit

Provider prompt caches are usually scoped per provider, region or account, so route requests sharing a long prefix (same tenant and system prompt) to the same upstream: consistent hashing on a prefix hash, balanced against load. Exact-match caching suits repeated deterministic calls; semantic caching should be opt-in per route with conservative thresholds and sampled precision checks. A fast regex and NER pass replaces PII with placeholders and restores it in the response, the pattern Uber describes. Uber 2024 Every completed or aborted request emits a usage event (request ID, tenant, model, input, cached and output tokens, status) to Kafka; a Flink job deduplicates on request ID, aggregates per tenant and model per minute, and writes a ledger with a unique constraint per window, so replays can't double-bill. Raw events go to ClickHouse for dashboards and to S3 for reconciliation with provider invoices; audit records go to write-once object storage.

Common mistake

Retrying at every layer: the SDK retries 3 times, the gateway retries 3 times, the provider client retries 3 times. One provider hiccup becomes 27× load on the struggling provider. Retry at exactly one layer (the gateway), respect retry-after, and enforce a retry budget.

May be out of date

Provider rate-limit structures, prompt-cache pricing and whether cached tokens count toward limits differ between providers and change often. Uber's figures (about 30 teams, 16M queries a month, peak 25 QPS) are from mid-2024. Uber 2024

Go deeper

Case 4: Agent execution platform for long-running agents

Prompt: "Customers define agents that run for minutes to days, call tools, run code and sometimes wait for human approval. Design the platform that executes them reliably at scale." The agent loop is in B4 and the coding-agent product in B6 case 3; this case is the runtime.

Clarifying questions and requirements

Functional: start, observe, pause, resume and cancel runs; tools with declared side-effect levels; approval steps; sandboxed code; per-run budgets; event history with replay. Non-functional: no lost runs, no duplicated side effects, resume within seconds of a crash, waiting runs cost nothing, tenants isolated.

Estimates

Assumptions: 100k runs/day averaging 30 model steps and 20 active minutes; 25% pause for approval (median 8 hours); 40% use code execution for ~10 minutes.

Run APIstart/cancel Approval UI→ signal/update Temporal clusterevent history per rundurable timerstask queues Workflow workersagent loop (deterministic) llm-q workersI/O bound, high conc. tool-q workersidempotency keys sandbox-q workersbounded by pool LLM gateway External APIsemail, CRM, pay microVM poolwarm, no egress Redisbudgets S3context blobs Kafka: run-eventskey = run_id ClickHouse + S3steps, payloads Replay UI A parked run is just history in the Temporal database plus a timer: no worker, thread or sandbox is held while a human decides.
Agent runtime on a durable workflow engine. Separate task queues let LLM, tool and sandbox capacity scale independently.

Data flow

Starting a run creates a workflow with workflow_id = run_id; duplicate starts are rejected, so the start API is idempotent. The workflow is the agent loop: it calls an LLM activity with a reference to the context in S3 (history must stay small), checks budget and policy on the proposed tool calls, and dispatches tool activities, appending results to the context and step events to Kafka. For a tool marked "requires approval", it records a pending approval, notifies the approver, and waits on a signal with a durable timer (say 72 hours, then escalate or cancel). Temporal message passing Timers Every N steps it continues-as-new with compact state. Continue-As-New

May be out of date

Durable-execution engines limit history length, history size and payload size, and the values depend on engine, version and hosting. The stable principle: keep large prompts, documents and tool outputs in object storage and pass references. Check current limits before quoting numbers.

Deep dive 1: tool-call idempotency

Activities run at least once: a worker can finish a tool call and crash before reporting it, and the engine retries. Temporal activities Give each tool call a key derived from run_id and step_id (stable across retries) and pass it to APIs that accept idempotency keys. Stripe idempotent requests Otherwise use check-then-act (look up an external reference you set before creating) or an intent log with a unique constraint. Classify tools as read, reversible or irreversible, and require approval for irreversible ones in code, not in the prompt. Model decisions are non-deterministic but replay must be deterministic; the engine handles this by recording each LLM activity's result in history, so a replay reuses the recorded output instead of calling the model again.

Deep dive 2: scheduling, autoscaling and the sandbox pool

Use separate task queues: LLM activities are I/O-bound (hundreds of concurrent calls per worker), tools vary, and sandbox activities are capped by the microVM pool. Autoscale each pool on schedule-to-start latency rather than CPU, and apply per-tenant concurrency limits at dispatch so one tenant's 10,000-run batch queues instead of taking every sandbox. Run model-written code in microVMs (Firecracker) or a user-space kernel (gVisor), with no network by default and an allowlisting egress proxy. Firecracker gVisor Keep a warm pool sized by Little's law (request rate × provisioning time, plus headroom), and never reuse a sandbox across tenants. Temporal's integration with the OpenAI Agents SDK runs each agent invocation as an activity on this model. Temporal 2025

Deep dive 3: cost limits and event-sourced history

Budgets: the workflow tracks per-run tokens, dollars, steps and wall time, checking before each model call; tenant budgets use reserve-then-reconcile in Redis as in case 3. An exhausted budget pauses the run for a human rather than failing silently; a kill switch cancels a tenant's runs; step caps and loop detection (same tool, same arguments, three times) stop runaways. Event sourcing: each step appends an immutable event keyed by run_id (llm.completed with tokens and an output reference, tool.called, approval.granted), giving an audit trail, live progress and two kinds of replay: faithful (recorded outputs, to reproduce a bug) and counterfactual (fork from step k with a new prompt or model). Fowler: event sourcing Keep this log separate from the engine's history: one is for humans and analytics, the other for execution.

Failure modes and alternatives

A lighter design (state in Postgres, steps as queue messages, cron for timeouts) works at small scale but tends to reinvent retries, timers and versioning. Framework checkpointing such as LangGraph's persistence suffices for many products but keeps durability inside the application process. LangGraph durable execution

Interview angle

"The agent sent the refund, then the worker crashed before recording it. What happens?" The activity retries with the same idempotency key and the payment API returns the original result. Without API support, check by external reference first. The workflow records the decision once; idempotency keys make the effect happen once.

Go deeper

Case 5: LLM observability and tracing pipeline at scale

Prompt: "Build the backend for an LLM observability product (or an internal platform): SDKs send traces of LLM calls, tool calls and agent steps; users search traces, see dashboards, and run online evaluations." B5 covers what to trace and evaluate. This case is the data pipeline.

Clarifying questions and requirements

Non-functional: ingestion never blocks or slows the customer's application (SDKs batch and send asynchronously, and drop data rather than block); a trace is queryable within ~30 s; dashboard queries return in about a second over 30 days of data; tenant isolation; scrubbing applied before persistence.

Estimates

Assumptions: 2B spans/day across all tenants (≈ 23k/s average, ~100k/s peak); 30% are LLM spans carrying prompts and completions averaging 8 KB; the rest average 1 KB.

SDKsOTel, batched Ingest APIauth, quota, 202 S3 raw batches Kafka: spanskey = trace_id Stream processingPII scrubcost = tokens × pricetrace assembly (windows)tail samplingmetric rollups ClickHousespans, traces, scores S3 payload store Kafka: trace.completed+ metrics topic Online eval workerssample, budget, judge LLM gateway (low pri) Alerting Query API + UIdashboards, search scores Raw batches in S3 make the whole pipeline replayable (for example, after a scrubbing-rule fix or a pricing-table correction).
Observability pipeline. The ingest edge only authenticates, validates and enqueues; all enrichment happens in stream processors; evaluators are just another consumer.

Data model

spans ENGINE = ReplacingMergeTree(version) PARTITION BY toYYYYMMDD(start_time) ORDER BY (project_id, toStartOfHour(start_time), trace_id, span_id) TTL start_time + retention_for_plan project_id, trace_id, span_id, parent_span_id, session_id, name, kind start_time, end_time, status, model, provider input_tokens, cached_tokens, output_tokens, cost_usd input_preview, output_preview (first N KB) payload_uri (s3://… full body) attributes Map(String, String) version scores ORDER BY (project_id, trace_id, evaluator) value, label, judge_model, rubric_version, created_at

Leading the sort key with project_id keeps each tenant's data contiguous, so every query, which always filters by project, reads only that tenant's parts. ReplacingMergeTree handles a span that arrives twice (start and end, or a retried batch) by keeping the highest version per sort key after background merges. Merges are eventual, so queries that need exact deduplication pay for it at read time. ClickHouse MergeTree Attribute names follow the OpenTelemetry GenAI semantic conventions where possible, so standard SDKs work unchanged. OTel GenAI conventions

Deep dive 1: large payloads and PII

Prompts and completions dominate bytes and are rarely read in full, so store a truncated, searchable preview inline and the full body in S3 under project/trace/span (the claim-check pattern). Langfuse moved raw events to S3, keeping only references in Redis, after Redis memory and serialization cost became a bottleneck. Langfuse 2024 Extract base64 attachments at ingest so they never bloat Kafka. PII scrubbing runs in the stream processor before any sink, with per-project rules; raw S3 batches are unscrubbed, so give them short retention and strict access, or scrub at the edge for customers who require it.

Deep dive 2: trace assembly and sampling

A trace's spans arrive out of order from many processes. Keying by trace_id puts them in one partition, so a keyed operator can buffer them and emit trace.completed once the root span has ended and a quiet gap (say 30 s, watermark-driven) has passed; late spans still update aggregates. Flink time Head sampling (decided at the root) is cheap but drops failures as often as successes; tail sampling keeps every error, slow trace and thumbs-down plus a share of the rest, at the cost of buffering. OTel sampling A common compromise for LLM products: keep all span metadata for cost accounting, and sample payload retention.

Deep dive 3: online evals as a stream consumer

Evaluator workers consume trace.completed, match project evaluator configs (filters, sample rates), and enqueue judge jobs, so a slow judge never back-pressures ingestion. Jobs carry per-project budgets, run at the lowest gateway priority or through a batch API, and are idempotent on (trace, evaluator, rubric version). Scores land in ClickHouse; windowed aggregates (judge score per prompt version per hour, error rate, p95 latency, cost per trace) feed alerts. Heuristic evaluators (JSON validity, regex) are nearly free and should run on all traffic.

Failure modes and trade-offs

Interview angle

"Why not just store traces in Postgres?" Answer with access patterns: append-only, huge volume, mostly aggregate queries over time ranges, rare point lookups, and retention by time. That's a columnar OLAP workload. Then mention the exceptions you'd keep in Postgres: projects, API keys, evaluator configs, prompt versions (transactional, small). Langfuse made exactly this split. Langfuse 2024

Go deeper

Case 6: Document processing pipeline

Prompt: "We receive millions of PDFs a day (invoices, receipts, forms). Extract structured data into our system of record, with humans reviewing what the model isn't sure about." B6 discusses the review-routing decision (question 20 there). This case is the pipeline.

Clarifying questions and requirements

Functional: ingest from email, upload and API; classify, split and extract; validate; route to human review by confidence; post to the system of record; reprocess on demand. Non-functional: every document reaches a terminal state (posted, rejected or in review), never lost; each document posted exactly once; p95 under 10 minutes for the auto path; cost per document tracked.

Estimates

Assumptions: 5M documents/day, 3 pages each, so 15M pages/day ≈ 175 pages/s average and maybe 900/s at a month-end peak. A model cascade: 50% handled by a cheap tier (layout model plus rules for known suppliers' templates), 40% by a small vision-language model, 10% escalated to a frontier model; 4% of documents end up in human review.

Intakeemail/upload/API S3key = sha256 Doc workflowone per document Classify + split Page queuefan-out model cascade T0: OCR + template rules T1: small VLM extraction T2: frontier VLM (hard) Fan-in + validateschema, sums, PO match Confidence routerthresholds by field Review queuepriority, SLA Results DB+ outbox (1 txn) ERP posteridempotency key DLQ + replay Corrections→ eval + training Escalate only on low confidence or failed validation; the expensive tier sees about a tenth of documents.
Document processing with fan-out per page, a model cascade, confidence routing and an exactly-once posting path.

Data flow

  1. Intake stores the file in S3 under its SHA-256 (identical resubmissions dedupe for free) and starts a per-document workflow keyed by tenant and hash.
  2. The workflow classifies, splits pages and fans out page tasks (queue or child workflows), each running OCR or a layout model, then the cheapest extraction tier likely to work.
  3. Fan-in waits for all pages. With plain queues, track completion as a set of page IDs, not a counter, so a duplicate result can't trigger aggregation early.
  4. Validation: schema checks, arithmetic (line items sum to subtotal; subtotal plus tax equals total), and business cross-checks (supplier exists, PO open, amounts match). Failures escalate a tier, with the errors as feedback, or go to review.
  5. The confidence router auto-approves or creates a review task; corrections become labelled data. Approved results post via an outbox with idempotency key doc_id.

Uber's TextSense describes a similar staged pipeline (ingestion, pre-processing, OCR, LLM extraction, post-processing, validation with human review, monitoring) and reported 2× less manual effort and 70% less handling time. Uber 2025

Deep dive 1: exactly-once results into the system of record

The pipeline is at-least-once everywhere: queue redelivery, workflow retries, a reviewer double-clicking "approve". Exactly-once is built at the edges. The results table has a unique constraint on (doc_id, extraction_version), so a duplicate write is a no-op. Posting goes through an outbox written in the same transaction as the approved result, so "approved but never posted" and "posted but not recorded" can't happen. The poster passes doc_id as an idempotency key or external reference. If the ERP doesn't support idempotency, the poster first queries by external reference, and a nightly reconciliation job compares ERP records with the results table. Separately, add business-level duplicate detection (same supplier, invoice number and amount arriving as different files), because duplicate invoice payment is a classic fraud and error pattern that file hashing can't catch.

Deep dive 2: the model cascade and cost

Route each document to the cheapest tier that's likely to succeed, and escalate on evidence: low field-level confidence, validation failures, or disagreement between two cheap extractions. Known suppliers with stable layouts often need no LLM at all. Non-urgent volume (back-office batches) goes through provider batch APIs at about half price. OpenAI Batch Because human review dominates cost, tune thresholds against a labelled set for a target error rate among auto-approved documents, and measure per-field: totals and bank details get strict thresholds, memo fields loose ones.

Failure modes

Interview angle

Interviewers like "how do you know the auto-approved ones are right?" Say: sample a small percentage of auto-approved documents for blind human review continuously, report error rate per field and per supplier, and treat that number as the SLO that sets thresholds. Without the audit sample, straight-through rate is a vanity metric.

Go deeper

Case 7: Semantic search and recommendations for a commerce or content platform

Prompt: "Add semantic search and personalized recommendations to a marketplace with 100M listings. Product wants LLM-quality relevance and 'why this result' explanations. Search must stay fast."

Clarifying questions and requirements

Functional: ranked, filterable search; home-feed and similar-item recommendations; optional explanations. Non-functional: p99 under 300 ms; new items searchable in minutes; out-of-stock hidden in seconds; ranking changes gated by A/B tests.

Estimates

OFFLINE / NEARLINE Warehouse/lakeevents, catalog Spark: train + embedtwo-tower, LLM labels Kafka: catalog+ user events Flink: embed changeditems, session features ANN indexreplicated shards Feature storeonline (Redis/KV) Lexical indexBM25 Explanation cacheprecomputed (head) ONLINE (p99 300 ms) Query Understandingcached rewrite, 20 ms RetrieveANN ∥ BM25, 40 ms Filterstock, policy RankGBDT/DNN, 40 ms Re-rank top 50cross-encoder 40 ms Results render; explanation streams after(cache hit, or async small-LLM call)
Semantic search and recommendations. LLMs do most of their work offline (labels, enrichment, rewrites, precomputed explanations); the online path is classic multi-stage retrieval and ranking.

Data flow

Offline, a two-tower model is trained on engagement data and Spark batch-embeds the catalog and users nightly into a new index version. This candidate-generation-then-ranking split is the structure YouTube described. Covington+ 2016 Nearline, catalog changes flow through Kafka: price and stock updates change filterable attributes in place (out-of-stock goes to a fast-path filter), content changes re-embed the item, and session events update short-term user features. Pinterest's PinnerSage represents users with several embeddings to capture multiple interests. Pal+ 2020 Online, the query is rewritten (cached for head queries), ANN and lexical retrieval run in parallel, filters drop unavailable items, a ranker scores a few hundred candidates with online features, and a small cross-encoder re-ranks the top 50.

Deep dive 1: the p99 latency budget, and where the LLM goes

StageBudget (p99)Notes
Network, auth, request fan-out20 ms
Query understanding20 msCache of LLM rewrites for head queries; small classifier for the tail
Retrieval (ANN ∥ lexical)40 msParallel; take what returns by the deadline
Feature fetch30 msBatched multi-get from the online store
Ranking (~500 candidates)40 msGBDT or a compact DNN
Cross-encoder re-rank (top 50)40 msSmall model on GPU; skip under load
Assembly and business rules20 ms
Headroom90 msTail latency of fan-out calls eats this

A frontier LLM can't sit in this path: its TTFT alone can exceed the whole budget, and the token volume is enormous. Put LLMs where latency doesn't count: offline (attribute enrichment of listings, synthetic queries and relevance labels for training and evaluation, rewrite tables for the top queries), cached (explanations precomputed for frequent query-item pairs), and after first render (an explanation streamed into the UI a few hundred milliseconds after results appear, from a small model). Stages also need deadlines: if re-ranking would push past budget, return the ranker's order. That's graceful degradation, with a metric on how often it happens.

Deep dive 2: feature store and training/serving skew

Ranking features exist twice: in the warehouse for training and in a low-latency store for serving. A feature store manages both from one definition and does point-in-time correct joins, so training never sees feature values from after the label event. Feast The classic bug is training/serving skew, where a feature is computed one way in Spark and another online. Compute streaming features once and write them to both stores, and log the features actually served so training can use them.

Deep dive 3: experimentation

Randomize by user, not request. Report a primary metric (purchases, long dwell) with guardrails (latency, zero-result rate). For ranking changes, interleaving (mixing two rankers' results and crediting clicks) needs far less traffic than A/B tests. LLM-judged relevance labels give fast offline NDCG, but calibrate them against human judgments and online outcomes before letting them gate launches.

Failure modes

Interview angle

"How do you put an LLM in a p99 under 300 ms path?" Mostly you don't: distil it into a small cross-encoder, cache it for head queries, precompute enrichment offline, or defer it until after first render. If pressed, budget a tiny model on a handful of candidates with a hard deadline and a fallback, and show the token arithmetic.

Go deeper

Case 8: Proactive notifications and an AI assistant scheduler

Prompt: "Our assistant should proactively tell users things: 'your flight is delayed, leave 30 minutes later', 'this email from your landlord needs a reply today'. Events come from email, calendar, travel and shopping integrations. Design it without spamming people." B8 covers how always-on assistants do proactivity (heartbeats, cron) at personal scale. This is the multi-tenant version.

Clarifying questions and requirements

Functional: ingest events; decide notify now, schedule, digest or drop; generate the message; deliver; let users snooze, mute and give feedback. Non-functional: urgent notifications within a minute of the triggering event; never duplicate; a hard per-user frequency cap; scheduled notifications fire within a minute of their due time and are re-validated before firing.

Estimates

Assumptions: 20M connected users; 500M source events/day (≈ 5.8k/s); deterministic rules discard 95%, leaving 25M candidates; the LLM triages candidates in per-user micro-batches, about 8M calls/day (~90/s) at ~1.5k input and 100 output tokens, so ~12B input tokens/day on a small model. A cap of 3 notifications per user per day bounds delivery at 60M, but expect ~15M actually sent, with many scheduled for a local time, which concentrates timers at the top of each hour in each time zone.

Email/Calendar Travel/Shopping Kafka: eventskey = user_id Normalizededupe event_id Rules engineprefs, quiet hours LLM decision (per-user batch)now / schedule / digest / drop + text Policy gateper-user cap, semantic dedup Timer servicebucketed sorted sets / workflows Fire-time re-validatestill relevant? cap still ok? Delivery (APNs/FCM/email)idempotent on notification_id Feedback eventsopen, dismiss, mute → models urgent path: rules only
Proactive notifications. The LLM sits in the middle of a deterministic funnel: cheap rules before it, hard policy after it.

Data flow

Connectors publish to Kafka keyed by user_id, so one user's events are processed in order. A normalizer maps them to a common schema and deduplicates by source event ID (Redis SET NX with TTL). Rules apply opt-outs, quiet hours and cheap relevance filters. Survivors accumulate per user for a few minutes (rule-flagged urgent events skip the wait); then one LLM call sees the user's pending candidates, a compact profile and recent notification history, and returns structured decisions. A policy gate enforces a per-user token bucket (say 3 a day, an hour apart) and drops near-duplicates by topic key or embedding similarity; scheduled items go to the timer service. LinkedIn's Concourse has the same funnel (Kafka, Samza, ML scoring), with its Air Traffic Controller as the final throttling and dedup gate. LinkedIn 2018 LinkedIn reported that centralizing communications under ATC cut member complaints by over 65% and email volume by half. LinkedIn 2015

Deep dive 1: delayed jobs and timers at scale

Kafka has no delayed delivery, and SQS delays are capped at minutes, so millions of "fire at 08:00 local" timers need their own service. Options:

Two details matter whichever you choose. Add jitter to "08:00" (spread over ±10 minutes) or every time zone's top of the hour becomes a thundering herd at the push provider. And re-validate at fire time: the meeting may have been cancelled, the user may have opened the app and seen the information, or the daily cap may have been reached by an urgent notification. Timer firing is at-least-once, so delivery dedupes on notification_id.

Deep dive 2: the decision step and notification fatigue

The LLM is good at judging relevance and writing the message, and bad at enforcing invariants, so caps, quiet hours, opt-outs and dedup live in code after it. Per-user batching cuts cost and improves decisions: the model can pick the one item that matters and fold the rest into a digest. Feedback (opened, dismissed, muted) tunes per-user budgets and labels data for a cheaper classifier that eventually replaces the LLM for common event types.

Common mistake

Putting "don't send more than 3 notifications a day" in the prompt. The model can't see what other workers sent in parallel, and prompts aren't guarantees. Caps are a per-user counter or token bucket checked atomically right before delivery.

Go deeper

Case 9: Real-time moderation and safety pipeline for user-generated content

Prompt: "Moderate posts, comments and images on a social platform in near real time, using LLMs where they help, without blowing the budget or the latency." B5 covers guardrails for an LLM's own inputs and outputs. This case moderates users' content.

Clarifying questions and requirements

Non-functional: pre-publish check p99 under ~150 ms; high-severity content actioned within minutes; precision high enough that false removals don't drive users away; every decision explainable and auditable; policies updatable within a day.

Estimates

Assumptions: 500M posts and comments/day (≈ 5.8k/s, ~20k/s at peak), 20% with images. A cascade where layer 1 decides confidently on about 97% of items, an LLM sees about 2.5%, and humans see about 0.2%:

synchronous, p99 ~150 ms Post L0: hashes,blocklists, velocity L1: classifierdistilled, GPU, ~5 ms Publish (97%) Block / hold Kafka: borderline~2.5% L2: LLM policy modelpolicy text in prompt Human reviewseverity × reach Label storedecisions, appeals Retrain / distill L1active learning new model version
A moderation cascade. Each layer is an order of magnitude more expensive and sees an order of magnitude less traffic. Human decisions train the cheap layer.

Deep dive 1: latency path vs throughput path

Only L0 and L1 run synchronously; they fit a 150 ms budget at 20k/s, with the classifier batching requests on the GPU over a few-millisecond window. Borderline items publish with reduced distribution (or are held, on high-risk surfaces such as new accounts) and enter Kafka for L2, which runs on a fixed GPU pool behind a queue. A spike grows the queue, not posting latency, and items are dequeued by severity × current view velocity, not FIFO, so a viral borderline post jumps the line. Open-weight policy-reasoning models such as gpt-oss-safeguard read the policy text at inference time, so a policy change ships as a prompt change; Llama Guard is an option for a fixed taxonomy. OpenAI: gpt-oss-safeguard Llama Guard 4

Deep dive 2: the feedback loop

Human decisions, appeal outcomes and L2 verdicts become labelled examples tagged with the policy version. Active learning sends the items L1 was least sure about to humans, which improves the model fastest per label. A new policy category starts as an L2 prompt (fast to ship, expensive per item), gathers labels, and is distilled into L1 over weeks, moving traffic back to the cheap path. Track per-category precision and recall on a golden set, L1 score drift and appeal overturn rates. Roblox describes a similar loop with golden sets, uncertainty sampling, AI-assisted red teaming and synthetic data. Roblox 2025

Failure modes and trade-offs

Interview angle

Expect "why not send everything to the LLM?" Do the arithmetic: 500M items × 800 tokens is 400B tokens a day, and the latency won't fit a pre-publish budget. Then say what the LLM is good for: borderline judgement, explanations for reviewers and appeals, labelling data for the cheap model, and quick response to new policies.

Go deeper

Case 10 (short): AI support platform integrated with an existing ticketing system

Prompt: "Our AI support agent must work inside the customer's existing ticketing system (a Zendesk-like SaaS): read new tickets, draft or send replies, update fields, and hand off to humans." The agent itself is B6 case 1. The interesting part here is keeping two systems in sync.

TicketingSaaS Webhook receiverHMAC, dedupe, 200 fast Queue AI workeragent, per-ticket lock Local DBticket mirror + outbox Sync workerrate limit, If-Match Reconciliation pollerincremental by updated_at
Two-way sync with an external system: webhooks in (fast ack, dedupe, enqueue), outbox out (idempotent, version-checked, rate-limited), and a reconciliation loop for whatever webhooks miss.
Interview angle

The probe is "what if the webhook arrives twice, or never?" Twice: event-ID dedupe plus idempotent processing. Never: reconciliation polling. Out of order: re-fetch current state and compare versions. The design should be correct with webhooks switched off, just slower.

Part 3. Cross-cutting deep dives

Handling LLM provider outages and rate limits system-wide

Treat model capacity as a shared, scarce resource that a central gateway allocates, not something each service grabs on its own.

Cost controls

Start with attribution: the gateway tags every call with tenant, feature, prompt version and model. Then add budgets per tenant, feature and run (reserve, then reconcile), alerts on spend rate, and kill switches. The levers, roughly in order of effort: prompt caching (stable content first, routing for cache affinity); output caps and context trimming; routing and cascades (cases 6 and 9); batch APIs for anything that can wait, at about half price Anthropic Batches; and finally distillation or self-hosting for steady, high-volume, narrow tasks, after doing the break-even arithmetic (tokens per day at API prices vs GPU-hours at realistic utilization). Notion's vector-search write-up shows such savings compounding across storage, engine and embedding-generation changes. Notion 2026

Evaluating quality in production with streaming evals

Offline suites gate changes; production evals catch what they miss: new query types, provider-side model changes, drift. The case 5 pipeline is the mechanism. The design choices are run cheap deterministic checks on 100% of traffic, LLM judges on a sample, and human review on a smaller sample that also calibrates the judges; stratify sampling toward important or rare segments; pin judge model and rubric versions; join real outcomes (thumbs, edits, escalations, refunds) by request ID; and alert on change against each version's own baseline, with enough samples to be confident. Methodology is in B5.

Data privacy: encryption, tenant isolation and residency

May be out of date

Provider data-retention options, regional endpoints and rate-limit structures change often and vary by contract tier. Check current documentation and your agreements.

Capacity planning for GPU and API capacity

Forecast tokens per minute at peak by model, with cached and uncached input separated, since they may count differently against limits. Anthropic rate limits Use Little's law for concurrency (streams, sandboxes, GPU slots) and queues to smooth bursts. Schedule batch into the valleys of interactive traffic to raise utilization. GPU autoscaling is slow (provisioning plus loading weights can take minutes), so keep warm headroom for interactive serving and use spot capacity for idempotent, checkpointed batch work. API capacity is partly a procurement problem with lead time: tier upgrades, committed throughput, multiple accounts and regions. Serving internals are in A6; cluster-level ML infrastructure is C3.

Schema evolution for event payloads

Events outlive the code that wrote them, especially in replayable logs and event-sourced agent histories. Use a schema registry (Avro, Protobuf or JSON Schema) with an enforced compatibility mode: backward (upgrade consumers first), forward (upgrade producers first) or full, optionally transitive. Confluent: schema evolution Shopify's CDC platform used Confluent Schema Registry with Avro. Shopify 2021 Add optional fields with defaults, never rename or reuse fields, version the event type when its meaning changes, and upcast old events in one place. AI-specific: every event with model output or vectors carries model_version, prompt_version and, for embeddings, model and dimension. That's what makes replays, migrations and audits possible later. With raw CDC, a column rename is a breaking change for every consumer, which is an argument for outbox contracts.

Interview angle

"What breaks at 10× scale?" Name the first bottleneck in this design (usually provider quota, then the human review queue, then index memory), the metric that would reveal it, and the specific change you'd make.

Interview question bank

1. How do you guarantee the vector index doesn't serve deleted documents?

Separate "can't be served" from "physically removed". On delete, write a versioned tombstone to the document registry (so a late update can't resurrect the doc) and add the ID to a deny-list whose TTL exceeds worst-case pipeline lag. The query service drops deny-listed hits, then vectors are deleted by doc_id filter. Give query caches short TTLs or invalidate them, and for sensitive corpora re-check top results against the source. Verify with a canary document that's deleted and then queried.

2. Exactly-once embedding: is it necessary?

No. You need an exactly-once result, which comes from at-least-once delivery plus idempotent writes: deterministic vector IDs (doc_id:chunk_hash), monotonic version checks in a registry, and a content-hash embedding cache so duplicates rarely even call the API. Kafka transactions wouldn't cover the embedding API or the vector database anyway. A duplicate costs a few redundant embedding calls.

3. How do you re-embed 500M chunks when switching models, with zero downtime?

Blue/green. Create a new index tagged with the new model; dual-write the live stream; backfill from a consistent snapshot in a low-priority lane via a batch API or burst GPU pool. At 500 tokens per chunk that's 250B tokens: weeks on a typical per-minute quota, hours on a few dozen GPUs. Evaluate on a golden set and shadow traffic, flip an alias atomically while switching the query-embedding model, and keep the old index for rollback.

4. How do you put an LLM in a p99 under 300 ms path?

Mostly you don't call a large one synchronously. Precompute (enrichment, rewrite tables for head queries), cache (explanations for frequent pairs), distil (a small cross-encoder trained on LLM labels), or defer (render, then stream the generated part). If an online call is unavoidable, use a small model on short input with a hard deadline and a non-LLM fallback. Show the arithmetic: thousands of QPS × hundreds of tokens is millions of tokens per second.

5. A consumer calls an LLM that takes 20–40 s per record. What goes wrong with a naive Kafka consumer?

Sequential processing exceeds max.poll.interval.ms, the consumer is evicted, a rebalance follows, and uncommitted records are reprocessed elsewhere, possibly hitting the same timeout. Throughput is also one call per partition at a time. Fix: poll small batches, process concurrently in an internal pool, commit the highest contiguous completed offset, and pause partitions when the pool is full; or use a queue model with per-record acks. Keep processing idempotent either way.

6. CDC or transactional outbox to drive an AI pipeline?

CDC to capture every change to existing tables without application changes, for example indexing a product database for RAG. Outbox for deliberate domain events with a controlled schema ("TurnCompleted"). Both avoid dual writes by deriving events from the commit. Raw CDC couples consumers to table schemas; the outbox costs an extra row per transaction. Many systems use both.

7. Design per-tenant token quotas when output length isn't known in advance.

Token buckets for requests and tokens in Redis, updated atomically by a Lua script. Charge an estimate at admission (tokenized input plus, say, the tenant's p90 output for the route, capped by max_tokens) and reconcile with actual usage at the end. For hot tenants, shard buckets or lease token blocks to pods. Define fail-open behaviour for when Redis is down, and mirror the provider's own limits so you don't send traffic that's certain to be rejected.

8. When can an LLM gateway transparently retry a streaming request?

Only before the first token reaches the client: retry on another deployment of the same model, then a pre-approved fallback, within a retry budget and honouring retry-after. After streaming starts, a silent retry would produce a different answer. Emit an error event, bill partial usage, and let the client regenerate. Cancel upstream on client disconnect.

9. How would you meter LLM usage accurately enough to bill on?

One usage event per request (aborted streams included) with a unique request ID and provider-reported token counts, published to Kafka. A stream job deduplicates on request ID, aggregates per tenant, model and window, and writes a ledger with a unique constraint per window, so replays can't double-count. Keep raw events for audit, reconcile daily with provider invoices, and post late events as corrections.

10. An agent run waits three days for a human approval. What does it consume?

Essentially only storage. In a durable workflow engine the run is an event history plus a durable timer; no worker, process or sandbox is held. The approval arrives as a signal that wakes the workflow, which replays history and continues; the timer escalates or cancels if nobody answers. Large context stays in object storage and is passed by reference.

11. How do you make agent tool calls idempotent when activities are retried?

Derive a key from stable IDs (run plus step, not attempt number) and pass it to APIs that support idempotency keys. Otherwise check by an external reference you control, or keep an intent log with a unique constraint. Require approval for irreversible tools. The engine records the decision once; keys make the effect happen once.

12. How do you run online LLM-as-judge evals without slowing ingestion or blowing the budget?

Make evaluators a separate consumer of a "trace completed" stream feeding a work queue, so judge latency never back-pressures ingestion. Sample (stratified), enforce per-project budgets, run at lowest priority or via batch, and make jobs idempotent on (trace, evaluator, rubric version). Run deterministic checks on everything, judges on a sample, and calibrate judges against humans.

13. How do you guarantee each extracted invoice is posted to the ERP exactly once?

Accept at-least-once internally; make the edges idempotent. Results have a unique constraint on (document, extraction version); the approved result and an outbox row commit together; the poster uses the document ID as idempotency key or external reference, checking before creating if the ERP lacks idempotency; a nightly reconciliation catches drift. Also detect business duplicates (same supplier, number and amount in different files).

14. Read-your-writes in a chat product: where can it break?

When history is read from a lagging replica, when the previous assistant message isn't persisted because the user sent a new message mid-stream, or when the next turn relies on an async derived view such as memory. Write the user message synchronously, read history with consistent reads or a write-through tail cache, serialize turns per conversation, and build context from raw recent messages.

15. A user deletes their account. Walk through the deletion.

A durable workflow records the request and writes a tombstone that blocks reads and writes, then fans out idempotent, retried activities driven by a data inventory: messages, caches, the per-user archive prefix, memories, the vector namespace, search, analytics, traces and eval sets. Kafka data expires or is tombstoned; backups age out or are crypto-shredded via per-user keys. Every step is audited, and a canary test verifies deletion end to end.

16. How do you schedule 15M "send at 8 am local" notifications a day?

A timer service, since Kafka has no delayed delivery and SQS delays are capped at minutes: time buckets sharded by user hash in Redis sorted sets (mirrored to a database) or in a time-partitioned table, with pollers that claim due items atomically. Add jitter to avoid top-of-hour herds, re-validate at fire time, and make delivery idempotent on notification ID because firing is at-least-once.

17. Why does moderation use a cascade instead of an LLM for everything?

Arithmetic and latency. Hundreds of millions of items × hundreds of tokens is hundreds of billions of tokens a day, and pre-publish checks have a budget of about 100 ms. Hashes and a distilled classifier decide most items synchronously; an LLM policy model handles the borderline few percent asynchronously; humans take the uncertain and severe, prioritized by reach. Their labels retrain the cheap layer.

18. Your LLM provider is down for two hours. What happens?

The gateway's circuit breakers trip and traffic moves to eval-validated substitutes. Features without a substitute enter their predefined degradation mode. Background work is shed to save capacity for interactive users, and durable queues and workflows hold async work and resume automatically. Afterwards, drop backlog items that are no longer useful instead of processing them blindly. All of this works only if the fallbacks were exercised before the outage.

19. How do you evolve event schemas without breaking consumers or replays?

Use a schema registry with enforced backward or full compatibility. Add only optional fields with defaults, never rename or reuse fields, and version the event type when its meaning changes. Consumers upcast old versions in one place. Put provenance (model, prompt version, embedding model and dimension) on AI events, and prefer outbox contracts over raw CDC for widely consumed streams.

20. Embedding-pipeline consumer lag keeps growing. How do you diagnose it?

Check lag in time, and whether it's on all partitions (a throughput problem) or a few (hot keys or a poison record). Then find where time goes: 429s mean you're quota-bound (adding consumers just adds 429s), slow upserts mean index churn, and frequent rebalances mean slow processing. Depending on the cause: scale up to the partition count, process concurrently, debounce hot docs, pause backfill, or move poison records to a DLQ.