记录:464
场景:在Spring Boot微服务集成Kafka客户端kafka-clients-3.0.0操作Kafka集群的Topic的创建和删除。
版本:JDK 1.8,Spring Boot 2.6.3,kafka_2.12-2.8.0,kafka-clients-3.0.0。
Kafka集群安装:https://blog.csdn.net/zhangbeizhen18/article/details/131156084
1.微服务中****配置Kafka信息
1.1在pom.xml添加依赖
pom.xml文件:
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>3.0.0</version>
</dependency>
解析:使用原生的kafka-clients,版本:3.0.0。操作kafka集群的Topic。
2.使用AdminClient创建Kafka集群的Topic
AdminClient全称:org.apache.kafka.clients.admin.AdminClient
(1)示例代码
@RestController
@RequestMapping("/hub/example/cluster/topic")
@Slf4j
public class UseKafkaClusterTopicController {
//定义Kafka的Topic
private final String topicName = "hub-topic-city-info-002";
@GetMapping("/f01_1")
public Object f01_1() {
try {
//1.配置Kafka集群
Map<String, Object> configs = new HashMap<>();
Collection<String> cluster = Lists.newArrayList("192.168.19.161:29092",
"192.168.19.162:29092",
"192.168.19.163:29092");
configs.put("bootstrap.servers", cluster);
//2.创建客户端AdminClient
AdminClient adminClient = KafkaAdminClient.create(configs);
//3.获取Kafka集群中Topic清单
Set<String> topicSet = adminClient.listTopics().names().get();
log.info("在Kafka集群已建Topic数量: {} ,清单:", topicSet.size());
topicSet.forEach(System.out::println);
//4.在Kafka集群创建Topic
if (!topicSet.contains(topicName)) {
log.info("新建Topic: {}", topicName);
// Topic名称,分区Partition数目,复制因子(replication Factor)
NewTopic newTopic = new NewTopic(topicName, 1, (short) 1);
Collection<NewTopic> newTopics = Lists.newArrayList(newTopic);
adminClient.createTopics(newTopics);
ThreadUtil.sleep(1000);
topicSet = adminClient.listTopics().names().get();
log.info("创建后,在Kafka集群已建Topic数量: {} ,清单:", topicSet.size());
topicSet.forEach(System.out::println);
}
} catch (Exception e) {
log.info("创建Topic异常.");
e.printStackTrace();
}
return "创建成功";
}
}
(2)解析代码
操作Kafka集群的Topic需要先创建AdminClient,使用AdminClient的API创建Topic。
创建Topic一般只需指定Topic名称,分区Partition数目,复制因子(replication Factor)就行。
3.使用AdminClient删除Kafka集群的Topic
AdminClient全称:org.apache.kafka.clients.admin.AdminClient
(1)示例代码
@RestController
@RequestMapping("/hub/example/cluster/topic")
@Slf4j
public class UseKafkaClusterTopicController {
//定义Kafka的Topic
private final String topicName = "hub-topic-city-info-002";
@GetMapping("/f01_2")
public Object f01_2() {
try {
//1.获取Kafka集群配置信息
Map<String, Object> configs = new HashMap<>();
Collection<String> cluster = Lists.newArrayList("192.168.19.161:29092",
"192.168.19.162:29092",
"192.168.19.163:29092");
configs.put("bootstrap.servers", cluster);
//2.创建客户端AdminClient
AdminClient adminClient = KafkaAdminClient.create(configs);
//3.获取Kafka集群中Topic清单
Set<String> topicSet = adminClient.listTopics().names().get();
log.info("在Kafka集群已建Topic数量: {} ,清单:", topicSet.size());
topicSet.forEach(System.out::println);
//4.在Kafka集群删除Topic
if (topicSet.contains(topicName)) {
log.info("删除Topic: {}", topicName);
Collection<String> topics = Lists.newArrayList(topicName);
DeleteTopicsResult deleteTopicsResult = adminClient.deleteTopics(topics);
deleteTopicsResult.all().get();
ThreadUtil.sleep(1000);
topicSet = adminClient.listTopics().names().get();
log.info("删除后,在Kafka集群已建Topic数量: {} ,清单:", topicSet.size());
topicSet.forEach((topic) -> {
System.out.println(topic);
});
}
} catch (Exception e) {
log.info("删除Topic异常.");
e.printStackTrace();
}
return "删除成功";
}
}
(2)解析代码
操作Kafka集群的Topic需要先创建AdminClient,使用AdminClient的API删除Topic。
创建Topic一般只需指定Topic名称就行。
4.测试
创建请求RUL:http://127.0.0.1:18210/hub-210-kafka/hub/example/cluster/topic/f01_1
删除请求RUL:http://127.0.0.1:18210/hub-210-kafka/hub/example/cluster/topic/f01_2
以上,感谢。
2023年6月18日
版权归原作者 zhangbeizhen18 所有, 如有侵权,请联系我们删除。