add topic console

This commit is contained in:
许晓东
2021-09-08 21:18:12 +08:00
parent 0b81f40b3f
commit fad17302c8
10 changed files with 196 additions and 2 deletions

View File

@@ -0,0 +1,26 @@
package com.xuxd.kafka.console.beans.vo;
import lombok.Data;
import org.apache.kafka.clients.admin.TopicDescription;
/**
* kafka-console-ui.
*
* @author xuxd
* @date 2021-09-08 21:12:36
**/
@Data
public class TopicDescriptionVO {
private String name;
private boolean internal;
private int partitions;
public static TopicDescriptionVO from(TopicDescription description) {
TopicDescriptionVO vo = new TopicDescriptionVO();
vo.setName(description.name());
vo.setInternal(description.isInternal());
vo.setPartitions(description.partitions().size());
return vo;
}
}

View File

@@ -2,6 +2,7 @@ package com.xuxd.kafka.console.config;
import kafka.console.KafkaAclConsole;
import kafka.console.KafkaConfigConsole;
import kafka.console.TopicConsole;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@@ -23,4 +24,9 @@ public class KafkaConfiguration {
public KafkaAclConsole kafkaAclConsole(KafkaConfig config) {
return new KafkaAclConsole(config);
}
@Bean
public TopicConsole topicConsole(KafkaConfig config) {
return new TopicConsole(config);
}
}

View File

@@ -0,0 +1,26 @@
package com.xuxd.kafka.console.controller;
import com.xuxd.kafka.console.service.TopicService;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
/**
* kafka-console-ui.
*
* @author xuxd
* @date 2021-09-08 20:28:35
**/
@RestController
@RequestMapping("/topic")
public class TopicController {
@Autowired
private TopicService topicService;
@GetMapping("/list")
public Object getTopicList() {
return topicService.getTopicList();
}
}

View File

@@ -0,0 +1,17 @@
package com.xuxd.kafka.console.service;
import com.xuxd.kafka.console.beans.ResponseData;
/**
* kafka-console-ui.
*
* @author xuxd
* @date 2021-09-08 20:01:49
**/
public interface TopicService {
ResponseData getTopicNameList();
ResponseData getTopicList();
}

View File

@@ -0,0 +1,39 @@
package com.xuxd.kafka.console.service.impl;
import com.xuxd.kafka.console.beans.ResponseData;
import com.xuxd.kafka.console.beans.vo.TopicDescriptionVO;
import com.xuxd.kafka.console.service.TopicService;
import java.util.Arrays;
import java.util.Comparator;
import java.util.List;
import java.util.Set;
import kafka.console.TopicConsole;
import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.admin.TopicDescription;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
/**
* kafka-console-ui.
*
* @author xuxd
* @date 2021-09-08 20:02:56
**/
@Slf4j
@Service
public class TopicServiceImpl implements TopicService {
@Autowired
private TopicConsole topicConsole;
@Override public ResponseData getTopicNameList() {
return ResponseData.create().data(topicConsole.getTopicNameList()).success();
}
@Override public ResponseData getTopicList() {
List<TopicDescription> topicDescriptions = topicConsole.getTopicList(topicConsole.getTopicNameList());
topicDescriptions.sort(Comparator.comparing(TopicDescription::name));
return ResponseData.create().data(topicDescriptions.stream().map(d -> TopicDescriptionVO.from(d))).success();
}
}

View File

@@ -7,7 +7,8 @@ server:
kafka:
config:
# kafka broker地址多个以逗号分隔
bootstrap-server: 'localhost:9092'
# bootstrap-server: 'localhost:9092'
bootstrap-server: '10.100.64.48:9092,10.100.77.250:9092,10.100.73.154:9092'
request-timeout-ms: 60000
# 服务端是否启用acl如果不启用下面的几项都忽略即可
enable-acl: true
@@ -16,7 +17,7 @@ kafka:
# 超级管理员用户名在broker上已经配置为超级管理员
admin-username: admin
# 超级管理员密码
admin-password: admin
admin-password: admin!QAZ
# 启动自动创建配置的超级管理员用户
admin-create: true
# broker连接的zk地址

View File

@@ -27,6 +27,14 @@ class KafkaConsole(config: KafkaConfig) {
}
}
protected def withAdminClientAndCatchError(f: Admin => Any, eh: Exception => Any): Any = {
try {
withAdminClient(f)
} catch {
case er: Exception => eh(er)
}
}
protected def withZKClient(f: AdminZkClient => Any): Any = {
val zkClient = KafkaZkClient(config.getZookeeperAddr, false, 30000, 30000, Int.MaxValue, Time.SYSTEM)
val adminZkClient = new AdminZkClient(zkClient)

View File

@@ -0,0 +1,37 @@
package kafka.console
import java.util
import java.util.concurrent.TimeUnit
import java.util.{Collections, List, Set}
import com.xuxd.kafka.console.config.KafkaConfig
import org.apache.kafka.clients.admin.{ListTopicsOptions, TopicDescription}
/**
* kafka-console-ui.
*
* @author xuxd
* @date 2021-09-08 19:52:27
* */
class TopicConsole(config: KafkaConfig) extends KafkaConsole(config: KafkaConfig) with Logging {
def getTopicNameList(): Set[String] = {
withAdminClientAndCatchError(admin => admin.listTopics(new ListTopicsOptions().listInternal(true)).names()
.get(3000, TimeUnit.MILLISECONDS),
e => {
log.error("listTopics error.", e)
Collections.emptySet()
}).asInstanceOf[Set[String]]
}
def getTopicList(topics: Set[String]): List[TopicDescription] = {
if (topics == null || topics.isEmpty) {
Collections.emptyList()
} else {
withAdminClientAndCatchError(admin => new util.ArrayList[TopicDescription](admin.describeTopics(topics).all().get().values()), e => {
log.error("describeTopics error.", e)
Collections.emptyList()
}).asInstanceOf[List[TopicDescription]]
}
}
}