Apache Kafka 是一個分散式、可 replay 的 log。Producers append records。Consumers 按 offset pull。Partitions 把 log 分到多個 brokers,並保住一把 key 的順序。
這篇 note 依 Evan 的 walkthrough。通用 queue pattern —— outbox、idempotent consumers、backpressure 與 broker choice —— 見 System Design 裡的 Message Queues。Production 的 Hono/MSK 購票路徑見 打造 Event-Driven 票務 Backend。這篇 note 講的是 log。
Pattern Map
| Pattern | Kafka shape | Reach for it when |
|---|---|---|
| Work queue | 一個 consumer group;每個 partition 由一個 member 持有 | Jobs 需要 buffering 與 per-key order;retry 由你負責 |
| Event stream | Retained log + offsets | Continuous processing、replay 與 ordered history 很重要 |
| Pub/sub | 每個 subscriber 一個 consumer group | 多個獨立 services 都要讀到同一個 event |
Kafka 三種都能表達,不代表三種都該預設用它。如果需求只有「執行這份 job 並 retry」,SQS 用更少 machinery 就給 delayed retries 與 DLQ。只有 retention、replay、high throughput 或多個獨立 readers 是產品要求時,Kafka 的成本才值得。
1. Kafka 為什麼存在
想像一個 World Cup 網站發佈 goals、bookings 與 substitutions。一個 producer 寫 events。一個 consumer 更新 live page。
Match feed → Producer → Queue → Consumer → Live site現在讓一千場比賽同時進行。一個 queue 與一個 consumer 成了 bottleneck。把 events 隨機分到多個 queues 雖然加了 capacity,卻丟了 causality:goal 可能出現在 kickoff 之前。
Kafka 的答案是按 game ID partition。一個 game 的每個 event 都 hash 到同一個 partition,所以比賽內的順序保住了。Consumer group 可以加 workers,又不會讓 group 內兩個 workers 同時持有同一個 partition。Topics 把 soccer 與 basketball 分開。
topic: match-events
partition 0 game-17: kickoff → goal → booking
partition 1 game-42: kickoff → substitution
partition 2 game-91: kickoff → goalGuarantee 刻意很窄:partition 內有序。Partition 0 與 partition 1 之間沒有有用的 global order。
Failure: 說「Kafka 保證順序」,卻隨機給每個 event 選 key。Kafka 只保住一個 partition 內的 append order。必須有序的 records 必須共用一把 key。
2. Brokers、topics 與 partitions
Kafka cluster 由 brokers 組成。Broker 是存 partition replicas、回應 producer 與 consumer requests 的 server。
Topic: match-events
Broker A Broker B Broker C
P0 leader P1 leader P2 leader
P1 follower P2 follower P0 follower- Topic 是 clients publish 與 subscribe 的 logical stream。
- Partition 是一條有序、immutable、append-only 的 log。它是 storage、ordering、replication 與 consumer parallelism 的單位。
- Record 是 log 裡的一項。Kafka 在 consume 後仍保留它;read 不會 delete。
- Offset 是 record 在一個 partition 內的位置。它不是 global,也不是 timestamp。
一個 consumer 可以持有多個 partitions。同一 consumer group 裡,一個 partition 同一時間只能由一個 consumer 持有。十二個 consumers 配六個 partitions,最多只有六個 active consumers。
Failure: Topic 仍只有一個 partition,卻一直加 consumers。十一個 workers 都在 idle。Partitions 決定 consumer group 內 parallelism 的上限。
3. Record shape 與 write path
一條 record 有四個實用 fields:key、value、timestamp 與 headers。Value 是 payload。Headers 放 schema version、correlation ID 與 traceparent 等 metadata。Key 通常決定 partition。
await producer.send({
topic: "match-events",
messages: [
{
key: matchId,
value: JSON.stringify({ type: "goal", playerId }),
headers: { "schema-version": "2", traceparent },
},
],
})常用的 mental model:
partition = hash(key) % partitionCount
producer
→ fetch cluster metadata
→ choose partition from key
→ send to that partition's leader
→ leader appends
→ followers replicate
→ consumer polls and advances its offset沒有 key 時,現代 producers 會把 batches 分散到多個 partitions。Distribution 變好;related-record ordering 消失。Custom partitioner 可以編碼另一套 policy,但所有 producers 都必須同意它。
改變 partition count 也會改變 hash(key) % partitionCount。舊 key 的新 records 可能移到另一個 partition,所以加 partitions 會打破變更前後的 per-key order。
Failure: 把 timestamp 當 ordering guarantee,或在 live ordered topic 上加 partitions 卻沒有規劃 key remapping。權威順序是一個 partition 內的 offset。
4. 什麼時候用 Kafka
當 log 本身解決一項 requirement 時,Kafka 才有價值。
| Requirement | Example | Why Kafka fits |
|---|---|---|
| Async processing | Transcode uploaded video | Ingest 與 workers 獨立 scale |
| Per-key order | 從 waiting room 放用戶進入 | 一把 key 留在一個 partition |
| Stream processing | 用 Flink aggregate ad clicks | Consumers 處理 continuous retained flow |
| Pub/sub | 把 live comments 送給多個 services | 每個 consumer group 都拿到自己的 copy |
| Replay | 修 bug 後重建 projection | Reset offsets,再讀一次 retained history |
做 video transcoding 時,Kafka 帶一條小 event,裡面只有 videoId 與 S3 URL。S3 才裝 video。Worker poll 到 event 後再下載。
Upload → S3
→ Kafka { videoId, s3Url } → Transcoder → renditionsKafka 不是 business row 的 source of truth,也不是 blob storage。如果 HTTP write 與 event 必須一致,就在一個 database transaction 裡 persist business row 與 outbox row,再 relay event。Outbox 見 System Design 裡的 Message Queues。
Failure: 因為 transcoder 是 async,就把 1 GB video 放進 Kafka。Blob 放 object storage,log 裡只放 pointer。也不要只因為 replay 聽起來有用,就給一次性 email job 選 Kafka。
5. Scale 從 partition key 開始
Scale 前先估 records per second、average record size、retention、replication factor 與 consumer work。Evan 給的 interview baselines 刻意很粗:records 儘量低於約 1 MB,一台配置不錯的 broker 大致按 1 TB 與 10k messages per second 想。真實 capacity 會隨 hardware、record size、replication、acknowledgements、compression 與 workload 巨幅變化。
ingress bytes/day
= records/sec × average bytes × 86,400
stored bytes
≈ ingress bytes/day × retention days × replication factor加 brokers 只增加潛在 storage 與 network capacity,不會神奇地拆開現有 topic。Topic 要有足夠 partitions,replicas 也要 reassign,新 brokers 才能接住它的 load。
Partition key 是最主要的 design decision:
- 它必須保住產品真正需要的 order:
matchId、orderId或accountId。 - 它應該 high-cardinality,而且 traffic distribution 均勻。
- 它不能比 ordering boundary 更寬。如果只有每個
invoiceId需要有序,用tenantId會把一個大 tenant 壓在一個 partition 上。 - 它是 data contract 的一部分。改 key 就會改 ordering 與 stateful consumer behavior。
Failure: 只說「我們會加 brokers」,卻說不出 key、partition count、throughput 或 retention。再多機器也救不了一個 hot partition。
6. Hot partitions
Hot partition 收到遠多於 peers 的 traffic。Ad-click topic 按 adId key 看似均勻,直到一個 campaign viral。一個 leader saturated,cluster average 卻仍健康。
| Strategy | What it buys | What it costs |
|---|---|---|
| No key | 一段時間後分佈均勻 | 沒有 per-entity order |
| Random salt | 一把 hot key 變 N 把 keys | Consumers 要 merge N 條 partial streams |
| Compound key | 按 region 等真實 dimension 分散 | Order 變成 per compound key |
| Backpressure | 保護 brokers 與 downstream systems | Producer latency 更高,或拒絕 work |
如果一個 celebrity account 的 exact order 不能退讓,它就是 serial workload。沒有 partitioning trick 能讓一條 ordered sequence 無限 parallel。減少每條 record 的工作、batch,或改變 requirement。
監控 per-partition bytes、requests 與 consumer lag。Cluster averages 會藏住 skew。
Failure: Salt key 後仍承諾 unsalted entity 的 total order。Load 之所以散開,是因為 order boundary 已經變了。
7. Durability 是 configuration
每個 partition 有一個 leader 與位於其他 brokers 的 follower replicas。Producers 寫 leader。Followers fetch 它的 log。Controller 追蹤 broker health,在 failure 後從 in-sync replicas (ISR) 裡 elect 新 leader。
Replication factor 為 3,意思是三份總 copies:一個 leader、兩個 followers。
Producer acks | Success means | Trade-off |
|---|---|---|
0 | Producer 發出了 request | Latency 最低;loss 可能 silent |
1 | Leader append 完成 | Leader failure 可能丟 unreplicated record |
all | 每個 required in-sync replica 都 ack | Durability 最強;latency 更高 |
acks=all 不是等待每個 configured replica,無論它健不健康。它等的是 ISR requirement。要配一個有意義的 min.insync.replicas;否則「all」仍可能只有一個 live replica。Durability 重要時要 disable unclean leader election,否則 out-of-date replica 可能成為 leader,丟掉已 acknowledged data。
Kafka 不是魔法般 always available。Interview 裡有用的回答是 failure domain:一個 broker 可以 fail,而 replicated partitions 繼續工作。整個 cluster、region、壞 configuration 或 operator 仍可能讓它 fail。
Failure: 說 replication factor 3 就一定安全承受任意兩台 broker failures,卻不看 ISR、replica placement、acks 與 min.insync.replicas。Write 沒有等 copies,copies 就不是 guarantee。
8. Consumer failure、offsets 與 rebalancing
Consumers pull records。每個 consumer group 存自己的 committed offsets,所以兩個 groups 能獨立、以不同速度讀同一個 topic。
poll record
→ perform durable side effect
→ commit offsetSide effect 前 commit,crash 就會丟 work。Side effect 後 commit,effect 與 commit 之間 crash 就會 reprocess record。因此 Kafka default 是 at-least-once。Consumer 必須用 event ID、unique inbox record 與 idempotent business write 讓 duplicate 安全。
Consumer join、leave 或停止 polling 時,group 會 rebalance,在 members 之間移動 partitions。受影響 partitions 會暫停 processing。Cooperative rebalancing 減少 disruption,但沒有消除保持 poll loop healthy 與 handler bounded 的要求。
- 看每個 partition 的 consumer lag,不只看 average lag。
- 保持 consumer unit of work 小。
- 按真實 processing time 設置 poll 與 session timeouts。
- Shutdown 前先停止 polling,finish 或 abandon in-flight work,只 commit 完成的 records。
Failure: Database write 還沒完成就 auto-commit offsets。Dashboard 說 caught up;business event 已經消失。相反的 failure 是 commit 後做 non-idempotent side effect,recovery 時無法安全 replay。
9. Retry 是 design 的一部分
Producer request 會 ambiguously fail:Kafka 可能已經 append record,只是 acknowledgment 沒回到 producer。啟用 idempotence 與 retries,讓 broker 能 deduplicate 同一 producer session 的 retransmissions。
const producer = kafka.producer({
idempotent: true,
retry: { retries: 5, initialRetryTime: 100 },
})Producer idempotence 不會讓整個 business workflow exactly once。它不會 deduplicate 兩個 application requests,也不會讓 database side effect 與 offset commit atomic。
Kafka 沒有 SQS 那種 queue-native delayed consumer retries。常見 design:
main topic
→ consumer fails
→ retry topic(s) with attempt and next-at metadata
→ retry consumer
→ dead-letter topic after the ceiling使用 bounded attempts、帶 jitter 的 exponential backoff,並觀測 dead-letter topic。保留 original event ID 與 error context。Poison record 不能永遠堵住 partition。
Failure: Catch exception 後 commit offset,再 log「will retry」,但根本沒有 retry event。另一個是把 permanent schema error 永遠 retry 在一個 ordered partition 的頭部。
10. Throughput 與 retention
Kafka 靠 sequential appends、batching、compression 與 parallel partitions 拿 throughput。
- Batch records,讓一次 network request 與 disk append 帶多個 events。更大 batches 提升 throughput,卻增加 waiting latency。
- Compress batch,可選 LZ4、Snappy、Zstd 或 Gzip。更少 network 與 disk,代價是 CPU。
- Partition evenly. Batching 救不了一個 hot leader。
- Keep payloads small. Large records 會同時吃 broker memory、network、replication bandwidth 與 consumer fetch capacity。
Kafka 通過 retention.ms 與 retention.bytes 按時間和/或 partition size 保留 records。常見 broker default 是七天。無論 consumers 有沒有讀過,retention 都照樣執行。
更長 retention 能 replay 更久之前的 history,卻會放大 storage、recovery time 與 cost。Log compaction 是另一種 policy:保留每把 key 的 latest value 與 tombstones,而不是永遠保留每個 event。它適合 rebuildable current state,不適合 immutable audit history。
retention answers: how much history can be replayed?
compaction answers: what is the latest value for each key?Failure: 承諾 90-day replay,disk 卻只按一天來 size;或者給 audit stream 開 compaction,還以為每個 intermediate event 都會留下。