美文网首页
kafka——AdminClient API

kafka——AdminClient API

作者: 小波同学 | 来源:发表于2020-11-30 00:31 被阅读0次

    一、Kafka 核心 API

    下图是官方文档中的一个图,形象的描述了能与 Kafka集成的客户端类型


    Kafka的五类客户端API类型如下:

    • AdminClient API:允许管理和检测Topic、broker以及其他Kafka实例,与Kafka自带的脚本命令作用类似。
    • Producer API:发布消息到1个或多个Topic,也就是生产者或者说发布方需要用到的API。
    • Consumer API:订阅1个或多个Topic,并处理产生的消息,也就是消费者或者说订阅方需要用到的API。
    • Stream API:高效地将输入流转换到输出流,通常应用在一些流处理场景。
    • Connector API:从一些源系统或应用程序拉取数据到Kafka,如上图中的DB。

    本文中,我们将主要介绍 AdminClient API。

    二、Topic 创建与删除

    2.1、创建 topic

    创建 topic 的序列图如下所示:


    • 1、controller 在 ZooKeeper 的 /brokers/topics 节点上注册 watcher,当 topic 被创建,则 controller 会通过 watch 得到该 topic 的 partition/replica 分配。

    • 2、controller从 /brokers/ids 读取当前所有可用的 broker 列表,对于 set_p 中的每一个 partition:

      • 2.1、从分配给该 partition 的所有 replica(称为AR)中任选一个可用的 broker 作为新的 leader,并将AR设置为新的 ISR
      • 2.2、将新的 leader 和 ISR 写入 /brokers/topics/[topic]/partitions/[partition]/state
      1. controller 通过 RPC 向相关的 broker 发送 LeaderAndISRRequest。

    2.2、删除 topic

    删除 topic 的序列图如下所示:


    • 1、controller 在 zooKeeper 的 /brokers/topics 节点上注册 watcher,当 topic 被删除,则 controller 会通过 watch 得到该 topic 的 partition/replica 分配。

    • 2、若 delete.topic.enable=false,结束;否则 controller 注册在 /admin/delete_topics 上的 watch 被 fire,controller 通过回调向对应的 broker 发送 StopReplicaRequest。

    三、AdminClient API

    3.1、导入相关依赖

    <dependency>
        <groupId>org.apache.kafka</groupId>
        <artifactId>kafka-clients</artifactId>
        <version>2.6.0</version>
    </dependency>
    

    3.2、构建AdminClient

    public static AdminClient adminClient(){
        Properties properties = new Properties();
        properties.setProperty(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG,"192.168.174.128:9092");
        AdminClient adminClient = AdminClient.create(properties);
        return adminClient;
    }
    

    3.3、创建Topic实例

    private static final String TOPIC_NAME = "yibo_topic";
    
    /**
     * 创建Topic实例
     */
    public static void createTopic(){
        AdminClient adminClient = AdminSample.adminClient();
        //副本因子
        Short re = 1;
        NewTopic newTopic = new NewTopic(TOPIC_NAME,1,re);
        CreateTopicsResult createTopicsResult = adminClient.createTopics(Arrays.asList(newTopic));
        System.out.println("CreateTopicsResult : " + createTopicsResult);
        adminClient.close();
    }
    

    3.4、创建Topic实例

    private static final String TOPIC_NAME = "yibo_topic";
    
    /**
     * 获取topic列表
     */
    public static void topicList() throws Exception {
        AdminClient adminClient = adminClient();
    
        //是否查看Internal选项
        ListTopicsOptions options = new ListTopicsOptions();
        options.listInternal(true);
    
        //ListTopicsResult listTopicsResult = adminClient.listTopics();
        ListTopicsResult listTopicsResult = adminClient.listTopics(options);
        Set<String> names = listTopicsResult.names().get();
    
        //打印names
        names.stream().forEach(System.out::println);
    
        Collection<TopicListing> topicListings = listTopicsResult.listings().get();
        //打印TopicListing
        topicListings.stream().forEach((topicList) -> {
            System.out.println(topicList.toString());
        });
        adminClient.close();
    }
    

    3.5、删除topic

    private static final String TOPIC_NAME = "yibo_topic";
    
    /**
     * 删除topic
     */
    public static void delTopic() throws Exception {
        AdminClient adminClient = adminClient();
        DeleteTopicsResult deleteTopicsResult = adminClient.deleteTopics(Arrays.asList(TOPIC_NAME));
        deleteTopicsResult.all().get();
    }
    

    3.6、描述topic

    private static final String TOPIC_NAME = "yibo_topic";
    
    /**
     * 描述topic
     * name: yibo_topic
     * desc: (name=yibo_topic,
     *      internal=false,
     *      partitions=
     *          (partition=0,
     *          leader=192.168.174.128:9092 (id: 0 rack: null),
     *          replicas=192.168.174.128:9092 (id: 0 rack: null),
     *          isr=192.168.174.128:9092 (id: 0 rack: null)),
     *          authorizedOperations=null)
     * @throws Exception
     */
    public static void describeTopic() throws Exception {
        AdminClient adminClient = adminClient();
        DescribeTopicsResult describeTopicsResult = adminClient.describeTopics(Arrays.asList(TOPIC_NAME));
        Map<String, TopicDescription> descriptionMap = describeTopicsResult.all().get();
        descriptionMap.forEach((key,value) -> {
            System.out.println("name: " + key+" desc: " + value);
        });
    }
    

    3.7、查询配置信息

    private static final String TOPIC_NAME = "yibo_topic";
    
    /**
     * 查询配置信息
     * ConfigResource(type=TOPIC, name='yibo_topic')
     * Config(
     *      entries=
     *          [ConfigEntry(name=compression.type, value=producer, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=leader.replication.throttled.replicas, value=, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=message.downconversion.enable, value=true, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=min.insync.replicas, value=1, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=segment.jitter.ms, value=0, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=cleanup.policy, value=delete, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=flush.ms, value=9223372036854775807, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=follower.replication.throttled.replicas, value=, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=segment.bytes, value=1073741824, source=STATIC_BROKER_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=retention.ms, value=604800000, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=flush.messages, value=9223372036854775807, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=message.format.version, value=2.6-IV0, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=max.compaction.lag.ms, value=9223372036854775807, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=file.delete.delay.ms, value=60000, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=max.message.bytes, value=1048588, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=min.compaction.lag.ms, value=0, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=message.timestamp.type, value=CreateTime, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=preallocate, value=false, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=min.cleanable.dirty.ratio, value=0.5, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=index.interval.bytes, value=4096, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=unclean.leader.election.enable, value=false, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=retention.bytes, value=-1, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=delete.retention.ms, value=86400000, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=segment.ms, value=604800000, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=message.timestamp.difference.max.ms, value=9223372036854775807, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
     *          ConfigEntry(name=segment.index.bytes, value=10485760, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[])])
     * @throws Exception
     */
    public static void describeConfig() throws Exception {
        AdminClient adminClient = adminClient();
        //TODO 这里做一个预留,集群时会讲到
        //ConfigResource configResource = new ConfigResource(ConfigResource.Type.BROKER,TOPIC_NAME);
    
        ConfigResource configResource = new ConfigResource(ConfigResource.Type.TOPIC,TOPIC_NAME);
        DescribeConfigsResult describeConfigsResult = adminClient.describeConfigs(Arrays.asList(configResource));
        Map<ConfigResource, Config> resourceConfigMap = describeConfigsResult.all().get();
        resourceConfigMap.forEach((key,value) -> {
            System.out.println(key + " " + value);
        });
    }
    

    3.8、修改配置信息 老版API

    private static final String TOPIC_NAME = "yibo_topic";
    
    /**
     * 修改配置信息 老版API
     * @throws Exception
     */
    public static void alterConfig1() throws Exception {
        AdminClient adminClient = adminClient();
        Map<ConfigResource,Config> configMap = new HashMap<>();
        ConfigResource configResource = new ConfigResource(ConfigResource.Type.TOPIC,TOPIC_NAME);
        Config config = new Config(Arrays.asList(new ConfigEntry("preallocate","true")));
        configMap.put(configResource,config);
        AlterConfigsResult alterConfigsResult = adminClient.alterConfigs(configMap);
        alterConfigsResult.all().get();
    }
    

    3.9、修改配置信息 新版API

    private static final String TOPIC_NAME = "yibo_topic";
    
    /**
     * 修改配置信息 新版API
     * @throws Exception
     */
    public static void alterConfig2() throws Exception {
        AdminClient adminClient = adminClient();
        Map<ConfigResource, Collection<AlterConfigOp>> configMap = new HashMap<>();
        ConfigResource configResource = new ConfigResource(ConfigResource.Type.TOPIC,TOPIC_NAME);
        AlterConfigOp alterConfigOp = new AlterConfigOp(new ConfigEntry("preallocate","false"),AlterConfigOp.OpType.SET);
        configMap.put(configResource,Arrays.asList(alterConfigOp));
        AlterConfigsResult alterConfigsResult = adminClient.incrementalAlterConfigs(configMap);
        alterConfigsResult.all().get();
    }
    

    3.10、增加partitions数量

    private static final String TOPIC_NAME = "yibo_topic";
    
    /**
     * 增加partitions数量
     * @param partitions
     * @throws Exception
     */
    public static void incrPartitions(int partitions) throws Exception {
        AdminClient adminClient = adminClient();
        Map<String,NewPartitions> partitionsMap = new HashMap<>();
        NewPartitions newPartitions = NewPartitions.increaseTo(partitions);
        partitionsMap.put(TOPIC_NAME,newPartitions);
        CreatePartitionsResult partitionsResult = adminClient.createPartitions(partitionsMap);
        partitionsResult.all().get();
    }
    

    参考:
    https://www.cnblogs.com/cyfonly/p/5954614.html

    https://www.cnblogs.com/L-Test/p/13439049.html

    相关文章

      网友评论

          本文标题:kafka——AdminClient API

          本文链接:https://www.haomeiwen.com/subject/weilwktx.html