Kafka 是 LinkedIn 在 2011 年開源的分散式系統,設計目標是每秒處理數百萬條訊息。它常被歸類成「訊息佇列」,但更準確的說法是可持久化、可重播的事件日誌平台——訊息被消費後不會消失,這一點決定了它跟傳統 MQ 幾乎所有的差異。
三大角色:生產者、Broker 叢集、消費者。訊息由上往下流動。
理解 Kafka 必須先掌握這 8 個概念。
| 比較項目 | Apache Kafka | RabbitMQ (AMQP) |
|---|---|---|
| 設計理念 | 事件日誌 / 串流平台 | 訊息佇列 / 任務分發 |
| 訊息保留 | 保留一段時間(可重播) | 消費後即刪除 |
| 消費模式 | Consumer 主動 Pull | Broker 主動 Push |
| 吞吐量 | 極高(百萬 msgs/sec) | 中等(數十萬) |
| 訊息順序 | Partition 內有序 | 單一 Queue 有序 |
| 多消費者 | 多個 Group 獨立消費 | Round-Robin 競爭 |
| 適合場景 | 事件溯源、Log 蒐集、實時分析 | 任務佇列、RPC、IoT 指令 |
從實體儲存一路往上,看 Kafka 的每個設計是為了解決什麼問題。點標題展開內容。
Kafka 的高效能來自對磁碟的順序讀寫(Sequential I/O)——不需要磁頭來回尋道,吞吐量可以逼近記憶體隨機存取。而順序讀寫之所以成立,是因為它把資料切成一層層的實體檔案:
order-topic-0/)。訊息在 Partition 內嚴格按 Offset 有序排列,也是並行消費的最小單位。00000000000000001024.log)。達到大小上限或時間限制後就開新的一個,舊 Segment 可以整批刪除,效率極高。00000000000001024.log 代表這個檔案從 Offset 1024 開始——所以要找某個 Offset 的訊息,二分搜檔名就能定位到檔案。__consumer_offsets。latest(預設):從最新訊息開始,適合正式環境;earliest:從最舊訊息開始,適合初次部署或資料重跑;none:沒有既有 Offset 就直接拋例外。
Kafka 不會無限保留訊息。可以按時間或空間設定保留策略:
log.retention.ms = 604800000log.retention.bytes = 1073741824特殊策略:Log Compaction(日誌壓縮)——針對有 Key 的訊息,只保留每個 Key 的最新版本,類似資料庫的 upsert。適合維護「最新狀態」的場景(如帳戶餘額),設定 cleanup.policy=compact。
一條訊息從 Producer 到 Broker 要經過三個階段:
hash(key) % numPartitions),相同 Key 必定進同一個 Partition,順序性因此得到保證(如同一個 userId 的訂單按序處理)。若沒有 Key,新版採 Sticky 策略——填滿一個批次才換 Partition,提升批次效率。
batch.size,預設 16KB)。達到 batch.size 或等滿 linger.ms(預設 0ms)後才批次送出。把 linger.ms 調到 5–10ms 通常能大幅提升吞吐量,代價是多了幾毫秒延遲。
compression.type=lz4 或 snappy,把整個 batch 壓縮後傳送,顯著降低網路頻寬消耗。壓縮是批次級別的,所以 batch 越大壓縮比越高——這與上一步的 linger.ms 是同一個取捨。
enable.idempotence=true):Broker 端自動去重,解決 Producer 重試導致的訊息重複。Kafka 3.0 後預設開啟。Consumer Group 的核心規則:同一個 Group 內,每個 Partition 只能被一個 Consumer 消費。這條規則同時決定了並行度的上限。
Rebalance 的觸發時機:
session.timeout.ms(預設 45s)未回應心跳session.timeout.ms 與 heartbeat.interval.ms,並盡量使用 Cooperative Sticky 分配策略(Kafka 2.4+),減少不必要的 Partition 轉移。每個 Partition 有一個 Leader 和多個 Follower。所有讀寫都走 Leader,Follower 的任務只是同步 Leader 的資料。
replica.lag.time.max.ms(預設 30s)內的副本才算數,落後太多會被踢出 ISR 名單。unclean.leader.election.enable 控制此行為,預設 false 即禁止。點卡片查看不同確認等級的安全性與效能取捨。
分散式系統中,「訊息究竟能被保證傳遞幾次」是核心挑戰。
拖動滑桿,直觀理解 Partition 與 Consumer 的分配關係。
生產速度大於消費速度時,Lag 會累積上去且不會自己消失。
從依賴引入到錯誤處理的完整生產級配置,範例是 Spring Boot 的訂單服務。
implementation 'org.springframework.kafka:spring-kafka' // spring-boot-starter-parent 已管理版本,不需手動指定
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency> 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 照樣推進,那條訊息就永遠不會再被消費到。
@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; // 兩條訊息要麼全成功,要麼全回滾
});
}
} @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,否則會永遠卡住
}
}
} // 在 @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 供人工審查
} INSERT IGNORE 防守,確保重複訊息不造成資料錯誤。max.message.bytes)。大檔案(如圖片)應存在 S3/OSS,訊息只傳遞路徑。這些指標決定你的 Kafka 叢集是否健康。
session.timeout.ms(如 60s)heartbeat.interval.ms(如 10s)max.poll.interval.ms(消費業務耗時較長時)acks=all + retries>0 + enable.idempotence=truereplication.factor≥3 + min.insync.replicas=2 + unclean.leader.election.enable=falseenable.auto.commit=false(業務成功後手動 commit Offset)enable.idempotence=true 可解決此問題。acks=all + min.insync.replicas=1:效果等同 acks=1(只有 Leader 確認)acks=all + min.insync.replicas=2(推薦):Leader + 至少 1 個 Follower 確認,即使 Leader 掛了資料也不遺失