kafka管控平臺推薦使用 滴滴開源 的 Kafka運維管控平臺(戳我呀) 更符合國人的操作習慣 、更強大的管控能力 、更高效的問題定位能力 、更便捷的集群運維能力 、更專業的資源治理 、更友好的運維生態 、
BliBli視頻: 石臻臻的雜貨鋪
kafka的動態配置
文章目錄
- 原始碼分析
- 1. Broker啟動加載動態配置
- 1.1 啟動加載動態配置總流程
- 1. 2 加載Topic動態配置
- 1.3 加載Broker動態配置
- 2. 查詢動態配置 流程 `--describe`
- 3. 新增/修改/洗掉/動態配置 的流程
- Topic配置
- 其他的型別都一樣
- 4. Broker監聽/config/changes的變更
- 原始碼總結
- Q&A
- 如果我想在我的專案中獲取kafka的所有配置該怎么辦?
- 是否可以直接在zk中寫入動態配置?
- 為什么不直接監聽 `/config/`下面的配置?

今天這篇文章,給大家分享一下最近看kafka中的動態配置,不需要重啟Broker,即時生效的配置 歡迎留言一起探討!
kafka中的配置
- Broker靜態配置
.properties檔案- ZK中的動態配置
全域 default配置- ZK中動態配置 指定配置
優先級從底到高
原始碼分析
1. Broker啟動加載動態配置
KafkaServer.startup
1.1 啟動加載動態配置總流程
1. 動態配置初始化
config.dynamicConfig.initialize(zkClient)
- 構造當前組態檔
currentConfig, 然后從zk中獲取節點/config/brokers/<default>資訊,然后更新配置updateDefaultConfig; (動態默認配置覆寫靜態配置) - 從節點
/config/brokers/{當前BrokerId}獲取配置, 如果配置中有ConfigType=PASSWORD的配置(例如ssl.keystore.password)存在,接著判斷 是否存在password.encoder.old.secret配置,(這個配置是用來加解密ConfigType=PASSWORD的舊的秘鑰),嘗試用舊秘鑰解密秘鑰; 然后將這些配置重新加密回寫入/config/brokers/{當前BrokerId}; 然后回傳配置 (這里主要是動態配置里面有密碼型別配置的時候需要做一次解密加密處理) - 將上面得到的配置(password型別修改之后) 更新記憶體總的配置;優先級 靜態配置<動態默認配置<指定動態配置
2. 注冊可變更配置監聽器
如果有對應的配置變更了,那么相應的監聽器就會收到通知去修改自己相應的配置;
config.dynamicConfig.addReconfigurables(this)
DynamicBrokerConfig.addReconfigurables
// .........
def addReconfigurables(kafkaServer: KafkaServer): Unit = {
kafkaServer.authorizer match {
case Some(authz: Reconfigurable) => addReconfigurable(authz)
case _ =>
}
addReconfigurable(new DynamicMetricsReporters(kafkaConfig.brokerId, kafkaServer))
addReconfigurable(new DynamicClientQuotaCallback(kafkaConfig.brokerId, kafkaServer))
addBrokerReconfigurable(new DynamicThreadPool(kafkaServer))
if (kafkaServer.logManager.cleaner != null)
addBrokerReconfigurable(kafkaServer.logManager.cleaner)
addBrokerReconfigurable(new DynamicLogConfig(kafkaServer.logManager, kafkaServer))
addBrokerReconfigurable(new DynamicListenerConfig(kafkaServer))
addBrokerReconfigurable(kafkaServer.socketServer)
}
3. 動態配置啟動監聽
// Create the config manager. start listening to notifications
dynamicConfigManager = new DynamicConfigManager(zkClient, dynamicConfigHandlers)
dynamicConfigManager.startup()
- 注冊節點處理器
change-notification-/config/changes= stateChangeHandler - 注冊節點處理器
/config/changes= zNodeChildChangeHandler - 獲取
/config/changes所有子節點看看有哪些變更 - 遍歷所有節點并截取節點的編號, 判斷一下是不是大于上一次執行過變更的節點ID
lastExecutedChange(啟動的時候是-1) - 上個條件滿足的話,則執行通知操作;不同entity執行的操作不一樣,具體請看下面每個型別
- 更新
lastExecutedChange - 清除過期的通知節點, 默認過期時間
15 * 60 * 1000(15分鐘)就是洗掉/config/changes /下面的過期節點
1. 2 加載Topic動態配置
TopicConfigHandler.processConfigChanges

- 獲取節點的
data資料, 如果獲取到了則執行通知流程notificationHandler.processNotification(d),處理器是ConfigChangedNotificationHandler; 它先決議節點的json資料,根據版本資訊不同呼叫不同的處理方法; 下面是version=2的處理方式; - 根據json資料可以得到
entityType和entityName; 那么久可以去對應的zk資料里面getData獲取資料; 并且將獲取到的資料Decode成Properties物件entityConfig; - 將key為下圖中的屬性 隱藏掉; 替換成value: [hidden]

- 呼叫EntityHandler; 這里是
TopicConfigHandler.processConfigChanges來進行處理,方法里面再看看流程-> - 從動態配置
entityConfig里面獲取message.format.version配置訊息格式版本號; 如果當前Broker的版本inter.broker.protocol.version小于message.format.version配置; 則將message.format.version配置 排除掉 - 呼叫
TopicConfigHandler.updateLogConfig來更新指定Topic的所有TopicPartition的配置,其實是將TP正在加載或初始化的狀態標記為沒有完成初始化,這將會在后續程序中促成TP重新加載并初始化 - 將動態配置和并覆寫Server的默認配置為新的 newConfig, 然后根據Topic獲取對應的Logs物件; 遍歷Logs去更新newConfig;并嘗試執行
initializeLeaderEpochCache; (需要注意的是:這里的動態配置不是支持所有的配置引數,請看【kafka運維】Kafka全網最全最詳細運維命令合集(精品強烈建議收藏!!!)的附件部分) - 當然特殊配置如
leader.replication.throttled.replicas,follower.replication.throttled.replicas這兩個限流相關;決議配置之后,然后通過quotaManager.markThrottled/quotaManager.removeThrottle更新/移除對應的限流磁區集合 - 如果動態配置了
unclean.leader.election.enable=true(允許非同步副本選主 );那么就會執行TopicUncleanLeaderElectionEnable方法來讓它改變選舉策略(前提是當前Broker是Controller角色)
1.3 加載Broker動態配置
BrokerConfigHandler.processConfigChanges
假設我們配置了默認配置; zk里面的節點是<default>
sh bin/kafka-configs.sh --bootstrap-server xxxxx:9090 --alter --entity-type brokers
--entity-default--add-config log.segment.bytes=88888888

- 從zk節點
/config/changes里面獲取變更節點的json資料.然后去對應的 /config/{entityType}/{entituName}獲取對應的資料 - 如果是
<default>節點,說明有配置動態默認配置; 則按照 靜態配置<動態默認配置<動態指定配置 的順序重新加載覆寫一下; 如果 新舊配置有變更(有可能執行了一次命令但是引數并沒有變化的情況,修改了個寂寞)的情況下 才會做更新的; 并且 通知到所有的BrokerReconfigurable; 這個就是上面啟動時候 1.1 啟動加載動態配置總流程的第2步驟 (注冊可變更配置監聽器) 注冊的; - 如果是指定BrokerId, 則除了上面2重新加載覆寫之外, 相關限流 配置
leader.replication.throttled.rate、follower.replication.throttled.rate、replica.alter.log.dirs.io.max.bytes.per.second都會被更新一下quotaManagers.leader/leader/alterLogDirs.updateQuota;如果這些配置沒有配置的話,則用Long.MaxValue(相當于是不限流)來更新
2. 查詢動態配置 流程 --describe
-
簡單檢驗
-
根據型別查詢
entities; type是topics就獲取所有topic; type是broker|broker-loggers則查詢所有Broker節點 -
遍歷
entities獲取配置 ;做些簡單校驗;然后想Broker發起describeConfigs請求; 節點策略是LeastLoadedNodeProvider
節點呼叫方法KafkaApis.handleDescribeConfigsRequest- 未經授權配置不查詢
- 經過授權的配置開始查詢 ;
- 當查詢的是
topics時, 去zk節點/confgi/型別/型別名,獲取到動態配置資料之后, 然后將其覆寫本地跟Log相關的靜態配置, 完事之后組裝一下回傳;(1.資料為空過濾2.敏感資料設定value=null; ConfigType=PASSWORD和不知道型別是啥的都是敏感資料 3. 組裝所有的同義配置(靜態默認配置、本地靜態、默認動態配置、指定動態配置、等等多個配置))
回傳的資料型別如下:


-
如果有
broker|broker-loggers節點, 則在 獲取到資料之后 然后指定nodeId節點發起describeBrokerConfigs請求- 如果查詢的是
brokers

- 如果查詢的是
broker-loggers

- 如果查詢的是
3. 新增/修改/洗掉/動態配置 的流程
1. 發起請求
- 查詢當前的型別配置; 這里的查詢 跟上面的
--describe流程是一樣的 - 相關校驗;如果有
delete-config配置, 需要校驗一下當前配置有沒有;如果沒有拋出例外; - 計算出需要變更的配置之后, 發起請求
incrementalAlterConfigs;如果請求型別是brokers/broker-loggers則發起請求的接收方是 指定的Broker 節點; 否則就是LeastLoadedNodeProvider(當前負載最少的節點)
2. incrementalAlterConfigs 增量修改配置
KafkaApis.handleIncrementalAlterConfigsRequest
- 通過請求引數決議 配置 configs
- 過濾一下未授權的配置
- 如果配置中有重復的項則拋出例外
Topic配置
- 獲取節點 /config/topics/{topicName} 中的配置資料;
- 然后根據請求引數的屬性 ,組裝好變更后的配置是什么樣的
configs; - 簡單校驗一下, 并且支持自定義校驗,如果有
alter.config.policy.class.name=配置(默認null)的話,則會實體化指定的類(需要繼承AlterConfigPolicy類);并呼叫他的validate方法來校驗; - 呼叫寫入zk配置的介面, 將動態配置重新寫入(SetDataRequest)到介面
/config/topics/{topicName}中; - 創建并寫入配置變更記錄順序節點
/config/changes/config_change_序列號中; 這個節點主要是讓Broker們來監聽這個節點的來了解到哪個配置有變更的;
其他的型別都一樣
省略
4. Broker監聽/config/changes的變更
在 1. Broker啟動加載動態配置 中我們了解到有對節點/config/change注冊一個子節點變更的監聽處理器

那么對動態配置做出修改之后, 這個節點就會新增一條資料,那么所有的Broker都會收到這個通知;

所以我們就要來看一看收到通知之后又做了哪些事情
這個流程是又回到了上面的 1. 2 加載Topics/Brokers動態配置 的流程中了;
原始碼總結
原理部分講解比較詳細的可以看 : Kafka動態配置實作原理決議 - 李志濤 - 博客園


Q&A
如果我想在我的專案中獲取kafka的所有配置該怎么辦?
- 啟動的時候加載一次所有Broker的配置
- 監聽節點
/config/change節點的變化
是否可以直接在zk中寫入動態配置?
不可以,因為Broker是監聽
/config/changes/里面的Broker節點,來實時得知有資料變更;
為什么不直接監聽 /config/下面的配置?
沒有必要,這樣監聽的資料資料太多了,而且 你不知道具體是改了哪個配置,所以每次都要全部更新一遍,無緣無故的加重負擔了, 用
/config/change節點來得知哪個型別的資料變更, 只變更這個相關資料就可以了
轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/296533.html
標籤:其他
上一篇:RabbitMq (一)理論篇部分 MQ作用是什么 MQ的優缺點 RabbitMQ的基礎架構 RabbitMQ 五種常用作業模式 RabbitMQ訊息確認機制
