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

Spring Boot集群管理工具KafkaAdminClient使用方法解析

 更新時間:2020年02月25日 10:25:56   作者:---WeiGeH  
這篇文章主要介紹了Spring Boot集群管理工具KafkaAdminClient使用方法解析,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下

原理介紹

在Kafka官網(wǎng)中這么描述AdminClient:The AdminClient API supports managing and inspecting topics, brokers, acls, and other Kafka objects. 具體的KafkaAdminClient包含了一下幾種功能(以Kafka1.0.0版本為準(zhǔn)):

  • 創(chuàng)建Topic:createTopics(Collection<NewTopic> newTopics)
  • 刪除Topic:deleteTopics(Collection<String> topics)
  • 羅列所有Topic:listTopics()
  • 查詢Topic:describeTopics(Collection<String> topicNames)
  • 查詢集群信息:describeCluster()
  • 查詢ACL信息:describeAcls(AclBindingFilter filter)
  • 創(chuàng)建ACL信息:createAcls(Collection<AclBinding> acls)
  • 刪除ACL信息:deleteAcls(Collection<AclBindingFilter> filters)
  • 查詢配置信息:describeConfigs(Collection<ConfigResource> resources)
  • 修改配置信息:alterConfigs(Map<ConfigResource, Config> configs)
  • 修改副本的日志目錄:alterReplicaLogDirs(Map<TopicPartitionReplica, String> replicaAssignment)
  • 查詢節(jié)點(diǎn)的日志目錄信息:describeLogDirs(Collection<Integer> brokers)
  • 查詢副本的日志目錄信息:describeReplicaLogDirs(Collection<TopicPartitionReplica> replicas)
  • 增加分區(qū):createPartitions(Map<String, NewPartitions> newPartitions)

其內(nèi)部原理是使用Kafka自定義的一套二進(jìn)制協(xié)議來實(shí)現(xiàn),詳細(xì)可以參見Kafka協(xié)議。主要實(shí)現(xiàn)步驟:

客戶端根據(jù)方法的調(diào)用創(chuàng)建相應(yīng)的協(xié)議請求,比如創(chuàng)建Topic的createTopics方法,其內(nèi)部就是發(fā)送CreateTopicRequest請求。
客戶端發(fā)送請求至Kafka Broker。

Kafka Broker處理相應(yīng)的請求并回執(zhí),比如與CreateTopicRequest對應(yīng)的是CreateTopicResponse。
客戶端接收相應(yīng)的回執(zhí)并進(jìn)行解析處理。

和協(xié)議有關(guān)的請求和回執(zhí)的類基本都在org.apache.kafka.common.requests包中,AbstractRequest和AbstractResponse是這些請求和回執(zhí)類的兩個基本父類。

代碼如下

@Component
public class KafkaConfig{

   // 配置Kafka
  public Properties getProps(){
    Properties props = new Properties();
    props.put("bootstrap.servers", "localhost:9092");
/*    props.put("retries", 2); // 重試次數(shù)
    props.put("batch.size", 16384); // 批量發(fā)送大小
    props.put("buffer.memory", 33554432); // 緩存大小,根據(jù)本機(jī)內(nèi)存大小配置
    props.put("linger.ms", 1000); // 發(fā)送頻率,滿足任務(wù)一個條件發(fā)送*/
    props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
    props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
    return props;
  }

}
@RestController
public class KafkaTopicManager {

  @Autowired
  private KafkaConfig kafkaConfig;

  @GetMapping("createTopic")
  public void createTopic(){
    AdminClient adminClient = KafkaAdminClient.create(kafkaConfig.getProps());

    NewTopic newTopic = new NewTopic("test1",4, (short) 1);
    Collection<NewTopic> newTopicList = new ArrayList<>();
    newTopicList.add(newTopic);
    adminClient.createTopics(newTopicList);

    adminClient.close();
  }
  @GetMapping("deleteTopic")
  public void deleteTopic(){
    AdminClient adminClient = KafkaAdminClient.create(kafkaConfig.getProps());
    adminClient.deleteTopics(Arrays.asList("test1"));
    adminClient.close();
  }
  @GetMapping("listAllTopic")
  public void listAllTopic(){
    AdminClient adminClient = KafkaAdminClient.create(kafkaConfig.getProps());
    ListTopicsResult result = adminClient.listTopics();
    KafkaFuture<Set<String>> names = result.names();
    try {
      names.get().forEach((k)->{
        System.out.println(k);
      });
    } catch (InterruptedException | ExecutionException e) {
      e.printStackTrace();
    }
    adminClient.close();
  }
  @GetMapping("getTopic")
  public void getTopic(){
    AdminClient adminClient = KafkaAdminClient.create(kafkaConfig.getProps());

    DescribeTopicsResult describeTopics = adminClient.describeTopics(Arrays.asList("syn-test"));

    Collection<KafkaFuture<TopicDescription>> values = describeTopics.values().values();

    if(values.isEmpty()){
      System.out.println("找不到描述信息");
    }else{
      for (KafkaFuture<TopicDescription> value : values) {
        System.out.println(value);
      }
    }
    adminClient.close();
  }
}

以上就是本文的全部內(nèi)容,希望對大家的學(xué)習(xí)有所幫助,也希望大家多多支持腳本之家。

相關(guān)文章

  • java使用RSA加密方式實(shí)現(xiàn)數(shù)據(jù)加密解密的代碼

    java使用RSA加密方式實(shí)現(xiàn)數(shù)據(jù)加密解密的代碼

    這篇文章給大家分享java使用RSA加密方式實(shí)現(xiàn)數(shù)據(jù)加密解密,通過實(shí)例代碼文字相結(jié)合給大家介紹的非常詳細(xì),具有一定的參考借鑒價值,需要的朋友參考下
    2019-11-11
  • 通過反射實(shí)現(xiàn)Java下的委托機(jī)制代碼詳解

    通過反射實(shí)現(xiàn)Java下的委托機(jī)制代碼詳解

    這篇文章主要介紹了通過反射實(shí)現(xiàn)Java下的委托機(jī)制代碼詳解,具有一定借鑒價值,需要的朋友可以參考下。
    2017-12-12
  • MyBatis生成UUID的實(shí)現(xiàn)

    MyBatis生成UUID的實(shí)現(xiàn)

    這篇文章主要介紹了MyBatis生成UUID的實(shí)現(xiàn),文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2020-12-12
  • JAVA使用ElasticSearch查詢in和not in的實(shí)現(xiàn)方式

    JAVA使用ElasticSearch查詢in和not in的實(shí)現(xiàn)方式

    今天小編就為大家分享一篇關(guān)于JAVA使用Elasticsearch查詢in和not in的實(shí)現(xiàn)方式,小編覺得內(nèi)容挺不錯的,現(xiàn)在分享給大家,具有很好的參考價值,需要的朋友一起跟隨小編來看看吧
    2018-12-12
  • 一篇文章教你如何用多種迭代寫法實(shí)現(xiàn)二叉樹遍歷

    一篇文章教你如何用多種迭代寫法實(shí)現(xiàn)二叉樹遍歷

    這篇文章主要介紹了C語言實(shí)現(xiàn)二叉樹遍歷的迭代算法,包括二叉樹的中序遍歷、先序遍歷及后序遍歷等,是非常經(jīng)典的算法,需要的朋友可以參考下
    2021-08-08
  • MybatisPlus查詢條件為空字符串或null問題及解決

    MybatisPlus查詢條件為空字符串或null問題及解決

    這篇文章主要介紹了MybatisPlus查詢條件為空字符串或null問題及解決方案,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2022-06-06
  • java判斷是否空最簡單的方法

    java判斷是否空最簡單的方法

    在本篇文章里小編給大家整理的一篇關(guān)于java判斷是否空最簡單的方法,有興趣的讀者們可以參考下。
    2019-12-12
  • ShardingSphere數(shù)據(jù)分片算法及測試實(shí)戰(zhàn)

    ShardingSphere數(shù)據(jù)分片算法及測試實(shí)戰(zhàn)

    這篇文章主要為大家介紹了ShardingSphere數(shù)據(jù)分片算法及測試實(shí)戰(zhàn)示例詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2023-03-03
  • SpringCloud Hystrix-Dashboard儀表盤的實(shí)現(xiàn)

    SpringCloud Hystrix-Dashboard儀表盤的實(shí)現(xiàn)

    這篇文章主要介紹了SpringCloud Hystrix-Dashboard儀表盤的實(shí)現(xiàn),文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2019-08-08
  • 深度理解SpringMVC中的HandlerMapping

    深度理解SpringMVC中的HandlerMapping

    這篇文章主要介紹了深度理解SpringMVC中的HandlerMapping,HandlerMapping的作用根據(jù)request找到對應(yīng)的處理器Handler,在HandlerMapping接口中有一個唯一的方法getHanler,需要的朋友可以參考下
    2023-09-09

最新評論

修武县| 普定县| 桑植县| 苗栗市| 都江堰市| 柯坪县| 乐清市| 铁力市| 南陵县| 安平县| 衡阳县| 黔西县| 湘西| 长寿区| 介休市| 奉化市| 张家港市| 锦州市| 文安县| 宁德市| 五台县| 太仓市| 苏尼特右旗| 镇远县| 称多县| 迭部县| 恩平市| 安阳县| 赣州市| 南涧| 沙雅县| 南召县| 新宾| 黔西| 腾冲县| 桂平市| 阿坝县| 余干县| 鹤岗市| 凤山县| 揭东县|