kafka——AdminClient API
一、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
- 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.jianshu.com/p/3f7a81e51c4f