You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用Java(Spring Boot)获取Kafka主题列表及日志压缩配置?

如何用Java判断Kafka主题是否启用日志压缩

当然可以用Java代码实现这个需求!你当前的代码通过Consumer获取主题列表,但Consumer的listTopics()只能拿到分区拓扑信息,没办法直接获取主题的配置项(比如cleanup.policy)。要获取主题配置,推荐使用Kafka官方提供的AdminClient——这是专门用于Kafka集群管理的客户端工具,比Consumer更适合这类运维/配置查询场景。

步骤1:配置AdminClient Bean

首先在Spring Boot中配置AdminClient,只需要指定Kafka Broker地址即可(不需要直接连ZooKeeper,AdminClient会通过Broker获取集群配置,这也是Kafka新版本推荐的方式):

import org.apache.kafka.clients.admin.AdminClient;
import org.apache.kafka.clients.admin.AdminClientConfig;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.beans.factory.annotation.Value;
import java.util.HashMap;
import java.util.Map;

@Configuration
public class KafkaAdminConfig {

    private final String kafkaBrokerList;

    public KafkaAdminConfig(@Value("${spring.kafka.bootstrap-servers}") String kafkaBrokerList) {
        this.kafkaBrokerList = kafkaBrokerList;
    }

    @Bean
    public AdminClient adminClient() {
        Map<String, Object> configs = new HashMap<>();
        configs.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaBrokerList);
        return AdminClient.create(configs);
    }
}

步骤2:实现判断主题是否启用日志压缩的方法

接下来,编写工具方法来查询指定主题的cleanup.policy配置,判断是否包含compact:

import org.apache.kafka.clients.admin.AdminClient;
import org.apache.kafka.clients.admin.Config;
import org.apache.kafka.clients.admin.ConfigResource;
import org.springframework.stereotype.Component;
import java.util.Collections;
import java.util.Optional;

@Component
public class KafkaTopicUtils {

    private final AdminClient adminClient;

    public KafkaTopicUtils(AdminClient adminClient) {
        this.adminClient = adminClient;
    }

    /**
     * 判断指定主题是否启用日志压缩(cleanup.policy=compact)
     * @param topicName 主题名称
     * @return true=启用压缩,false=未启用或主题不存在
     */
    public boolean isTopicCompacted(String topicName) {
        try {
            // 构造主题配置资源
            ConfigResource resource = new ConfigResource(ConfigResource.Type.TOPIC, topicName);
            // 查询主题配置
            Config config = adminClient.describeConfigs(Collections.singleton(resource))
                    .all()
                    .get() // 异步转同步,实际项目中建议用异步回调避免阻塞
                    .get(resource);

            // 获取cleanup.policy配置值,默认是delete
            Optional<String> cleanupPolicy = config.entries().stream()
                    .filter(entry -> entry.name().equals("cleanup.policy"))
                    .map(entry -> entry.value())
                    .findFirst();

            // 判断是否包含compact(可能多个值用逗号分隔,比如"delete,compact")
            return cleanupPolicy.isPresent() && cleanupPolicy.get().contains("compact");
        } catch (Exception e) {
            // 处理主题不存在、网络异常等情况
            e.printStackTrace();
            return false;
        }
    }
}

步骤3:结合原有逻辑获取带压缩标记的主题列表

你可以把原有获取主题列表的逻辑也换成AdminClient(更高效),然后批量判断每个主题是否启用压缩:

import org.apache.kafka.clients.admin.ListTopicsResult;
import org.springframework.stereotype.Component;
import java.util.ArrayList;
import java.util.List;
import java.util.Set;

@Component
public class KafkaTopicService {

    private final KafkaTopicUtils topicUtils;
    private final AdminClient adminClient;
    private static final String CONSUMER_OFFSETS = "__consumer_offsets";

    public KafkaTopicService(KafkaTopicUtils topicUtils, AdminClient adminClient) {
        this.topicUtils = topicUtils;
        this.adminClient = adminClient;
    }

    public List<String> getCompactedTopics() {
        List<String> compactedTopics = new ArrayList<>();
        try {
            // 获取所有主题列表
            ListTopicsResult topicsResult = adminClient.listTopics();
            Set<String> allTopics = topicsResult.names().get();
            allTopics.remove(CONSUMER_OFFSETS);

            // 逐个判断是否启用压缩
            for (String topic : allTopics) {
                if (topicUtils.isTopicCompacted(topic)) {
                    compactedTopics.add(topic);
                }
            }
        } catch (Exception e) {
            e.printStackTrace();
        }
        return compactedTopics;
    }

    // 保留原有获取所有主题的方法(可选)
    public List<String> getAllTopics() {
        try {
            Set<String> allTopics = adminClient.listTopics().names().get();
            allTopics.remove(CONSUMER_OFFSETS);
            return new ArrayList<>(allTopics);
        } catch (Exception e) {
            e.printStackTrace();
            return new ArrayList<>();
        }
    }
}

关键说明

  • 为什么用AdminClient而不是直接连ZooKeeper?:Kafka从2.2版本开始就推荐用AdminClient替代直接ZooKeeper访问,ZooKeeper的内部存储格式可能随版本变化,而AdminClient是官方提供的稳定API,兼容性更好。
  • 配置查询的细节:cleanup.policy可能同时设置多个值(比如delete,compact),所以判断时要用contains("compact")而不是直接等于。
  • 异步处理:示例中用get()将异步操作转为同步,实际生产环境中如果主题数量多,建议使用异步回调(whenComplete)来提升性能,避免阻塞线程。

内容的提问来源于stack exchange,提问作者Debapriyo Majumder

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.14 06:32:58