技術筆記 Kafka
技術筆記

Kafka 核心知識

事件日誌Partition / OffsetConsumer GroupISRExactly-Once

Kafka 是 LinkedIn 在 2011 年開源的分散式系統,設計目標是每秒處理數百萬條訊息。它常被歸類成「訊息佇列」,但更準確的說法是可持久化、可重播的事件日誌平台——訊息被消費後不會消失,這一點決定了它跟傳統 MQ 幾乎所有的差異。

1M+
msgs/sec 單叢集
<10ms
端到端延遲
80%+
Fortune 100 使用

Kafka 解決了什麼問題?

服務解耦
傳統架構中 A 服務直接呼叫 B、C、D,形成複雜網格,加一個下游就得改上游。Kafka 站在中間,發布者只管寫進 Topic,訂閱者各自來讀,兩邊互不依賴。
流量削峰
促銷活動瞬間爆量,後端來不及處理。Kafka 當緩衝層把訊息先收下來,消費端依自己的速度消化,不會被瞬間流量直接壓垮。
資料重播
傳統 MQ 訊息消費後即刪除,Kafka 則保留一段時間(預設 7 天)。消費端可以把 Offset 倒回去重播過去的事件,用來除錯或重建下游資料。

架構全景圖

三大角色:生產者、Broker 叢集、消費者。訊息由上往下流動。

生產者 Producers
訂單服務Producer
用戶服務Producer
支付服務Producer
↓發送訊息(按 Topic 路由)
Kafka Cluster(由 ZooKeeper / KRaft 管理)
Broker 1Leader · P0, P1
Broker 2Follower · ISR
Broker 3Follower · ISR
order-events (P=4, R=3) user-events (P=2, R=3)
↓訂閱並消費(依 Consumer Group 分工)
Consumer Group A(訂單處理)
Pod 1P0
Pod 2P1
Consumer Group B(通知服務)
Pod 1P0, P1
重點:同一個 Topic 可以被多個 Consumer Group 獨立消費,各自維護進度、彼此不影響。這是 Kafka 與傳統 MQ 最大的差異之一。

核心術語詞典

理解 Kafka 必須先掌握這 8 個概念。

Producer(生產者)
負責將訊息發送到指定 Topic 的應用程式。可設定傳送確認等級 (acks)、分區路由策略、批次大小等參數。
Consumer(消費者)
從 Topic 的 Partition 中讀取訊息的應用程式。一個 Consumer 通常屬於某個 Consumer Group,只消費其被分配到的 Partition。
Broker(代理節點)
Kafka 叢集中的一台伺服器,負責儲存訊息、接收 Producer 訊息、提供 Consumer 讀取。多個 Broker 組成一個 Cluster。
Topic(主題)
訊息的邏輯分類標籤(類似資料庫的「表」)。Producer 指定 Topic 發送,Consumer 訂閱 Topic 接收。一個 Topic 由多個 Partition 組成。
Partition(分區)
Topic 的實體分割單位。訊息在 Partition 內有序排列,但跨 Partition 無序。Partition 數量決定了最大的並行消費數量。
Offset(偏移量)
訊息在 Partition 中的唯一序號(從 0 開始)。Consumer 透過記錄 Offset 來追蹤自己讀到哪裡,實現斷線重連後的續讀。
Consumer Group(消費者群組)
多個 Consumer 組成的群組,共同消費一個 Topic。Kafka 保證同一 Partition 在同一時刻只會被群組中的一個 Consumer 處理。
Replica(副本)
每個 Partition 在不同 Broker 上的備份。有一個 Leader(讀寫主力)和多個 Follower(備用)。副本數量由 replication.factor 設定。

Kafka vs 傳統 MQ(RabbitMQ)

比較項目Apache KafkaRabbitMQ (AMQP)
設計理念事件日誌 / 串流平台訊息佇列 / 任務分發
訊息保留保留一段時間(可重播)消費後即刪除
消費模式Consumer 主動 PullBroker 主動 Push
吞吐量極高(百萬 msgs/sec)中等(數十萬)
訊息順序Partition 內有序單一 Queue 有序
多消費者多個 Group 獨立消費Round-Robin 競爭
適合場景事件溯源、Log 蒐集、實時分析任務佇列、RPC、IoT 指令
選擇原則:需要高吞吐、資料重播、多消費者獨立讀取 → Kafka;需要複雜路由、任務優先級、即時 Push → RabbitMQ。

從實體儲存一路往上,看 Kafka 的每個設計是為了解決什麼問題。點標題展開內容。

Kafka 的高效能來自對磁碟的順序讀寫(Sequential I/O)——不需要磁頭來回尋道,吞吐量可以逼近記憶體隨機存取。而順序讀寫之所以成立,是因為它把資料切成一層層的實體檔案:

Cluster(叢集)
多台 Broker 伺服器組成的整體。由 ZooKeeper(舊)或 KRaft(新)協調管理。
Broker(節點)
叢集中的一台機器,負責儲存與提供讀寫服務。各 Partition 的 Leader 分散在不同 Broker 上以達到負載平衡。
Topic(主題)
邏輯概念,相當於資料庫的「表」。訊息按 Topic 分類,建立時要指定 Partition 數量與副本係數。
Partition(分區)
實際的資料夾(如 order-topic-0/)。訊息在 Partition 內嚴格按 Offset 有序排列,也是並行消費的最小單位。
Segment(段)
Partition 中的實體檔案(如 00000000000000001024.log)。達到大小上限或時間限制後就開新的一個,舊 Segment 可以整批刪除,效率極高。
Segment 命名規則:檔名就是該 Segment 第一條訊息的 Offset。例如 00000000000001024.log 代表這個檔案從 Offset 1024 開始——所以要找某個 Offset 的訊息,二分搜檔名就能定位到檔案。
  • Offset 是一個從 0 開始遞增的整數,代表訊息在 Partition 中的位置。
  • Consumer 消費完訊息後,會把當前 Offset 提交(commit)到 Kafka 內建的特殊 Topic:__consumer_offsets。
  • 這讓 Consumer 重啟後能從上次位置繼續消費,而不是從頭重讀。
  • 不同 Consumer Group 各自維護獨立的 Offset,互不干擾——這就是「同一份資料可以被多組人各讀各的」的實作基礎。
Offset 的三種重置策略:
latest(預設):從最新訊息開始,適合正式環境;
earliest:從最舊訊息開始,適合初次部署或資料重跑;
none:沒有既有 Offset 就直接拋例外。

Kafka 不會無限保留訊息。可以按時間或空間設定保留策略:

時間保留(預設 7 天)
log.retention.ms = 604800000
超過時間的 Segment 會被整批刪除。
空間保留
log.retention.bytes = 1073741824
Partition 超過指定大小後,從最舊的 Segment 開始刪。

特殊策略:Log Compaction(日誌壓縮)——針對有 Key 的訊息,只保留每個 Key 的最新版本,類似資料庫的 upsert。適合維護「最新狀態」的場景(如帳戶餘額),設定 cleanup.policy=compact。

一條訊息從 Producer 到 Broker 要經過三個階段:

1
Partitioner(分區選擇)
若訊息帶有 Key,會對 Key 做 Hash 取模(hash(key) % numPartitions),相同 Key 必定進同一個 Partition,順序性因此得到保證(如同一個 userId 的訂單按序處理)。若沒有 Key,新版採 Sticky 策略——填滿一個批次才換 Partition,提升批次效率。
2
Batch 批次累積
Producer 先把訊息放進本地緩衝區(batch.size,預設 16KB)。達到 batch.size 或等滿 linger.ms(預設 0ms)後才批次送出。把 linger.ms 調到 5–10ms 通常能大幅提升吞吐量,代價是多了幾毫秒延遲。
3
壓縮(Compression)
設定 compression.type=lz4 或 snappy,把整個 batch 壓縮後傳送,顯著降低網路頻寬消耗。壓縮是批次級別的,所以 batch 越大壓縮比越高——這與上一步的 linger.ms 是同一個取捨。
冪等性生產者(enable.idempotence=true):Broker 端自動去重,解決 Producer 重試導致的訊息重複。Kafka 3.0 後預設開啟。

Consumer Group 的核心規則:同一個 Group 內,每個 Partition 只能被一個 Consumer 消費。這條規則同時決定了並行度的上限。

正常情況(P=4, Consumer=2)
Consumer 1 → Partition 0, 1
Consumer 2 → Partition 2, 3
Consumer 超過 Partition(P=2, Consumer=3)
Consumer 3 完全閒置,不會有任何 Partition 分配給它。

Rebalance 的觸發時機:

  • 新增或移除 Consumer(如 Pod 擴縮容)
  • Consumer 超過 session.timeout.ms(預設 45s)未回應心跳
  • Topic 的 Partition 數量改變
Rebalance 的代價:期間所有 Consumer 會暫停消費(Stop-The-World)。生產環境要設定合理的 session.timeout.ms 與 heartbeat.interval.ms,並盡量使用 Cooperative Sticky 分配策略(Kafka 2.4+),減少不必要的 Partition 轉移。

每個 Partition 有一個 Leader 和多個 Follower。所有讀寫都走 Leader,Follower 的任務只是同步 Leader 的資料。

Leader
讀 / 寫
ISR F1
同步中
ISR F2
同步中
Out
落後
  • ISR:與 Leader 進度差距在 replica.lag.time.max.ms(預設 30s)內的副本才算數,落後太多會被踢出 ISR 名單。
  • acks=all:Producer 必須等 ISR 中所有副本都確認寫入,才算發送成功。
  • Leader 選舉:Leader 掛掉時從 ISR 名單中選出新的。若 ISR 名單為空,可能被迫選一個落後的副本,造成資料遺失——unclean.leader.election.enable 控制此行為,預設 false 即禁止。

生產者確認機制(acks)

點卡片查看不同確認等級的安全性與效能取捨。

0
Fire & Forget
發送後不管
1
Leader ACK
Leader 確認即回覆
all
ISR All ACK
全員確認(最安全)
acks=1(平衡點):等待 Leader 副本寫入後才回覆成功。風險:Leader 寫完立刻掛掉、而此時 Follower 還沒同步完成,這條訊息就遺失了。適合大多數非關鍵業務。

訊息傳遞語意

分散式系統中,「訊息究竟能被保證傳遞幾次」是核心挑戰。

At-Most-Once
最多一次
可能遺失,不重複
At-Least-Once
至少一次
不遺失,可能重複
Exactly-Once
精確一次
不遺失,不重複
At-Least-Once(最常用):Consumer 在業務處理成功後才 commit Offset。若中途失敗(如 Pod 重啟),該訊息會被重新消費。
因此業務邏輯必須實現冪等性:相同訊息處理多次,結果要跟處理一次完全相同。常見做法是資料庫唯一鍵,或事先記錄訊息 ID。

Partition 分配模擬器

拖動滑桿,直觀理解 Partition 與 Consumer 的分配關係。

Partitions
Consumers

Consumer Lag 視覺化

生產速度大於消費速度時,Lag 會累積上去且不會自己消失。

一次流量尖峰之後,Lag 為什麼回不去
Consumer Lag 是維運 Kafka 最重要的監控指標。消費速度是一條水平線——尖峰過去之後 Lag 仍停在 670,因為消費端沒有「加速追趕」的能力。要嘛增加 Consumer 數量(不超過 Partition 數),要嘛增加 Partition 數量(需 Rebalance)。

從依賴引入到錯誤處理的完整生產級配置,範例是 Spring Boot 的訂單服務。

1. 引入依賴

build.gradle
implementation 'org.springframework.kafka:spring-kafka'
// spring-boot-starter-parent 已管理版本,不需手動指定
pom.xml xml
<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka</artifactId>
</dependency>

2. 核心配置(application.yml)

yaml
spring:
  kafka:
    bootstrap-servers: localhost:9092   # Broker 地址,生產環境填多個

    producer:
      acks: all                          # 最安全:等待 ISR 全部確認
      retries: 3                         # 失敗自動重試
      enable-idempotence: true           # 冪等性:自動去重(Kafka 3.0+ 預設開啟)
      batch-size: 32768                  # 批次大小 32KB(預設 16KB)
      linger-ms: 5                       # 等待 5ms 累積批次(提升吞吐量)
      compression-type: lz4              # 壓縮格式(lz4 速度最快)
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer

    consumer:
      group-id: order-service-group
      enable-auto-commit: false          # 必須手動提交!自動提交可能在業務失敗後仍 commit
      auto-offset-reset: latest          # latest=從最新訊息消費;earliest=從最舊訊息消費
      max-poll-records: 100              # 每次 poll 最多取 100 筆
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
      properties:
        spring.json.trusted.packages: "com.example.dto"  # 信任反序列化的包路徑

    listener:
      ack-mode: manual_immediate         # 配合 enable-auto-commit: false
enable-auto-commit: false 是這份設定裡最關鍵的一行。自動提交是按時間間隔 commit 的,跟你的業務有沒有處理成功完全無關——業務拋例外了 Offset 照樣推進,那條訊息就永遠不會再被消費到。

3. 生產者實作

java
@Service
@RequiredArgsConstructor
public class OrderEventProducer {

    private final KafkaTemplate<String, Object> kafkaTemplate;

    public void publishOrderCreated(Order order) {
        // 用 userId 作為 Key → 同一個用戶的訂單 Event 進入同一個 Partition
        // 保證該用戶的事件被同一個 Consumer 按序處理
        ListenableFuture<SendResult<String, Object>> future =
            kafkaTemplate.send("order-created", order.getUserId(), order);

        // 非同步 Callback:確認是否成功送達 Broker
        future.addCallback(
            result  -> log.info("訊息發送成功,Offset: {}", result.getRecordMetadata().offset()),
            failure -> log.error("訊息發送失敗: {}", failure.getMessage())
        );
    }

    // 事務性發送(需要 Exactly-Once)
    public void publishWithTransaction(Order order) {
        kafkaTemplate.executeInTransaction(t -> {
            t.send("order-created", order.getUserId(), order);
            t.send("audit-log",     order.getUserId(), new AuditLog(order));
            return true; // 兩條訊息要麼全成功,要麼全回滾
        });
    }
}

4. 消費者實作

java
@Service
@RequiredArgsConstructor
@Slf4j
public class OrderEventConsumer {

    private final OrderService orderService;

    // concurrency = 3 → 啟動 3 個 Consumer 執行緒(不超過 Partition 數)
    @KafkaListener(
        topics   = "order-created",
        groupId  = "order-service-group",
        concurrency = "3"
    )
    public void handleOrderCreated(
            @Payload Order order,
            @Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
            @Header(KafkaHeaders.OFFSET) long offset,
            Acknowledgment ack) {

        log.info("消費訊息 | partition={}, offset={}, orderId={}", partition, offset, order.getId());

        try {
            orderService.process(order); // 業務邏輯(必須具備冪等性!)
            ack.acknowledge();           // 業務成功才 commit Offset
        } catch (Exception e) {
            log.error("處理失敗,orderId={}, error={}", order.getId(), e.getMessage());
            // 不呼叫 ack.acknowledge() → Offset 不推進 → 下次仍會重試
            // 注意:重試次數達上限後應進入 DLQ,否則會永遠卡住
        }
    }
}

5. 錯誤處理與死信佇列(DLQ)

java
// 在 @Configuration 類別中配置
@Bean
public ConcurrentKafkaListenerContainerFactory<String, Object> kafkaListenerContainerFactory(
        ConsumerFactory<String, Object> consumerFactory,
        KafkaTemplate<String, Object> kafkaTemplate) {

    ConcurrentKafkaListenerContainerFactory<String, Object> factory =
        new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory);

    // 設定重試策略:最多重試 3 次,每次間隔 2 秒
    factory.setCommonErrorHandler(
        new DefaultErrorHandler(
            // 達到最大重試次數後,將訊息轉發到 DLQ Topic
            new DeadLetterPublishingRecoverer(kafkaTemplate,
                (record, ex) -> new TopicPartition(record.topic() + ".DLQ", record.partition())
            ),
            new FixedBackOff(2000L, 3L) // 間隔 2s,最多重試 3 次
        )
    );

    return factory;
}

// 監聽 DLQ 的 Consumer(人工介入或告警)
@KafkaListener(topics = "order-created.DLQ", groupId = "dlq-group")
public void handleDLQ(ConsumerRecord<String, Object> record) {
    log.error("DLQ 收到問題訊息,Key={}, Value={}", record.key(), record.value());
    // 可以:發告警、儲存到 DB 供人工審查
}

生產環境最佳實踐

冪等性設計
訊息至少會被消費一次(At-Least-Once)。資料庫層用唯一鍵(Unique Key)或 INSERT IGNORE 防守,確保重複訊息不造成資料錯誤。
Key 的選擇
需要順序保證的相關事件,務必使用相同 Key(如 userId、orderId)。沒有順序需求的訊息不要給 Key,讓 Kafka 自動做負載平衡。
Consumer 數量上限
Consumer 數量不能超過 Partition 數,超過的完全閒置。規劃 Partition 數量時要考慮未來擴容——上線後增加 Partition 會觸發 Rebalance,且已有 Key 的分區路由會改變。
避免大訊息
Kafka 預設訊息大小上限 1MB(max.message.bytes)。大檔案(如圖片)應存在 S3/OSS,訊息只傳遞路徑。

必監控的 5 大指標

這些指標決定你的 Kafka 叢集是否健康。

Consumer Lag kafka-consumer-groups.sh --describe --group <group>
最重要的一個。各 Partition 未消費訊息數量的總和。持續上升代表消費端跟不上生產端,要立即擴容或排查消費者效能。
Under-Replicated Partitions kafka.server:UnderReplicatedPartitions
值 > 0 代表有 Partition 的副本數量不足(Follower 落後或已掛)。長期不為 0 代表叢集有資料遺失風險。
Active Controller Count kafka.controller:ActiveControllerCount
叢集中應有且只有 1 個 Active Controller。值為 0 代表無 Controller(嚴重);> 1 代表腦裂(Split Brain)。
Request Latency kafka.network:RequestMetrics.TotalTimeMs
Producer/Consumer 請求的端到端延遲。突然上升可能是 Broker 負載過高、網路問題或磁碟 I/O 瓶頸。
Disk Usage Per Broker kafka.log:LogFlushRateAndTimeMs + OS disk usage
訊息存在磁碟上,用量超過 80% 可能導致 Broker 崩潰。要結合 Retention 策略控制成長,並設磁碟使用率告警。

常見問題排查 SOP

Consumer Lag 持續上升
1. 確認消費者邏輯是否有效能問題(查 GC、DB 慢查詢)
2. 增加 Consumer 數量(需 ≤ Partition 數)
3. 考慮增加 Partition 數(需 Rebalance,挑業務低峰進行)
頻繁 Rebalance
1. 調高 session.timeout.ms(如 60s)
2. 調低 heartbeat.interval.ms(如 10s)
3. 調高 max.poll.interval.ms(消費業務耗時較長時)
4. 升級 Kafka 並改用 Cooperative Sticky 策略
訊息重複消費
1. 確認業務邏輯有冪等保護(唯一鍵)
2. 避免在業務處理前就提前 commit Offset
3. 若要完全精確一次,才動用 Exactly-Once(開銷較高)
訊息積壓在 DLQ
1. 查 DLQ 訊息的錯誤原因(Schema 不符?業務邏輯異常?)
2. 修復問題後,把 DLQ 訊息重新發回主 Topic 重處理
3. DLQ 也要設定 Retention,避免一直累積

高頻面試題

三層保障同時配置才能確保不遺失:
Producer 端:acks=all + retries>0 + enable.idempotence=true
Broker 端:replication.factor≥3 + min.insync.replicas=2 + unclean.leader.election.enable=false
Consumer 端:enable.auto.commit=false(業務成功後手動 commit Offset)
Kafka 只保證同一個 Partition 內的有序性,跨 Partition 不保證順序。
若需要相關訊息有序處理(如同一個用戶的操作),需要:
1. 發送時帶相同的 Key(如 userId),Kafka 會將相同 Key 路由到同一個 Partition
2. Consumer 端的 concurrency 不超過 Partition 數,且每個 Partition 只由一個 Consumer 處理
重要:若 Producer 開啟 retries,在不開啟冪等性時可能因重試導致亂序(舊訊息後到)。開啟 enable.idempotence=true 可解決此問題。
會觸發 Rebalance(再平衡):
1. Group Coordinator(Broker 中的一個角色)偵測到新成員加入
2. 通知所有 Consumer 停止消費(Stop-The-World)
3. 重新分配 Partition 給所有 Consumer
4. 消費者繼續從各自被分配到的 Partition 消費
注意:若新增的 Consumer 數量已超過 Partition 數,多出的 Consumer 不會分配到任何 Partition,完全閒置。
acks=1:Leader 寫入成功即回覆。效能較好,但若 Leader 在 Follower 同步前就掛掉,這條訊息會遺失。
acks=all:必須等待 ISR 中所有副本都確認,才回覆成功。效能略低(多一次網路往返),但資料絕對不會遺失(只要有 min.insync.replicas 個副本存活)。
選擇:金融交易、訂單系統 → acks=all;日誌蒐集、埋點資料可容忍少量遺失 → acks=1。
冪等性定義:對同一個操作執行一次或多次,結果完全相同。
為什麼需要:Kafka 預設的傳遞語意是 At-Least-Once,在以下情況訊息會被重複消費:
- Consumer 處理完訊息但在 commit Offset 前崩潰重啟
- Rebalance 發生時,部分訊息可能被不同 Consumer 重複讀取
實作方式:資料庫層添加 messageId 唯一鍵,INSERT 時使用 ON DUPLICATE KEY 忽略重複;或先查後寫(check-then-act)。
沒有一個萬用公式,但有幾個考量原則:
1. 吞吐量:目標 QPS ÷ 單 Partition 吞吐量(通常 10~100MB/s)= 最低 Partition 數
2. Consumer 數量:Partition 數 ≥ 計畫的最大 Consumer 數(否則部分 Consumer 閒置)
3. 保留倍數:設定為 Consumer 數量的整數倍,負載最均衡(如 Consumer=4,Partition=8 或 12)
注意:Partition 數量只能增加不能減少(減少需重建 Topic)。設計時預留 2 倍擴充空間。一般建議起步 6~12 個,視業務增長再調整。
設計理念不同:RabbitMQ 是傳統訊息佇列(訊息消費後刪除,Push 模式);Kafka 是事件日誌平台(訊息保留可重播,Pull 模式)。
關鍵差異:
- 吞吐量:Kafka 可達百萬 msg/s,RabbitMQ 通常數十萬
- 多消費者:Kafka 支援多個 Consumer Group 獨立消費同一 Topic;RabbitMQ 競爭消費
- 訊息重播:Kafka 支援(保留一段時間);RabbitMQ 不支援(消費後刪除)
選擇:高吞吐量、事件溯源、多消費者 → Kafka;複雜路由、任務優先級、Pub-Sub 廣播 → RabbitMQ。
min.insync.replicas(ISR 最小同步副本數):設在 Broker 或 Topic 層級,規定 ISR 中至少要有幾個副本確認寫入,Producer 才能收到成功回覆。
與 acks=all 的搭配:
- acks=all + min.insync.replicas=1:效果等同 acks=1(只有 Leader 確認)
- acks=all + min.insync.replicas=2(推薦):Leader + 至少 1 個 Follower 確認,即使 Leader 掛了資料也不遺失
最佳實踐:replication.factor=3,min.insync.replicas=2。這樣可以容忍一台 Broker 掛掉,同時資料不遺失。