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 step | Classic version | What changes when AI components are involved |
|---|---|---|
| Functional requirements | Features and user flows. | Add what a correct output looks like and who judges it. Agree whether results can stream. |
| Non-functional requirements | Latency, 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. |
| Estimates | QPS, 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. |
| API | Request-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 model | Entities 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 design | Services, 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 bottlenecks | Sharding, 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.
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
- Slow: latency ≈ TTFT + output tokens ÷ decode speed. Plan for seconds and a long tail.
- Rate-limited in several dimensions: requests, input tokens and output tokens per minute, per model. Anthropic documents a continuously refilling token bucket and returns
429withretry-after. Anthropic rate limits Your own limiters should count tokens too. - Non-deterministic: the same input can give different outputs. Persist outputs instead of recomputing; key caches on (model, prompt version, input hash); decide explicitly whether a replay reuses stored outputs or regenerates.
- Wrong without erroring: a 200 can be a bad answer, so validate (schemas, business rules) and monitor quality, not just status codes.
- Linear in tokens: cost = input × input price + output × output price, minus cache discounts. Cache reads bill at a fraction of input price, and on Anthropic's API they don't count toward input-token limits for most models, which changes capacity math as well as cost. Prompt caching
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").
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.
| Operation | Rough figure | Design implication |
|---|---|---|
| Hosted LLM time-to-first-token, short prompt, non-reasoning model | ~0.2–1 s | Stream it. Reasoning models can take far longer. |
| TTFT with a very long prompt (~100k tokens), uncached | several seconds | Prompt 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 s | Batch; quota, not latency, bounds throughput. |
| Self-hosted small embedding model, one modern GPU | ~1k–10k chunks/s | Wins for big re-embeds. |
| Cross-encoder rerank of ~50 candidates on GPU | ~10–50 ms | Fits a 300 ms budget. |
| Small transformer classifier per item (batched, GPU) | ~1–10 ms | First stage of cascades. |
| In-memory HNSW query, millions of vectors, recall ~0.95 | ~1–10 ms | Memory is the cost: 1B × 768-d float32 ≈ 3 TB. ann-benchmarks |
| Object-storage-backed vector search (warm cache) | tens of ms; cold queries slower | Much cheaper per vector; Notion reported 50–70 ms. Notion 2026 |
| Redis GET/SET in the same AZ | sub-ms; ~100k+ ops/s per node | Fine 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/s | Shard before it saturates. |
| Kafka produce→consume, same cluster | a few ms to tens of ms; tens to hundreds of MB/s per broker | Consumers, not Kafka, bottleneck. |
| S3 request rate per prefix | at least 3,500 writes/s and 5,500 reads/s; small-object latency ~100–200 ms | Keep S3 off tight sync paths. AWS S3 docs |
| DynamoDB per-partition throughput | 3,000 read units/s and 1,000 write units/s | Key design matters. AWS DynamoDB docs |
| ClickHouse batch insert / aggregate scan | insert in batches of thousands+ rows; scans of hundreds of millions of rows/s per server | Never insert row by row. |
| Cross-region round trip | ~60–150 ms | Residency and failover add to TTFT. |
| Human review of one item | ~30 s–3 min | Often the largest cost line. |
- Designing Data-Intensive Applications (Kleppmann; the 2nd edition with Riccomini came out in 2026): the reference for logs, replication, partitioning and consistency.
- Anthropic rate limits: token-bucket limits, cache-aware input-token limits, headers to drive client-side throttling.
- Google SRE book, "Handling Overload": client-side throttling and criticality, which map directly onto LLM quota management.
- B6: AI System Design Case Studies: the AI-centric answering framework this page builds on.
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.
- Consumer groups: each partition goes to exactly one consumer in a group, so parallelism is capped at the partition count. Independent groups read the same topic, which is how one change stream feeds an embedder, a search indexer and analytics.
- Offsets: commit after processing and a crash means reprocessing (at-least-once); commit before and records can be lost. Almost every design picks at-least-once plus idempotent processing.
- Retention and compaction: topics keep data by time or size, so slow consumers catch up and new ones replay. A compacted topic keeps the latest record per key, and a null-value tombstone eventually removes the key, which makes a "current state of every document" topic that an index can be rebuilt from.
- Idempotent producers and transactions: producer IDs and sequence numbers stop broker-side retries from duplicating records; transactions atomically write outputs and commit input offsets, and
read_committedconsumers skip aborted writes. Confluent: exactly-once Confluent: transactions That's exactly-once within Kafka; it doesn't cover the embedding API call or the vector-database upsert. - Rebalancing: membership changes reassign partitions. Cooperative rebalancing and the newer server-driven protocol (KIP-848) reduce disruption. KIP-848 The AI-specific trap is slow processing: a consumer stuck on long LLM calls can exceed
max.poll.interval.ms, get evicted, and cause a rebalance plus duplicate work. - Consumer lag, ideally expressed as the age of the oldest unprocessed record, is the key health metric and the natural autoscaling signal for AI worker pools.
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:
- CDC: Debezium reads the replication log (for Postgres, logical decoding through a replication slot) and emits a change event per committed insert, update or delete, with before and after images. A delete is followed by a tombstone so compacted topics can drop the key. Debezium Postgres Incremental snapshots backfill existing rows while streaming continues, using the watermarking technique from Netflix's DBLog. Debezium DBLog
- Transactional outbox: write the business row and an outbox row describing the event in one transaction; a relay (often Debezium's outbox router) publishes the outbox. microservices.io Debezium outbox router Raw CDC makes your table schema a public contract; the outbox publishes deliberate domain events instead.
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
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) | |
|---|---|---|
| Consumption | Competing consumers; a message is deleted when acknowledged. | Offsets per group; records kept for the retention period. |
| Ordering | None (SQS standard) or per message group (FIFO). | Per partition, by key. |
| Per-message retry | Native: visibility timeout, redelivery, DLQ. SQS DLQ | DIY: a failed record blocks its partition unless moved to retry topics or a DLQ. Uber 2018 |
| Replay / parallelism | No replay; scale workers freely. | Replay by rewinding; parallelism capped by partitions. |
| Best AI fit | Independent 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
- Postgres + pgvector: vectors live beside their rows, filters and ACLs, in one transaction, which removes a class of sync bugs. It supports HNSW and IVFFlat. pgvector A strong default up to tens of millions of vectors.
- Choosing a vector DB: scale (vectors × dims × bytes), filtered-search recall, update/delete throughput under churn, multi-tenancy (namespaces vs a filtered shared index), hybrid search, and in-memory vs disk vs object-storage-backed architecture. Details in B2.
- OLAP (ClickHouse): columnar storage sorted by an
ORDER BYkey, with partitions and TTLs, built for big batch inserts and fast aggregates, which fits traces, usage and eval scores. ClickHouse MergeTree - Wide-column / key-value (Cassandra, ScyllaDB, DynamoDB): partition by entity, cluster by time, for append-heavy data like chat messages. Discord 2023
- Object storage: cheap, durable, high-throughput but not low-latency; used via pointers.
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
- Idempotency keys: the server stores the key with the result and returns the stored result on retry. Stripe Leach For LLM side effects this matters twice: repeats are expensive, and they return a different answer.
- Retries: exponential backoff with jitter, a cap, a retry budget, at one layer only, honouring
retry-after. AWS: retries and jitter - DLQs: after N attempts, park the message with its error history, alert, and provide replay. On Kafka, use tiered retry topics with growing delays. Uber 2018
- Backpressure and load shedding: bounded queues, pausing consumption, and shedding low-criticality work first. Stripe runs four limiters, from per-user request rate to a worker-utilization load shedder. Stripe In a backlog, old work may be worthless, so use message TTLs or temporarily serve newest-first. AWS: backlogs
- Rate limiting: token buckets (capacity plus refill rate, allowing bursts); for LLMs, count tokens as well as requests. Token bucket Add circuit breakers and per-dependency bulkheads. Fowler
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
| Need | Typical choice | Guarantee | Watch out for |
|---|---|---|---|
| Propagate DB changes to an index or memory | Debezium CDC → Kafka, or outbox | Commit-ordered per key, replayable | Slot growth, schema changes, big rows |
| Slow independent AI tasks | SQS / RabbitMQ / share groups | At-least-once, per-message retry, DLQ | Visibility timeout shorter than the task |
| Multi-step jobs with waits | Temporal or similar | Durable state, retries, timers, signals | Determinism, history size, big payloads |
| Windowed aggregates (metering, alerts) | Flink / Kafka Streams | Exactly-once state via checkpoints | Late events, state size |
| Backfills, re-embeds | Spark, batch APIs | Throughput, lower price | Cut-over with the live stream |
| Quotas, dedup keys, hot tails | Redis | Atomic scripts, sub-ms | Hot keys, fail-open policy |
| Chat history | DynamoDB / Cassandra / ScyllaDB | Partition-local ordered reads | Hot partitions |
| Vectors | pgvector → vector DB → object-storage engines | ANN recall/latency trade-off | Filtered recall, churn, memory |
| Traces, usage, scores | ClickHouse + object storage | Fast aggregates, cheap retention | Small inserts, updates |
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.
- Kafka design documentation: delivery semantics, idempotent producer, compaction.
- Debezium: the outbox pattern.
- Kleppmann, "Turning the database inside-out": the log-centric view behind derived data such as vector indexes.
- AWS Builders' Library: making retries safe with idempotent APIs.
- Notion's data lake: Debezium → Kafka → Hudi → S3, the base for later AI features.
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
- How fresh must results be: within a minute of an edit, or by tomorrow?
- Corpus size and churn? Small edits or whole-document rewrites?
- What must never happen? Usually: serving a deleted or permission-revoked document. That's stricter than freshness.
- Hosted embedding API or self-hosted model? How often will the model or chunker change? (More often than anyone expects.)
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.
- Change events: 1M/day ≈ 12/s, perhaps 120/s at a 10× peak.
- Embedding volume: re-embedding whole documents means 8M chunks/day (3.2B tokens); chunk-hash diffing cuts that to perhaps a quarter. At an assumed $0.02 per 1M tokens, streaming costs tens of dollars a day.
- Full re-embed: 160B tokens ≈ $3,200 at that price, but at an assumed 5M tokens/minute quota it takes ~32,000 minutes (about 22 days). A self-hosted model at an assumed ~3k chunks/s per GPU needs ~37 GPU-hours, about 70 minutes on 32 GPUs.
- Index: 400M × 1,024 dims × 4 bytes ≈ 1.6 TB float32, ~410 GB as int8, before graph overhead. That argues for quantization or a disk- or object-storage-backed engine.
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:
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.
Data flow
- A row commits; Debezium emits a change event keyed by
doc_idwith the commit's log position assource_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. - A worker loads the registry row and drops the event unless its version is newer than the applied one and newer than any tombstone.
- 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.
- 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. - 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
- Poison document: bounded retries, then a DLQ; never block the partition.
- Embedding provider outage: lag grows inside Kafka retention. Fail over only to the same model elsewhere, since other models' vectors are incompatible, and serve slightly stale results meanwhile.
- Stalled CDC connector: alert on replication-slot lag before the primary's disk fills.
- Chunker bug: every vector carries
chunker_version, so affected docs can be found and rebuilt by replay. - Hot document: debounce per
doc_idand process only the latest version.
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.
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.
- Notion: two years of vector search: batch plus Kafka dual path, three vector-store generations, cost numbers.
- Figma: the infrastructure behind AI search: debouncing, deduplication and quantization at billions of entries.
- Shopify: real-time ML embeddings with Dataflow: streaming GPU inference lessons.
- Debezium as part of your AI solution: CDC to embeddings to a vector store, including deletes.
- B2: Semantic Search & RAG: chunking, ANN index choices and retrieval evals.
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
- Users, activity, conversation length; do people return to old conversations?
- Retention, "temporary chat" modes, regions and residency.
- What "memory" means: user-visible facts across conversations, plus searchable past chats?
- Deletion granularity (message, conversation, memory, account) and SLA.
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.
- Writes: 240M messages/day ≈ 2.8k/s, about 8–9k/s at a 3× peak.
- Size: ~1.5 KB average (short user messages, ~2 KB replies, metadata), so ~360 GB/day and ~130 TB/year before replication; the hot last 30 days is ~11 TB.
- Reads: ~120M conversation-tail loads a day (~30 KB each), mostly from cache.
- Memory extraction once per session: 15M jobs/day ≈ 175/s; at an assumed 4k input tokens on a small model priced at an assumed $0.10 per 1M input tokens, about 60B tokens and ~$6k a day. Skipping trivial sessions could halve it.
API and data model
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.
Data flow for one turn
- 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. - 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.
- 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. - 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
- Memory poisoning: extract only from user-authored messages, never tool outputs or pasted pages, into a user-visible memory list.
- Extractor backlog after an outage: drop jobs older than a few days rather than processing stale work.
- Store choice: partitioned Postgres is simpler up to a point; DynamoDB or ScyllaDB fits this write rate. Discord's Cassandra-to-ScyllaDB move cut 177 nodes to 72 and lowered tail latency. Discord 2023
"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.
- Discord: how we store trillions of messages and the earlier billions of messages post: partitioning, hot partitions, request coalescing.
- DynamoDB partition-key design.
- B3: Agent Memory Systems: what to extract and how to consolidate it.
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
- Internal platform (teams) or external product (paying customers)? That changes whether metering must be invoice-grade.
- Which providers and self-hosted models? Is cross-provider fallback acceptable for quality and compliance?
- What interface? Mirroring a popular API (OpenAI-compatible) lowers adoption cost. Uber's GenAI Gateway made exactly that choice. Uber 2024
- Quotas per tenant in requests, tokens or dollars? Priority classes (interactive vs batch)?
- Data rules: which data may go to which provider; redaction requirements; audit retention.
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.
- Concurrency: 2,000 × 7 = 14,000 open streams at peak. If one async gateway pod comfortably holds ~2,000 streams, that's 7 pods; run about 15 across three zones for headroom.
- Upstream tokens: 2,000 × 2,850 ≈ 5.7M tokens/s ≈ 342M tokens/minute, far beyond one account's quota, so the router must spread load across providers, accounts, regions and committed capacity.
- Usage events: one per request, ~1 KB, so 2 MB/s and ~170M events/day at sustained peak. Trivial for Kafka; the hard part is exactness, not volume.
- Redis: 2–4 operations per request for quota checks is under 10k ops/s, but the largest tenants' keys are hot.
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.
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.
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
- Uber: GenAI Gateway: OpenAI-compatible interface, PII redaction, metering and audit.
- Stripe: scaling your API with rate limiters: the four limiter types and their Redis token buckets.
- AWS: fairness in multi-tenant systems.
- B1 for prompt caching and model routing; B5 for reliability patterns.
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
- How long do runs last, and how long can they wait for humans?
- Which tools have side effects (send email, issue refund)? Is there code execution?
- What must customers see: live progress, step history, replays? What cost controls?
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.
- Model calls: 3M/day ≈ 35/s (~100/s peak). At ~20k input tokens per step (context re-sent each step), 60B input tokens/day, so prompt caching is essential.
- Active runs: 100k × 20 min ÷ 1,440 ≈ 1,400 at once, perhaps 4,000 at peak. Parked runs: 25k × 8 h ÷ 24 h ≈ 8,000, which must hold no worker or sandbox.
- Sandboxes: 40k × 10 min ÷ 1,440 ≈ 280 concurrent (~850 peak); at an assumed 2 vCPU each, ~1,700 vCPUs at peak.
- Workflow history: ~8 events per step, ~240 per run, ~280 events/s cluster-wide. Comfortable, if payloads stay small.
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
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
- Provider outage mid-run: activity retries back off for minutes, then the run parks as "waiting on capacity" rather than failing.
- Deploying new agent logic with 8,000 runs parked: replay must follow the original code path, so version workflow code and test replay of production histories in CI.
- Sandbox abuse: no default egress, CPU and time limits, per-tenant quotas, anomaly detection.
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
"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.
- Temporal: workflows, activities and retry policies: the execution model in the project's own docs.
- Production-ready agents with the OpenAI Agents SDK + Temporal.
- Fowler: event sourcing.
- B4: Agent Architectures (durable execution section) and B5 (sandboxing).
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
- Volume and burstiness? Payload sizes (prompts can be 100k tokens; images may be attached)?
- Query patterns: find one trace by ID, filter by metadata, aggregate cost and latency over time, compare evaluator scores.
- Retention per plan; PII rules (must data be scrubbed before it's stored?).
- Online evals: which evaluators, at what sampling rate, and who pays for judge tokens?
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.
- Raw volume: 0.6B × 8 KB + 1.4B × 1 KB ≈ 6.2 TB/day ≈ 72 MB/s average, ~300 MB/s at peak. Kafka handles this with a modest cluster; ClickHouse needs batching.
- Stored: with columnar compression of ~5–10× on metadata and text, and large payloads offloaded to S3, roughly 1 TB/day in ClickHouse, so ~30 TB for 30-day retention, plus the S3 payload store.
- Online evals: judging 1% of LLM spans is 6M judge calls/day ≈ 70/s. At an assumed ~2k tokens per call on a small judge model, that's ~12B tokens/day, which is a real cost line. Hence per-project sampling budgets.
Data model
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
- "Too many parts" in ClickHouse from small inserts: always batch by count or time.
- One tenant floods ingestion: keying by trace_id spreads its load across partitions, so isolation must come from per-project edge quotas returning
429. - Clock skew: store client and server timestamps; order spans by parent-child structure.
- Alternative: Langfuse chose a Redis queue over Kafka for easy self-hosting, a fair trade when replay and multiple consumer groups matter less. Langfuse 2024
"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
- Langfuse: from zero to scale: the move to ClickHouse, S3 event storage and async workers.
- OpenTelemetry: sampling: head vs tail sampling trade-offs.
- ClickHouse MergeTree docs: sort keys, partitions, TTLs.
- B5: what to trace and how to design LLM judges.
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
- Volume, page counts and arrival pattern (month-end spikes are common for invoices)?
- Latency expectation: minutes, or by end of day? This decides whether batch APIs are allowed.
- Target schema per document type; accuracy bar per field; cost of an error (paying a wrong amount vs a typo in a memo field).
- System of record (an ERP) and whether it supports idempotent writes.
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.
- Human review: 200k documents/day × an assumed 1.5 minutes = 5,000 reviewer-hours/day, about 625 people on 8-hour shifts. At an assumed $20/hour, that's about $100k/day.
- Model cost: under assumed per-document prices ($0.002 small tier, $0.03 frontier tier), roughly 2M × $0.002 + 0.5M × $0.03 ≈ $19k/day, plus OCR.
- So humans dominate cost. Each percentage point of straight-through processing saves 50k reviews (~$25k/day under these assumptions), which justifies better models in the escalation tier if they reduce review volume.
Data flow
- 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.
- 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.
- 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.
- 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.
- 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
- Poison PDFs (encrypted, corrupt, 2,000 pages): size and page limits at intake, bounded retries, DLQ with reason codes and a replay tool.
- Month-end surge: queues absorb it. Prioritize by due date, autoscale OCR and extraction workers on backlog age, and divert low-priority tenants to batch.
- Review backlog breaching SLA: raise auto-approve thresholds only for low-risk fields, add reviewers, and alert on queue age by priority.
- Model drift after a provider update: canary evals on a fixed labelled set, with the extraction version recorded on every result so affected documents can be reprocessed.
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.
- Uber: advancing invoice processing with GenAI: model selection, accuracy tracking and human-in-the-loop UI.
- SQS dead-letter queues and SQS exactly-once processing: what the queue guarantees and what it doesn't.
- Docling: an open-source document parsing toolkit for the cheap tier.
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
- QPS per surface and the latency SLO (say p99 under 300 ms server-side)?
- How fast must catalog changes show: minutes for new items, seconds for out-of-stock?
- Which engagement signals and experimentation platform exist? Must explanations arrive with results?
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
- Traffic: assume 30M daily users and 20k search QPS at peak.
- Item vectors: 100M × 256 dims × 1 byte (int8) ≈ 26 GB, which fits in memory and can be replicated; 768-d float32 would be ~300 GB and need sharding. HF: embedding quantization
- ANN: at an assumed ~2–5k QPS per replica, 20k QPS needs roughly 5–10 replicas plus headroom.
- Churn: 5M item updates/day ≈ 60/s to re-embed. Shopify reported a streaming pipeline at around 2,500 embeddings per second. Shopify 2024
- LLM on the hot path: 20k QPS × 500 tokens ≈ 10M tokens/s. Ruled out on cost before latency even comes up.
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
| Stage | Budget (p99) | Notes |
|---|---|---|
| Network, auth, request fan-out | 20 ms | |
| Query understanding | 20 ms | Cache of LLM rewrites for head queries; small classifier for the tail |
| Retrieval (ANN ∥ lexical) | 40 ms | Parallel; take what returns by the deadline |
| Feature fetch | 30 ms | Batched multi-get from the online store |
| Ranking (~500 candidates) | 40 ms | GBDT or a compact DNN |
| Cross-encoder re-rank (top 50) | 40 ms | Small model on GPU; skip under load |
| Assembly and business rules | 20 ms | |
| Headroom | 90 ms | Tail 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
- Index rebuilds: blue/green versions with an alias flip, as in case 1.
- Selective filters wreck ANN recall: use filter-aware ANN, partition by a dominant filter, or brute-force score the filtered set. Qdrant: filtering
- Cold-start items: content-based item-tower embeddings make them retrievable on day one.
- Feedback loops: the model learns from what it showed, so reserve an exploration slice.
"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.
- Deep Neural Networks for YouTube Recommendations: the canonical two-stage architecture.
- PinnerSage and ItemSage: user and product embeddings in production at Pinterest.
- Feast documentation: offline/online feature stores and point-in-time joins.
- B2: hybrid search, rerankers and ANN index internals.
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
- Which sources, and push (webhooks) or pull (polling)? How many users?
- Channels: push, email, in-app? Quiet hours and time zones?
- What's the cost of a missed notification vs an annoying one? (Usually: annoying ones get the whole feature turned off.)
- Can a notification be scheduled (a "leave now" reminder) and later cancelled if circumstances change?
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.
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:
- Bucketed sorted sets in Redis: key
timers:{shard}:{minute}with members scored by due time. Pollers own shards, read due buckets, and claim items atomically with a Lua script, then hand them to delivery. Cheap and fast, but durability depends on Redis persistence, so mirror timers to a database for recovery. - A database with time-bucketed partitions: partition key
due_minute#shard(sharded to avoid a hot partition at 08:00), scanned by workers per bucket. Durable, a little slower. - Durable workflow timers: one workflow per scheduled notification with a timer; elegant when cancellation and re-validation logic is complex, at the cost of engine load. Temporal timers
- A priority-queue service: Netflix's Timestone is a priority queue with deadlines and earliest-deadline-first dequeue, built on Redis, Kafka, Flink and Elasticsearch. Netflix 2022
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.
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.
- LinkedIn: Concourse: near-real-time personalized notifications on Kafka and Samza.
- Netflix: Timestone: deadline-ordered queueing at scale.
- Apache Samza: Air Traffic Controller case study.
- B8: Personal AI Assistants: proactivity patterns at personal scale.
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
- Volume, content types (text, images, video, voice) and languages?
- Pre-publish blocking or post-publish takedown? Usually both, by surface and risk.
- Policy categories and their severity. Some (child safety, credible threats) require immediate action and legal reporting.
- Appeals process, human reviewer capacity and regulatory transparency 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%:
- LLM calls: 12.5M/day ≈ 145/s. At an assumed ~800 tokens per call (policy excerpt plus content), ~10B tokens/day, which suggests a self-hosted open-weight policy model rather than a frontier API.
- Human review: 1M items/day at an assumed 30 s each ≈ 8,300 reviewer-hours/day, which is too many. Review must be prioritized by severity × predicted reach, and low-severity, low-reach items resolved automatically.
- For scale reference, Roblox reported about 6.1B chat messages a day, text filters peaking above 750,000 requests per second, and a voice classifier peaking at 8,300 requests per second, using distilled and quantized transformer models with thousands of human experts for appeals and complex cases. Roblox 2025
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
- Adversarial adaptation (misspellings, text in images, coded language): OCR on images, character normalization, red-team data and fast policy-prompt updates.
- L2 outage: borderline items stay in reduced distribution; nothing borderline gets full reach until reviewed.
- Over-blocking: track false-positive rates by language and dialect; classifiers are often less accurate on under-represented ones.
- Partitioning: key content events by content ID for even load, and compute per-author velocity features (spam bursts) in a separate keyed aggregation.
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.
- Roblox: how we use AI to moderate content at massive scale.
- gpt-oss-safeguard user guide: writing policies for a policy-reasoning classifier.
- OpenAI moderation guide: a hosted first-pass classifier.
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.
- Inbound: verify the webhook signature, deduplicate on the provider's event ID, persist and enqueue, and return 2xx within the provider's timeout. Never do LLM work inside the webhook handler; providers retry slow handlers and you get duplicates. Webhooks can arrive out of order and some will never arrive, so treat them as hints: fetch the ticket's current state from the API before acting.
- Processing: take a per-ticket lease so two webhooks for the same ticket don't run the agent twice concurrently. Record the ticket version the agent saw.
- Outbound via outbox: the agent's decision (reply text, field changes) and an outbox row are written in one local transaction. Transactional outbox A sync worker sends them to the ticketing API with an idempotency mechanism (an idempotency key if supported; otherwise a marker such as an external ID or tag to check before posting), respecting the API's rate limits with a per-customer token bucket.
- Conflicts: a human agent may edit the ticket while the AI is drafting. Use optimistic concurrency if the API offers versions or ETags; otherwise re-read before writing and abort if the ticket changed in a way that matters (status, assignee, a new customer message). Define ownership rules: once a human takes the ticket, the AI stops writing.
- Reconciliation: a periodic incremental export (by
updated_atcursor) catches missed webhooks and drift, which is the standard safety net for any webhook integration.
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.
- Equivalence classes, validated by evals: for each use case, list the acceptable substitutes (the same model in another region or cloud; a different model that passed that use case's eval suite). Fallback chains come from this list and are never improvised.
- A degradation mode per feature, decided in advance: a smaller model, cached or precomputed results, deferral to async ("we'll email you the summary"), or a feature flag off.
- Priority shedding: under pressure, pause evals, enrichment and backfills first so interactive traffic keeps its quota.
- Circuit breakers and retry budgets at the gateway only, respecting
retry-after. Without a budget, a brown-out becomes a self-inflicted outage. Google SRE: cascading failures Fallback paths that never run are broken when you need them, so route a small share of traffic through them continuously.
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
- Encryption: TLS in transit; envelope encryption with per-tenant keys in a KMS at rest, which enables key revocation and crypto-shredding.
- Isolation wherever data flows: the tenant ID belongs in every cache key (a shared semantic cache without it leaks across tenants), in server-enforced vector-store filters or namespaces, in row-level security Postgres RLS, and in trace queries. Large or regulated tenants get siloed indexes or clusters.
- Residency: region-pinned stacks, routing by home region at the edge, same-region model endpoints. Derived stores (traces, eval sets, backups) carry the same obligations.
- Minimization: redact before third-party calls where the task allows; keep prompt payloads only as long as needed.
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.
"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.