制性能實(shí)測(cè):acks=all 在百萬級(jí)吞吐下的延遲代價(jià)有多大)
一、一個(gè)參數(shù)引發(fā)的架構(gòu)爭(zhēng)議Kafka 生產(chǎn)者發(fā)送消息時(shí)acks參數(shù)決定了消息算不算發(fā)送成功acks 值含義可靠性吞吐量0生產(chǎn)者發(fā)完就不管了最低可能丟數(shù)據(jù)最高1Leader 寫入成功即確認(rèn)中等Leader 宕機(jī)可能丟較高all/-1所有 ISR 副本都寫入成功最高最低爭(zhēng)論永遠(yuǎn)圍繞同一個(gè)問題acksall 到底慢多少值不值得有人說acksall 性能差 3 倍有人說幾乎沒影響。事實(shí)是——取決于你怎么測(cè)。本文用 6 組實(shí)驗(yàn)給你真實(shí)數(shù)據(jù)。二、acks 機(jī)制源碼分析2.1 生產(chǎn)者發(fā)送鏈路2.2 Broker 端處理邏輯在 Broker 端消息寫入的核心邏輯在ReplicaManager.appendRecords()中// Kafka ReplicaManager 源碼簡(jiǎn)化版展示 acks 邏輯 def appendRecords( timeout: Long, requiredAcks: Short, // 就是 acks 參數(shù) internalTopicsAllowed: Boolean, entriesPerPartition: Map[TopicPartition, MemoryRecords], responseCallback: Map[TopicPartition, PartitionResponse] Unit ): Unit { ? // 1. 參數(shù)校驗(yàn)acks 必須是 0, 1, 或 -1 if (requiredAcks ! 0 requiredAcks ! 1 requiredAcks ! -1) throw new InvalidRequiredAcksException(...) ? // 2. 寫入 Leader 本地日志 val localProduceResults appendToLocalLog( internalTopicsAllowed, entriesPerPartition, requiredAcks ) ? // 3. 根據(jù) acks 值決定響應(yīng)策略 requiredAcks match { case 0 // acks0不需要等待任何確認(rèn)立即回調(diào) responseCallback(localProduceResults.mapValues(_.toPartitionResponse)) ? case 1 // acks1Leader 寫入完成即可回調(diào) // 如果 Leader 寫入失敗返回錯(cuò)誤否則直接響應(yīng) responseCallback(localProduceResults.mapValues(_.toPartitionResponse)) ? case -1 // acksall需要等待所有 ISR 副本同步完成 // 創(chuàng)建 DelayedProduce掛起請(qǐng)求等待 follower 拉取 val produceMetadata ProduceMetadata(requiredAcks, localProduceResults) val delayedProduce new DelayedProduce( timeout, produceMetadata, this, responseCallback, localProduceResults ) // 放入 DelayedOperationPurgatory延遲操作隊(duì)列 delayedProducePurgatory.tryCompleteElseWatch(delayedProduce) } }2.3 acksall 的延遲等待機(jī)制acksall時(shí)消息不是寫完 Leader 就返回而是放入DelayedProducePurgatory等待所有 ISR 副本同步// DelayedProduce.tryComplete() 核心邏輯簡(jiǎn)化 override def tryComplete(): Boolean { produceMetadata.produceStatus.foreach { case (topicPartition, status) // 檢查每個(gè)分區(qū)的副本同步狀態(tài) val partition replicaManager.getPartition(topicPartition) val leaderHWIncremented partition.getReplica(leaderReplicaId) match { case Some(leaderReplica) // 獲取 Leader 的高水位HW val leaderHW leaderReplica.highWatermark ? // 檢查所有 ISR 副本是否都拉取到了這條消息 status.requiredOffset match { case Some(requiredOffset) // 關(guān)鍵只有當(dāng) HW requiredOffset 時(shí)才算所有副本同步完成 if (leaderHW requiredOffset) { true // 同步完成 } else { return false // 還有副本沒同步完繼續(xù)等待 } case None // Leader 寫入失敗不需要等待 true } case None true } } // 所有分區(qū)都滿足條件完成延遲操作 forceComplete() }關(guān)鍵點(diǎn)acksall的延遲取決于follower 副本拉取 Leader 數(shù)據(jù)的速度。如果 follower 落后很多延遲就會(huì)很高。2.4 ISR 與 acksall 的關(guān)系如果 ISR 中某個(gè)副本長(zhǎng)時(shí)間不拉取acksall就會(huì)一直等待直到超時(shí)request.timeout.ms。三、壓測(cè)環(huán)境與方案設(shè)計(jì)3.1 集群配置組件配置Broker 數(shù)量3 節(jié)點(diǎn)服務(wù)器規(guī)格8C 16G × 3SSD 500GBKafka 版本3.6.0Topic 配置6 分區(qū) × 3 副本min.insync.replicas2消息大小1KB典型業(yè)務(wù)消息消息格式JSON3.2 生產(chǎn)者配置// 壓測(cè)用 Producer 配置 Properties props new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, broker1:9092,broker2:9092,broker3:9092); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); ? // 核心變量 props.put(ProducerConfig.ACKS_CONFIG, acks); // 0 / 1 / all ? // 固定參數(shù) props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384); // 16KB batch props.put(ProducerConfig.LINGER_MS_CONFIG, 5); // 最多等 5ms 湊 batch props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, lz4); // LZ4 壓縮 props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 67108864); // 64MB 緩沖 props.put(ProducerConfig.RETRIES_CONFIG, 3); props.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 5); props.put(ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, 30000); // 30s 超時(shí)3.3 壓測(cè)場(chǎng)景場(chǎng)景編號(hào)場(chǎng)景名稱說明S1單分區(qū)串行1 分區(qū)1 個(gè) Producer 線程測(cè)試單線程極限S2多分區(qū)并行6 分區(qū)6 個(gè) Producer 線程測(cè)試并行吞吐S3高并發(fā)多分區(qū)6 分區(qū)20 個(gè) Producer 線程測(cè)試高并發(fā)S4弱副本集群ISR 僅 1 個(gè)副本模擬副本同步慢S5大消息壓力消息 10KB測(cè)試帶寬瓶頸S6持續(xù)寫入穩(wěn)定性連續(xù)寫入 30 分鐘觀察延遲分布3.4 監(jiān)控指標(biāo)// 使用 Kafka Producer Metrics API 采集指標(biāo) MapMetricName, ? extends Metric metrics producer.metrics(); ? // 核心指標(biāo) double recordSendRate getMetric(metrics, record-send-rate); // 發(fā)送速率條/秒 double recordAckRate getMetric(metrics, record-ack-rate); // 確認(rèn)速率條/秒 double requestLatencyAvg getMetric(metrics, request-latency-avg); // 平均請(qǐng)求延遲ms double requestLatencyMax getMetric(metrics, request-latency-max); // 最大請(qǐng)求延遲ms double recordQueueTimeAvg getMetric(metrics, record-queue-time-avg); // 隊(duì)列等待時(shí)間 double ioWaitTimeNsAvg getMetric(metrics, io-wait-time-ns-avg); // IO 等待時(shí)間 double batchSizeAvg getMetric(metrics, batch-size-avg); // 平均 batch 大小 double compressionRate getMetric(metrics, compression-rate-avg); // 壓縮率四、壓測(cè)數(shù)據(jù)對(duì)比4.1 場(chǎng)景 S1單分區(qū)串行指標(biāo)acks0acks1acksall吞吐量 (msg/s)85,20062,10038,400吞吐量 (MB/s)83.260.637.5平均延遲 (ms)1.22.87.6P99 延遲 (ms)3.16.518.2P99.9 延遲 (ms)5.812.335.6CPU 使用率 (%)455258分析單分區(qū)場(chǎng)景下acksall吞吐量?jī)H為acks0的 45%。因?yàn)閱畏謪^(qū)只有一個(gè) Leader 處理寫入acksall需要額外等待 2 個(gè) follower 同步串行等待效應(yīng)明顯。4.2 場(chǎng)景 S2多分區(qū)并行6 分區(qū) × 6 線程指標(biāo)acks0acks1acksall吞吐量 (msg/s)512,000438,000325,000吞吐量 (MB/s)500.0427.7317.4平均延遲 (ms)1.83.26.8P99 延遲 (ms)4.58.115.3P99.9 延遲 (ms)8.215.628.4CPU 使用率 (%)626872分析多分區(qū)后差距縮小。acksall吞吐量達(dá)到acks0的 63%。因?yàn)?6 個(gè)分區(qū)分布在 3 個(gè) Broker 上follower 同步可以并行進(jìn)行。4.3 場(chǎng)景 S3高并發(fā)多分區(qū)6 分區(qū) × 20 線程指標(biāo)acks0acks1acksall吞吐量 (msg/s)1,180,0001,020,000815,000吞吐量 (MB/s)1,152.3996.1795.9平均延遲 (ms)4.26.511.3P99 延遲 (ms)12.618.432.1P99.9 延遲 (ms)25.338.756.8CPU 使用率 (%)788285關(guān)鍵發(fā)現(xiàn)高并發(fā)下acksall突破百萬級(jí)吞吐815K msg/s差距進(jìn)一步縮小到 69%。原因是高并發(fā)下 batch 聚合更充分網(wǎng)絡(luò)往返成本被均攤。4.4 場(chǎng)景 S4弱副本集群ISR 同步慢模擬 follower 落后場(chǎng)景將其中一個(gè) Broker 的replica.fetch.max.bytes限制為 1KB/s。指標(biāo)acks0acks1acksall吞吐量 (msg/s)510,000435,0008,200平均延遲 (ms)1.83.13,650P99 延遲 (ms)4.68.028,000超時(shí)率 (%)00.0142.5重要結(jié)論當(dāng)副本同步異常時(shí)acksall性能斷崖式下降。吞吐量從 325K 跌到 8K延遲從 7ms 飆到 3.6 秒42.5% 的請(qǐng)求超時(shí)。這就是為什么min.insync.replicas和 ISR 監(jiān)控至關(guān)重要。4.5 場(chǎng)景 S5大消息10KB指標(biāo)acks0acks1acksall吞吐量 (msg/s)85,00072,00051,000吞吐量 (MB/s)830.1703.1498.0平均延遲 (ms)5.27.814.5P99 延遲 (ms)15.322.138.6分析大消息場(chǎng)景下acksall的性能比例60%與小消息多分區(qū)場(chǎng)景63%接近。瓶頸轉(zhuǎn)移到網(wǎng)絡(luò)帶寬上。4.6 場(chǎng)景 S630 分鐘持續(xù)寫入穩(wěn)定性6 分區(qū) × 10 線程持續(xù) 30 分鐘每 1 分鐘采樣一次指標(biāo)acks0acks1acksall平均吞吐量 (msg/s)980,000860,000690,000吞吐量標(biāo)準(zhǔn)差12,00018,00045,000最大延遲 (ms)3552180延遲波動(dòng)范圍 (ms)1-352-525-180關(guān)鍵發(fā)現(xiàn)acksall不僅平均延遲更高延遲波動(dòng)也更大。標(biāo)準(zhǔn)差是acks0的 3.75 倍。這是因?yàn)?follower 的 Fetch 請(qǐng)求存在調(diào)度抖動(dòng)偶發(fā)的慢拉取會(huì)拉長(zhǎng)整個(gè) batch 的確認(rèn)時(shí)間。4.7 綜合對(duì)比匯總五、延遲來源拆解acksall比acks0多出的延遲到底花在哪里5.1 延遲拆解模型延遲組件acks0acks1acksall說明隊(duì)列等待1.2 ms1.2 ms1.2 msbatch 聚合時(shí)間網(wǎng)絡(luò)發(fā)送0.5 ms0.5 ms0.5 msProducer → LeaderLeader 寫入0.3 ms0.3 ms0.3 msLeader 日志追加副本同步——4.0 msFollower Fetch 寫入網(wǎng)絡(luò)響應(yīng)0.2 ms0.5 ms0.5 msBroker → Producer總計(jì)2.2 ms2.5 ms6.5 ms副本同步是acksall延遲的主要來源占總延遲的 61.5%。5.2 副本同步為什么需要 4msFollower 同步 Leader 數(shù)據(jù)是通過Fetch 請(qǐng)求輪詢的不是 Leader 主動(dòng)推送// Follower 的 Fetch 線程ReplicaFetcherThread // 默認(rèn)每 500ms 輪詢一次replica.fetch.wait.max.ms while (true) { // 1. 向 Leader 發(fā)送 FetchRequest val fetchRequest buildFetchRequest(partitionMap) val response leaderBroker.fetch(fetchRequest) ? // 2. 將拉取到的數(shù)據(jù)寫入本地日志 response.records.foreach { records localLog.append(records) } ? // 3. 更新本地 High Watermark updateHighWatermark() ? // 4. 等待下一輪默認(rèn) 500ms 間隔 Thread.sleep(fetchWaitMaxMs) // replica.fetch.wait.max.ms }4ms 的分解Follower Fetch 請(qǐng)求到達(dá) Leader 0.3 ms Leader 返回?cái)?shù)據(jù)網(wǎng)絡(luò)傳輸 0.5 ms Follower 寫入本地日志 0.3 ms 等待下一輪 Fetch 輪詢平均等待半個(gè)周期 2.5 ms ← 最大開銷 Leader 檢測(cè) HW 更新并回調(diào) 0.4 ms關(guān)鍵優(yōu)化減少 Fetch 輪詢間隔可以顯著降低延遲# Broker 端配置 replica.fetch.wait.max.ms100 # 從 500ms 降到 100ms默認(rèn) 500 replica.fetch.min.bytes1 # 有數(shù)據(jù)立即拉取默認(rèn) 1 byte replica.fetch.max.bytes1048576 # 單次拉取最大 1MB調(diào)整后重新壓測(cè) S3 場(chǎng)景指標(biāo)默認(rèn)配置優(yōu)化配置改善平均延遲 (ms)11.37.2-36%P99 延遲 (ms)32.121.5-33%吞吐量 (msg/s)815,000845,0003.7%六、生產(chǎn)環(huán)境配置建議6.1 acks 選型決策表業(yè)務(wù)場(chǎng)景推薦 acks理由日志收集/監(jiān)控指標(biāo)0或1容忍少量丟數(shù)據(jù)追求吞吐用戶行為埋點(diǎn)1偶爾丟幾條不影響分析訂單/支付流水a(chǎn)ll不能丟數(shù)據(jù)審計(jì)/合規(guī)日志all必須不丟實(shí)時(shí)推薦特征1容忍秒級(jí)數(shù)據(jù)丟失IoT 設(shè)備告警all告警不能丟CDN 訪問日志0量大且非關(guān)鍵消息通知/IMall消息不能丟6.2 與 acksall 配套的必選配置# Producer 端 acksall retries2147483647 # 無限重試Integer.MAX_VALUE max.in.flight.requests.per.connection5 # 保障冪等性需 ≤5 enable.idempotencetrue # 開啟冪等生產(chǎn)者 request.timeout.ms30000 # 30s 超時(shí) delivery.timeout.ms120000 # 2min 投遞超時(shí) compression.typelz4 # 壓縮減少網(wǎng)絡(luò)開銷 ? # Broker 端 min.insync.replicas2 # 至少 2 個(gè)副本同步成功 default.replication.factor3 # 3 副本 replica.fetch.wait.max.ms100 # 縮短 Fetch 間隔 replica.lag.time.max.ms10000 # 10s 未同步則移出 ISR unclean.leader.election.enablefalse # 禁止非 ISR 副本成為 Leader6.3min.insync.replicas的陷阱Topic: 3 副本, min.insync.replicas2 ? 正常情況ISR[0,1,2], 3 個(gè)副本在線 → acksall 需要 3 個(gè)全部寫入正常返回 ? 一個(gè)副本宕機(jī)ISR[0,1], 2 個(gè)副本在線 → min.insync.replicas2滿足條件acksall 仍可用 ? 兩個(gè)副本宕機(jī)ISR[0], 1 個(gè)副本在線 → min.insync.replicas2不滿足條件 → acksall 直接返回 NotEnoughReplicasException → Topic 不可寫入這就是可靠性 vs 可用性的 trade-off。min.insync.replicas2保證至少 2 個(gè)副本有數(shù)據(jù)但如果 2 個(gè)副本都宕機(jī)Topic 就不可寫了。6.4 冪等生產(chǎn)者與 acksall// 冪等生產(chǎn)者Kafka 0.11 props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); // 冪等生產(chǎn)者自動(dòng)設(shè)置 // acks all // retries Integer.MAX_VALUE // max.in.flight.requests.per.connection ≤ 5 ? // 冪等性通過 PID (Producer ID) Sequence Number 實(shí)現(xiàn) // Producer 發(fā)送PID42, Seq0 → Broker 寫入 [PID42, Seq0] // Producer 超時(shí)重試PID42, Seq0 → Broker 發(fā)現(xiàn)重復(fù)跳過寫入 // Producer 發(fā)送下一條PID42, Seq1 → Broker 寫入 [PID42, Seq1]開啟冪等性后acksall的重試不會(huì)產(chǎn)生重復(fù)消息這是生產(chǎn)環(huán)境的標(biāo)準(zhǔn)配置。七、常見踩坑案例坑 1acksall 但 min.insync.replicas1# 危險(xiǎn)配置 acksall min.insync.replicas1 # ← 問題在這里min.insync.replicas1意味著只要 Leader 自己寫成功就算全部 ISR 同步完。如果 Leader 寫完后立即宕機(jī)follower 還沒拉取到數(shù)據(jù)消息就丟了。正確配置min.insync.replicas至少為 2???2acksall 但 unclean.leader.electiontrueacksall unclean.leader.election.enabletrue # ← 問題在這里unclean.leader.election.enabletrue允許非 ISR 中的副本成為 Leader。這意味著一個(gè)落后很多的副本可能成為新 Leader導(dǎo)致已確認(rèn)的消息丟失。正確配置unclean.leader.election.enablefalse???3高并發(fā)下延遲飆升現(xiàn)象acksall在并發(fā)量上去后P99 延遲從 30ms 飆到 500ms。原因linger.ms0導(dǎo)致每個(gè)小消息都單獨(dú)發(fā)送大量小請(qǐng)求擠滿 Broker 的請(qǐng)求隊(duì)列。修復(fù)linger.ms5 # 等待 5ms 湊 batch batch.size65536 # 64KB batch 大小 compression.typelz4 # 壓縮減少請(qǐng)求體坑 4跨機(jī)房部署延遲放大兩地三中心部署Leader 和 Follower 在不同機(jī)房網(wǎng)絡(luò) RTT 15ms。acksall延遲從 7ms 飆到 22ms增加了 15ms 跨機(jī)房 RTT。解決方案使用Conflent Multi-Region Clusters方案Leader 選舉優(yōu)先在本地機(jī)房或使用MirrorMaker2做跨機(jī)房異步同步本地集群用acks1八、總結(jié)維度acks0acks1acksall吞吐量百萬級(jí)100%86%69%平均延遲最低50%170%P99 延遲最低46%154%延遲穩(wěn)定性好中差數(shù)據(jù)可靠性最低中最高適用場(chǎng)景日志/指標(biāo)埋點(diǎn)/特征訂單/支付核心結(jié)論acksall 不是洪水猛獸在高并發(fā) 多分區(qū)場(chǎng)景下吞吐量只比 acks0 低 31%延遲增加約 7ms但前提是 ISR 健康一旦副本同步異常acksall 性能斷崖式下降325K → 8K msg/s必須配套使用min.insync.replicas2unclean.leader.electionfalseenable.idempotencetrue延遲可優(yōu)化縮短replica.fetch.wait.max.ms可將延遲降低 36%按業(yè)務(wù)選 acks不要全局一刀切按 Topic 可靠性需求分別配置下一篇預(yù)告本專欄下一篇文章將深入Spark 3.5 AQE自適應(yīng)查詢執(zhí)行原理與 10 個(gè)生產(chǎn)調(diào)優(yōu)案例——為什么同一份 Spark SQL開 AQE 后性能差 3 倍如果覺得有幫助點(diǎn)個(gè)贊和收藏關(guān)注專欄「AI大模型大數(shù)據(jù)硬件編程」不錯(cuò)過后續(xù)更新。