mirror of
https://github.com/didi/KnowStreaming.git
synced 2026-01-08 15:52:15 +08:00
增加返回Kafka内部Topic
This commit is contained in:
@@ -372,7 +372,7 @@ public class PartitionServiceImpl extends BaseVersionControlService implements P
|
|||||||
AdminClient adminClient = kafkaAdminClient.getClient(clusterPhy.getId());
|
AdminClient adminClient = kafkaAdminClient.getClient(clusterPhy.getId());
|
||||||
|
|
||||||
// 获取Topic列表
|
// 获取Topic列表
|
||||||
ListTopicsResult listTopicsResult = adminClient.listTopics(new ListTopicsOptions().timeoutMs(KafkaConstant.ADMIN_CLIENT_REQUEST_TIME_OUT_UNIT_MS));
|
ListTopicsResult listTopicsResult = adminClient.listTopics(new ListTopicsOptions().timeoutMs(KafkaConstant.ADMIN_CLIENT_REQUEST_TIME_OUT_UNIT_MS).listInternal(true));
|
||||||
for (String topicName: listTopicsResult.names().get()) {
|
for (String topicName: listTopicsResult.names().get()) {
|
||||||
DescribeTopicsResult describeTopicsResult = adminClient.describeTopics(
|
DescribeTopicsResult describeTopicsResult = adminClient.describeTopics(
|
||||||
Arrays.asList(topicName),
|
Arrays.asList(topicName),
|
||||||
|
|||||||
@@ -238,7 +238,7 @@ public class TopicServiceImpl implements TopicService {
|
|||||||
try {
|
try {
|
||||||
AdminClient adminClient = kafkaAdminClient.getClient(clusterPhy.getId());
|
AdminClient adminClient = kafkaAdminClient.getClient(clusterPhy.getId());
|
||||||
|
|
||||||
ListTopicsResult listTopicsResult = adminClient.listTopics(new ListTopicsOptions().timeoutMs(KafkaConstant.ADMIN_CLIENT_REQUEST_TIME_OUT_UNIT_MS));
|
ListTopicsResult listTopicsResult = adminClient.listTopics(new ListTopicsOptions().timeoutMs(KafkaConstant.ADMIN_CLIENT_REQUEST_TIME_OUT_UNIT_MS).listInternal(true));
|
||||||
|
|
||||||
List<Topic> topicList = new ArrayList<>();
|
List<Topic> topicList = new ArrayList<>();
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user