最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Kafka源碼系列教程之刪除topic

 更新時(shí)間:2018年08月19日 15:15:07   作者:浪尖  
這篇文章主要給大家介紹了關(guān)于Kafka源碼系列教程之刪除topic的相關(guān)資料,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧

前言

Apache Kafka發(fā)源于LinkedIn,于2011年成為Apache的孵化項(xiàng)目,隨后于2012年成為Apache的主要項(xiàng)目之一。Kafka使用Scala和Java進(jìn)行編寫。Apache Kafka是一個(gè)快速、可擴(kuò)展的、高吞吐、可容錯的分布式發(fā)布訂閱消息系統(tǒng)。Kafka具有高吞吐量、內(nèi)置分區(qū)、支持?jǐn)?shù)據(jù)副本和容錯的特性,適合在大規(guī)模消息處理場景中使用。

本文依然是以kafka0.8.2.2為例講解

一,如何刪除一個(gè)topic

刪除一個(gè)topic有兩個(gè)關(guān)鍵點(diǎn):

1,配置刪除參數(shù)

delete.topic.enable這個(gè)Broker參數(shù)配置為True。

2,執(zhí)行

bin/kafka-topics.sh --zookeeper zk_host:port/chroot --delete --topic my_topic_name

假如不配置刪除參數(shù)為true的話,topic其實(shí)并沒有被清除,只是被標(biāo)記為刪除。此時(shí),估計(jì)一般人的做法是刪除topic在Zookeeper的信息和日志,其實(shí)這個(gè)操作并不會清除kafkaBroker內(nèi)存的topic數(shù)據(jù)。所以,此時(shí)最佳的策略是配置刪除參數(shù)為true然后,重啟kafka。

二,重要的類介紹

1,PartitionStateMachine

該類代表分區(qū)的狀態(tài)機(jī)。決定者分區(qū)的當(dāng)前狀態(tài),和狀態(tài)轉(zhuǎn)移。四種狀態(tài)

  • NonExistentPartition
  • NewPartition
  • OnlinePartition
  • OfflinePartition

2,ReplicaManager

負(fù)責(zé)管理當(dāng)前機(jī)器的所有副本,處理讀寫、刪除等具體動作。

讀寫:寫獲取partition對象,再獲取Replica對象,再獲取Log對象,采用其管理的Segment對象將數(shù)據(jù)寫入、讀出。

3,ReplicaStateMachine

副本的狀態(tài)機(jī)。決定者副本的當(dāng)前狀態(tài)和狀態(tài)之間的轉(zhuǎn)移。一個(gè)副本總共可以處于一下幾種狀態(tài)的一種
NewReplica:Crontroller在分區(qū)重分配的時(shí)候可以創(chuàng)建一個(gè)新的副本。只能接受變?yōu)閒ollower的請求。前狀態(tài)可以是NonExistentReplica

OnlineReplica:新啟動的分區(qū),能接受變?yōu)閘eader或者follower請求。前狀態(tài)可以是NewReplica, OnlineReplica or OfflineReplica

OfflineReplica:死亡的副本處于這種狀態(tài)。前狀態(tài)可以是NewReplica, OnlineReplica

ReplicaDeletionStarted:分本刪除開始的時(shí)候處于這種狀態(tài),前狀態(tài)是OfflineReplica

ReplicaDeletionSuccessful:副本刪除成功。前狀態(tài)是ReplicaDeletionStarted

ReplicaDeletionIneligible:刪除失敗的時(shí)候處于這種狀態(tài)。前狀態(tài)是ReplicaDeletionStarted

NonExistentReplica:副本成功刪除之后處于這種狀態(tài),前狀態(tài)是ReplicaDeletionSuccessful

4,TopicDeletionManager

該類管理著topic刪除的狀態(tài)機(jī)

1),TopicCommand通過創(chuàng)建/admin/delete_topics/<topic>,來發(fā)布topic刪除命令。

2),Controller監(jiān)聽/admin/delete_topic子節(jié)點(diǎn)變動,開始分別刪除topic

3),Controller有個(gè)后臺線程負(fù)責(zé)刪除Topic

三,源碼徹底解析topic的刪除過程

此處會分四個(gè)部分:

A),客戶端執(zhí)行刪除命令作用

B),不配置delete.topic.enable整個(gè)流水的源碼

C),配置了delete.topic.enable整個(gè)流水的源碼

D),手動刪除zk上topic信息和磁盤數(shù)據(jù)

1,客戶端執(zhí)行刪除命令

bin/kafka-topics.sh --zookeeper zk_host:port/chroot --delete --topic my_topic_name

進(jìn)入kafka-topics.sh我們會看到

exec $(dirname $0)/kafka-run-class.sh kafka.admin.TopicCommand $@

進(jìn)入TopicCommand里面,main方法里面

else if(opts.options.has(opts.deleteOpt))
 deleteTopic(zkClient, opts)

實(shí)際內(nèi)容是

val topics = getTopics(zkClient, opts)
if (topics.length == 0) {
 println("Topic %s does not exist".format(opts.options.valueOf(opts.topicOpt)))
}
topics.foreach { topic =>
 try {
 ZkUtils.createPersistentPath(zkClient, ZkUtils.getDeleteTopicPath(topic))

在"/admin/delete_topics"目錄下創(chuàng)建了一個(gè)topicName的節(jié)點(diǎn)。

2,假如不配置delete.topic.enable整個(gè)流水是

總共有兩處listener會響應(yīng):

A),TopicChangeListener

B),DeleteTopicsListener

使用topic的刪除命令刪除一個(gè)topic的話,指揮觸發(fā)DeleteTopicListener。

var topicsToBeDeleted = {
 import JavaConversions._
 (children: Buffer[String]).toSet
}
val nonExistentTopics = topicsToBeDeleted.filter(t => !controllerContext.allTopics.contains(t))
topicsToBeDeleted --= nonExistentTopics
if(topicsToBeDeleted.size > 0) {
 info("Starting topic deletion for topics " + topicsToBeDeleted.mkString(","))
 // mark topic ineligible for deletion if other state changes are in progress
 topicsToBeDeleted.foreach { topic =>
 val preferredReplicaElectionInProgress =
  controllerContext.partitionsUndergoingPreferredReplicaElection.map(_.topic).contains(topic)
 val partitionReassignmentInProgress =
  controllerContext.partitionsBeingReassigned.keySet.map(_.topic).contains(topic)
 if(preferredReplicaElectionInProgress || partitionReassignmentInProgress)
  controller.deleteTopicManager.markTopicIneligibleForDeletion(Set(topic))
 }
 // add topic to deletion list 
 controller.deleteTopicManager.enqueueTopicsForDeletion(topicsToBeDeleted)
}

由于都會判斷delete.topic.enable是否為true,假如不為true就不會執(zhí)行,為true就進(jìn)入執(zhí)行

controller.deleteTopicManager.markTopicIneligibleForDeletion(Set(topic))
controller.deleteTopicManager.enqueueTopicsForDeletion(topicsToBeDeleted)

3,delete.topic.enable配置為true

此處與步驟2的區(qū)別,就是那兩個(gè)處理函數(shù)。

controller.deleteTopicManager.markTopicIneligibleForDeletion(Set(topic))
controller.deleteTopicManager.enqueueTopicsForDeletion(topicsToBeDeleted)

markTopicIneligibleForDeletion函數(shù)的處理為

if(isDeleteTopicEnabled) {
 val newTopicsToHaltDeletion = topicsToBeDeleted & topics
 topicsIneligibleForDeletion ++= newTopicsToHaltDeletion
 if(newTopicsToHaltDeletion.size > 0)
 info("Halted deletion of topics %s".format(newTopicsToHaltDeletion.mkString(",")))
}

主要是停止刪除topic,假如存儲以下三種情況

* Halt delete topic if -
* 1. replicas being down
* 2. partition reassignment in progress for some partitions of the topic
* 3. preferred replica election in progress for some partitions of the topic

enqueueTopicsForDeletion主要作用是更新刪除topic的集合,并激活TopicDeleteThread

def enqueueTopicsForDeletion(topics: Set[String]) {
 if(isDeleteTopicEnabled) {
 topicsToBeDeleted ++= topics
 partitionsToBeDeleted ++= topics.flatMap(controllerContext.partitionsForTopic)
 resumeTopicDeletionThread()
 }
}

在刪除線程DeleteTopicsThread的doWork方法中

topicsQueuedForDeletion.foreach { topic =>
// if all replicas are marked as deleted successfully, then topic deletion is done
 if(controller.replicaStateMachine.areAllReplicasForTopicDeleted(topic)) {
 // clear up all state for this topic from controller cache and zookeeper
 completeDeleteTopic(topic)
 info("Deletion of topic %s successfully completed".format(topic))
 }

進(jìn)入completeDeleteTopic方法中

// deregister partition change listener on the deleted topic. This is to prevent the partition change listener
// firing before the new topic listener when a deleted topic gets auto created
partitionStateMachine.deregisterPartitionChangeListener(topic)
val replicasForDeletedTopic = controller.replicaStateMachine.replicasInState(topic, ReplicaDeletionSuccessful)
// controller will remove this replica from the state machine as well as its partition assignment cache
replicaStateMachine.handleStateChanges(replicasForDeletedTopic, NonExistentReplica)
val partitionsForDeletedTopic = controllerContext.partitionsForTopic(topic)
// move respective partition to OfflinePartition and NonExistentPartition state
partitionStateMachine.handleStateChanges(partitionsForDeletedTopic, OfflinePartition)
partitionStateMachine.handleStateChanges(partitionsForDeletedTopic, NonExistentPartition)
topicsToBeDeleted -= topic
partitionsToBeDeleted.retain(_.topic != topic)
controllerContext.zkClient.deleteRecursive(ZkUtils.getTopicPath(topic))
controllerContext.zkClient.deleteRecursive(ZkUtils.getTopicConfigPath(topic))
controllerContext.zkClient.delete(ZkUtils.getDeleteTopicPath(topic))
controllerContext.removeTopic(topic)

主要作用是解除掉監(jiān)控分區(qū)變動的listener,刪除Zookeeper具體節(jié)點(diǎn)信息,刪除磁盤數(shù)據(jù),更新內(nèi)存數(shù)據(jù)結(jié)構(gòu),比如從副本狀態(tài)機(jī)里面移除分區(qū)的具體信息。

其實(shí),最終要的是我們的副本磁盤數(shù)據(jù)是如何刪除的。我們重點(diǎn)介紹這個(gè)部分。

首次清除的話,在刪除線程DeleteTopicsThread的doWork方法中

{
 // if you come here, then no replica is in TopicDeletionStarted and all replicas are not in
 // TopicDeletionSuccessful. That means, that either given topic haven't initiated deletion
 // or there is at least one failed replica (which means topic deletion should be retried).
 if(controller.replicaStateMachine.isAnyReplicaInState(topic, ReplicaDeletionIneligible)) {
 // mark topic for deletion retry
 markTopicForDeletionRetry(topic)
 }

進(jìn)入markTopicForDeletionRetry

val failedReplicas = controller.replicaStateMachine.replicasInState(topic, ReplicaDeletionIneligible)
info("Retrying delete topic for topic %s since replicas %s were not successfully deleted"
 .format(topic, failedReplicas.mkString(",")))
controller.replicaStateMachine.handleStateChanges(failedReplicas, OfflineReplica)

在ReplicaStateMachine的handleStateChanges方法中,調(diào)用了handleStateChange,處理OfflineReplica

// send stop replica command to the replica so that it stops fetching from the leader
brokerRequestBatch.addStopReplicaRequestForBrokers(List(replicaId), topic, partition, deletePartition = false)

接著在handleStateChanges中

brokerRequestBatch.sendRequestsToBrokers(controller.epoch, controllerContext.correlationId.getAndIncrement)

給副本數(shù)據(jù)存儲節(jié)點(diǎn)發(fā)送StopReplicaKey副本指令,并開始刪除數(shù)據(jù)

stopReplicaRequestMap foreach { case(broker, replicaInfoList) =>
 val stopReplicaWithDelete = replicaInfoList.filter(p => p.deletePartition == true).map(i => i.replica).toSet
 val stopReplicaWithoutDelete = replicaInfoList.filter(p => p.deletePartition == false).map(i => i.replica).toSet
 debug("The stop replica request (delete = true) sent to broker %d is %s"
 .format(broker, stopReplicaWithDelete.mkString(",")))
 debug("The stop replica request (delete = false) sent to broker %d is %s"
 .format(broker, stopReplicaWithoutDelete.mkString(",")))
 replicaInfoList.foreach { r =>
 val stopReplicaRequest = new StopReplicaRequest(r.deletePartition,
  Set(TopicAndPartition(r.replica.topic, r.replica.partition)), controllerId, controllerEpoch, correlationId)
 controller.sendRequest(broker, stopReplicaRequest, r.callback)
 }
}
stopReplicaRequestMap.clear()

Broker的KafkaApis的Handle方法在接受到指令后

case RequestKeys.StopReplicaKey => handleStopReplicaRequest(request)
val (response, error) = replicaManager.stopReplicas(stopReplicaRequest)

接著是在stopReplicas方法中

{
 controllerEpoch = stopReplicaRequest.controllerEpoch
 // First stop fetchers for all partitions, then stop the corresponding replicas
 replicaFetcherManager.removeFetcherForPartitions(stopReplicaRequest.partitions.map(r => TopicAndPartition(r.topic, r.partition)))
 for(topicAndPartition <- stopReplicaRequest.partitions){
 val errorCode = stopReplica(topicAndPartition.topic, topicAndPartition.partition, stopReplicaRequest.deletePartitions)
 responseMap.put(topicAndPartition, errorCode)
 }
 (responseMap, ErrorMapping.NoError)
}

進(jìn)一步進(jìn)入stopReplica方法,正式進(jìn)入日志刪除

getPartition(topic, partitionId) match {
 case Some(partition) =>
 if(deletePartition) {
  val removedPartition = allPartitions.remove((topic, partitionId))
  if (removedPartition != null)
  removedPartition.delete() // this will delete the local log
 }

以上就是kafka的整個(gè)日志刪除流水。

4,手動刪除zk上topic信息和磁盤數(shù)據(jù)

TopicChangeListener會監(jiān)聽處理,但是處理很簡單,只是更新了

val deletedTopics = controllerContext.allTopics -- currentChildren
controllerContext.allTopics = currentChildren

val addedPartitionReplicaAssignment = ZkUtils.getReplicaAssignmentForTopics(zkClient, newTopics.toSeq)
controllerContext.partitionReplicaAssignment = controllerContext.partitionReplicaAssignment.filter(p =>

四,總結(jié)

Kafka的topic的刪除過程,實(shí)際上就是基于Zookeeper做了一個(gè)訂閱發(fā)布系統(tǒng)。Zookeeper的客戶端創(chuàng)建一個(gè)節(jié)點(diǎn)/admin/delete_topics/<topic>,由kafka Controller監(jiān)聽到事件之后正式觸發(fā)topic的刪除:解除Partition變更監(jiān)聽的listener,清除內(nèi)存數(shù)據(jù)結(jié)構(gòu),刪除副本數(shù)據(jù),刪除topic的相關(guān)Zookeeper節(jié)點(diǎn)。

delete.topic.enable配置該參數(shù)為false的情況下執(zhí)行了topic的刪除命令,實(shí)際上未做任何動作。我們此時(shí)要徹底刪除topic建議修改該參數(shù)為true,重啟kafka,這樣topic信息會被徹底刪除,已經(jīng)測試。

一般流行的做法是手動刪除Zookeeper的topic相關(guān)信息及磁盤數(shù)據(jù)但是這樣的話會造成部分內(nèi)存數(shù)據(jù)未清除。至于是否會有隱患,未測試。

好了,以上就是這篇文章的全部內(nèi)容了,希望本文的內(nèi)容對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,如果有疑問大家可以留言交流,謝謝大家對腳本之家的支持。

相關(guān)文章

  • Java檢測網(wǎng)絡(luò)是否正常通訊

    Java檢測網(wǎng)絡(luò)是否正常通訊

    在網(wǎng)絡(luò)應(yīng)用程序中,檢測IP地址和端口是否通常是必要的,本文主要介紹了Java檢測網(wǎng)絡(luò)是否正常通訊,具有一定的參考價(jià)值,感興趣的可以了解一下
    2023-11-11
  • java基本教程之Thread中start()和run()的區(qū)別 java多線程教程

    java基本教程之Thread中start()和run()的區(qū)別 java多線程教程

    這篇文章主要介紹了Thread中start()和run()的區(qū)別,Thread類包含start()和run()方法,它們的區(qū)別是什么?下面將對此作出解答
    2014-01-01
  • 處理java異步事件的阻塞和非阻塞方法分析

    處理java異步事件的阻塞和非阻塞方法分析

    這篇文章主要介紹了處理java異步事件的阻塞和非阻塞方法分析,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,阻塞與非阻塞關(guān)注的是交互雙方是否可以彈性工作。,需要的朋友可以參考下
    2019-06-06
  • 詳解Java中-classpath和路徑的使用

    詳解Java中-classpath和路徑的使用

    本篇文章主要介紹了Java中-classpath和路徑的使用,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過來看看吧
    2017-04-04
  • Java如何不解壓讀取.zip的文件內(nèi)容

    Java如何不解壓讀取.zip的文件內(nèi)容

    這篇文章主要給大家介紹了關(guān)于Java如何不解壓讀取.zip的文件內(nèi)容的相關(guān)資料,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2021-03-03
  • SpringMVC框架搭建idea2021.3.2操作數(shù)據(jù)庫的示例詳解

    SpringMVC框架搭建idea2021.3.2操作數(shù)據(jù)庫的示例詳解

    這篇文章主要介紹了SpringMVC框架搭建idea2021.3.2操作數(shù)據(jù)庫,本文通過示例代碼給大家介紹的非常詳細(xì),對大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2022-04-04
  • spring單元測試之@RunWith的使用詳解

    spring單元測試之@RunWith的使用詳解

    這篇文章主要介紹了spring單元測試之@RunWith的使用詳解,@RunWith 就是一個(gè)運(yùn)行器,@RunWith(JUnit4.class) 就是指用JUnit4來運(yùn)行,
    @RunWith(SpringJUnit4ClassRunner.class),讓測試運(yùn)行于Spring測試環(huán)境,需要的朋友可以參考下
    2023-12-12
  • java  線程詳解及線程與進(jìn)程的區(qū)別

    java 線程詳解及線程與進(jìn)程的區(qū)別

    這篇文章主要介紹了java 線程詳解及線程與進(jìn)程的區(qū)別的相關(guān)資料,網(wǎng)上關(guān)于java 線程的資料很多,對于進(jìn)程的資料很是,這里就整理下,需要的朋友可以參考下
    2017-01-01
  • java如何判斷一個(gè)數(shù)是否是素?cái)?shù)(質(zhì)數(shù))

    java如何判斷一個(gè)數(shù)是否是素?cái)?shù)(質(zhì)數(shù))

    這篇文章主要介紹了java如何判斷一個(gè)數(shù)是否是素?cái)?shù)(質(zhì)數(shù)),具有很好的參考價(jià)值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-09-09
  • 不使用他人jar包情況下優(yōu)雅的進(jìn)行dubbo調(diào)用詳解

    不使用他人jar包情況下優(yōu)雅的進(jìn)行dubbo調(diào)用詳解

    這篇文章主要為大家介紹了不使用他人jar包情況下優(yōu)雅的進(jìn)行dubbo調(diào)用詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2022-09-09

最新評論

陇南市| 台湾省| 淳化县| 马公市| 原平市| 宝山区| 迁西县| 长治市| 武威市| 大厂| 永兴县| 桑日县| 桃源县| 桦南县| 石楼县| 左贡县| 鹤岗市| 奇台县| 横山县| 中江县| 城市| 内乡县| 石棉县| 西盟| 克拉玛依市| 大方县| 怀来县| 清流县| 台南县| 敦化市| 崇州市| 晋城| 牡丹江市| 景洪市| 中西区| 嘉善县| 南通市| 霞浦县| 扬中市| 武夷山市| 高邮市|