Spring Boot集群管理工具KafkaAdminClient使用方法解析
原理介绍
在Kafka官网中这么描述AdminClient:TheAdminClientAPIsupportsmanagingandinspectingtopics,brokers,acls,andotherKafkaobjects.具体的KafkaAdminClient包含了一下几种功能(以Kafka1.0.0版本为准):
- 创建Topic:createTopics(Collection
newTopics) - 删除Topic:deleteTopics(Collection
topics) - 罗列所有Topic:listTopics()
- 查询Topic:describeTopics(Collection
topicNames) - 查询集群信息:describeCluster()
- 查询ACL信息:describeAcls(AclBindingFilterfilter)
- 创建ACL信息:createAcls(Collection
acls) - 删除ACL信息:deleteAcls(Collection
filters) - 查询配置信息:describeConfigs(Collection
resources) - 修改配置信息:alterConfigs(Map
configs) - 修改副本的日志目录:alterReplicaLogDirs(Map
replicaAssignment) - 查询节点的日志目录信息:describeLogDirs(Collection
brokers) - 查询副本的日志目录信息:describeReplicaLogDirs(Collection
replicas) - 增加分区:createPartitions(Map
newPartitions)
其内部原理是使用Kafka自定义的一套二进制协议来实现,详细可以参见Kafka协议。主要实现步骤:
客户端根据方法的调用创建相应的协议请求,比如创建Topic的createTopics方法,其内部就是发送CreateTopicRequest请求。
客户端发送请求至KafkaBroker。
KafkaBroker处理相应的请求并回执,比如与CreateTopicRequest对应的是CreateTopicResponse。
客户端接收相应的回执并进行解析处理。
和协议有关的请求和回执的类基本都在org.apache.kafka.common.requests包中,AbstractRequest和AbstractResponse是这些请求和回执类的两个基本父类。
代码如下
@Component
publicclassKafkaConfig{
//配置Kafka
publicPropertiesgetProps(){
Propertiesprops=newProperties();
props.put("bootstrap.servers","localhost:9092");
/*props.put("retries",2);//重试次数
props.put("batch.size",16384);//批量发送大小
props.put("buffer.memory",33554432);//缓存大小,根据本机内存大小配置
props.put("linger.ms",1000);//发送频率,满足任务一个条件发送*/
props.put("key.serializer","org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer","org.apache.kafka.common.serialization.StringSerializer");
returnprops;
}
}
@RestController
publicclassKafkaTopicManager{
@Autowired
privateKafkaConfigkafkaConfig;
@GetMapping("createTopic")
publicvoidcreateTopic(){
AdminClientadminClient=KafkaAdminClient.create(kafkaConfig.getProps());
NewTopicnewTopic=newNewTopic("test1",4,(short)1);
CollectionnewTopicList=newArrayList<>();
newTopicList.add(newTopic);
adminClient.createTopics(newTopicList);
adminClient.close();
}
@GetMapping("deleteTopic")
publicvoiddeleteTopic(){
AdminClientadminClient=KafkaAdminClient.create(kafkaConfig.getProps());
adminClient.deleteTopics(Arrays.asList("test1"));
adminClient.close();
}
@GetMapping("listAllTopic")
publicvoidlistAllTopic(){
AdminClientadminClient=KafkaAdminClient.create(kafkaConfig.getProps());
ListTopicsResultresult=adminClient.listTopics();
KafkaFuture>names=result.names();
try{
names.get().forEach((k)->{
System.out.println(k);
});
}catch(InterruptedException|ExecutionExceptione){
e.printStackTrace();
}
adminClient.close();
}
@GetMapping("getTopic")
publicvoidgetTopic(){
AdminClientadminClient=KafkaAdminClient.create(kafkaConfig.getProps());
DescribeTopicsResultdescribeTopics=adminClient.describeTopics(Arrays.asList("syn-test"));
Collection>values=describeTopics.values().values();
if(values.isEmpty()){
System.out.println("找不到描述信息");
}else{
for(KafkaFuturevalue:values){
System.out.println(value);
}
}
adminClient.close();
}
}
以上就是本文的全部内容,希望对大家的学习有所帮助,也希望大家多多支持毛票票。
声明:本文内容来源于网络,版权归原作者所有,内容由互联网用户自发贡献自行上传,本网站不拥有所有权,未作人工编辑处理,也不承担相关法律责任。如果您发现有涉嫌版权的内容,欢迎发送邮件至:czq8825#qq.com(发邮件时,请将#更换为@)进行举报,并提供相关证据,一经查实,本站将立刻删除涉嫌侵权内容。