如何在给定Java代码中为输入Topic添加compact清理策略配置?
在Kafka AdminClient代码中为Topic添加cleanup.policy=compact配置
要给你的输入Topic设置cleanup.policy=compact,你只需要利用NewTopic类支持自定义配置的特性即可——它提供了带配置参数的构造方法,或者可以通过configs()方法来添加配置项。
修改后的完整代码
public void createTopics(Properties envProps) { Map<String, Object> config = new HashMap<>(); config.put("bootstrap.servers", envProps.getProperty("bootstrap.servers")); AdminClient client = AdminClient.create(config); List<NewTopic> topics = new ArrayList<>(); // 为输入Topic创建配置Map,添加cleanup.policy=compact Map<String, String> inputTopicConfigs = new HashMap<>(); inputTopicConfigs.put("cleanup.policy", "compact"); // 创建带自定义配置的输入Topic topics.add(new NewTopic( envProps.getProperty("input.topic.name"), Integer.parseInt(envProps.getProperty("input.topic.partitions")), Short.parseShort(envProps.getProperty("input.topic.replication.factor")), inputTopicConfigs )); // 输出Topic保持原有创建方式 topics.add(new NewTopic( envProps.getProperty("output.topic.name"), Integer.parseInt(envProps.getProperty("output.topic.partitions")), Short.parseShort(envProps.getProperty("output.topic.replication.factor")) )); client.createTopics(topics); client.close(); } public Properties loadEnvProperties(String fileName) throws IOException { Properties envProps = new Properties(); FileInputStream input = new FileInputStream(fileName); envProps.load(input); input.close(); return envProps; }
关键说明
- 我们创建了一个
Map<String, String>类型的inputTopicConfigs,把cleanup.policy设置为"compact"; - 然后使用
NewTopic的四参数构造方法(最后一个参数就是Topic的配置Map)来创建输入Topic; - 如果你更喜欢链式调用的风格,也可以这样写:
NewTopic inputTopic = new NewTopic( envProps.getProperty("input.topic.name"), Integer.parseInt(envProps.getProperty("input.topic.partitions")), Short.parseShort(envProps.getProperty("input.topic.replication.factor")) ); inputTopic.configs(inputTopicConfigs); topics.add(inputTopic);
这样就能成功为输入Topic添加日志压缩的清理策略了。
内容的提问来源于stack exchange,提问作者subrahmanyam b
相关产品推荐
相关产品推荐

