跳至主要內容

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

PatternKafka shapeReach for it when
Work queue一個 consumer group;每個 partition 由一個 member 持有Jobs 需要 buffering 與 per-key order;retry 由你負責
Event streamRetained log + offsetsContinuous 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。


text
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 分開。


text
topic: match-events

partition 0  game-17: kickoff → goal → booking
partition 1  game-42: kickoff → substitution
partition 2  game-91: kickoff → goal

Guarantee 刻意很窄: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。


text
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:keyvaluetimestampheaders。Value 是 payload。Headers 放 schema version、correlation ID 與 traceparent 等 metadata。Key 通常決定 partition。


ts
await producer.send({
  topic: "match-events",
  messages: [
    {
      key: matchId,
      value: JSON.stringify({ type: "goal", playerId }),
      headers: { "schema-version": "2", traceparent },
    },
  ],
})

常用的 mental model:


text
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 才有價值。


RequirementExampleWhy Kafka fits
Async processingTranscode uploaded videoIngest 與 workers 獨立 scale
Per-key order從 waiting room 放用戶進入一把 key 留在一個 partition
Stream processing用 Flink aggregate ad clicksConsumers 處理 continuous retained flow
Pub/sub把 live comments 送給多個 services每個 consumer group 都拿到自己的 copy
Replay修 bug 後重建 projectionReset offsets,再讀一次 retained history

做 video transcoding 時,Kafka 帶一條小 event,裡面只有 videoId 與 S3 URL。S3 才裝 video。Worker poll 到 event 後再下載。


text
Upload → S3
       → Kafka { videoId, s3Url } → Transcoder → renditions

Kafka 不是 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 TB10k messages per second 想。真實 capacity 會隨 hardware、record size、replication、acknowledgements、compression 與 workload 巨幅變化。


text
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:matchIdorderIdaccountId
  • 它應該 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 卻仍健康。


StrategyWhat it buysWhat it costs
No key一段時間後分佈均勻沒有 per-entity order
Random salt一把 hot key 變 N 把 keysConsumers 要 merge N 條 partial streams
Compound key按 region 等真實 dimension 分散Order 變成 per compound key
Backpressure保護 brokers 與 downstream systemsProducer 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 acksSuccess meansTrade-off
0Producer 發出了 requestLatency 最低;loss 可能 silent
1Leader append 完成Leader failure 可能丟 unreplicated record
all每個 required in-sync replica 都 ackDurability 最強;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、acksmin.insync.replicas。Write 沒有等 copies,copies 就不是 guarantee。



8. Consumer failure、offsets 與 rebalancing

Consumers pull records。每個 consumer group 存自己的 committed offsets,所以兩個 groups 能獨立、以不同速度讀同一個 topic。


text
poll record
  → perform durable side effect
  → commit offset

Side 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。


ts
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:


text
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.msretention.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。


text
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 都會留下。



Recap Q&A