kafka——AdminClient API

2020-11-30  本文已影响0人  小波同学

一、Kafka 核心 API

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


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

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

二、Topic 创建与删除

2.1、创建 topic

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


2.2、删除 topic

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


三、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

上一篇下一篇

猜你喜欢

热点阅读