不同消费者组的Kafka主题保留时间
创始人
2025-01-09 22:00:39
0

要为不同消费者组设置不同的Kafka主题保留时间,可以通过在Kafka配置文件中设置retention.ms参数来实现。下面是一个示例代码,展示了如何设置不同消费者组的Kafka主题保留时间:

import java.util.Properties;
import org.apache.kafka.clients.admin.AdminClient;
import org.apache.kafka.clients.admin.AdminClientConfig;
import org.apache.kafka.clients.admin.AlterConfigOp;
import org.apache.kafka.clients.admin.Config;
import org.apache.kafka.clients.admin.ConfigEntry;
import org.apache.kafka.clients.admin.NewPartitions;
import org.apache.kafka.clients.admin.NewTopic;
import org.apache.kafka.clients.admin.TopicDescription;
import org.apache.kafka.common.config.ConfigResource;
import org.apache.kafka.common.config.TopicConfig;

public class KafkaTopicRetentionTime {

    public static void main(String[] args) throws Exception {
        // Kafka broker配置
        Properties properties = new Properties();
        properties.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");

        // 创建AdminClient
        try (AdminClient adminClient = AdminClient.create(properties)) {

            // 创建一个新的Kafka主题
            String topicName = "test_topic";
            short replicationFactor = 1;
            int numPartitions = 1;
            NewTopic newTopic = new NewTopic(topicName, numPartitions, replicationFactor);
            adminClient.createTopics(Collections.singletonList(newTopic)).all().get();

            // 设置消费者组A的保留时间为24小时
            String consumerGroupA = "consumer_group_a";
            int retentionTimeA = 24 * 60 * 60 * 1000; // 24小时
            setTopicRetentionTime(adminClient, topicName, consumerGroupA, retentionTimeA);

            // 设置消费者组B的保留时间为48小时
            String consumerGroupB = "consumer_group_b";
            int retentionTimeB = 48 * 60 * 60 * 1000; // 48小时
            setTopicRetentionTime(adminClient, topicName, consumerGroupB, retentionTimeB);
        }
    }

    private static void setTopicRetentionTime(AdminClient adminClient, String topicName, String consumerGroup, int retentionTime) throws Exception {
        // 获取主题配置
        ConfigResource configResource = new ConfigResource(ConfigResource.Type.TOPIC, topicName);
        Config config = adminClient.describeConfigs(Collections.singleton(configResource)).all().get().get(configResource);

        // 设置保留时间配置
        ConfigEntry retentionEntry = new ConfigEntry(TopicConfig.RETENTION_MS_CONFIG, String.valueOf(retentionTime));
        Config newConfig = new Config(config.entries());
        newConfig.add(retentionEntry);

        // 修改主题配置
        adminClient.alterConfigs(Collections.singletonMap(configResource, newConfig)).all().get();

        // 输出修改后的主题配置
        TopicDescription topicDescription = adminClient.describeTopics(Collections.singletonList(topicName)).all().get().get(topicName);
        System.out.println("Topic: " + topicDescription.name());
        System.out.println("Retention Time for Consumer Group " + consumerGroup + ": " + topicDescription.configs().get(TopicConfig.RETENTION_MS_CONFIG).value());
    }
}

上述代码中,首先创建了一个AdminClient实例,然后使用createTopics方法创建一个名为test_topic的Kafka主题。

接下来,通过setTopicRetentionTime方法,分别为消费者组A和消费者组B设置了不同的主题保留时间。该方法首先使用describeConfigs方法获取主题的当前配置,然后为主题添加了一个新的保留时间配置TopicConfig.RETENTION_MS_CONFIG,并使用alterConfigs方法将新配置应用到主题上。

最后,使用describeTopics方法获取了修改后的主题配置,并输出了消费者组A和消费者组B的保留时间。

相关内容

热门资讯

4分钟外挂!边锋微信小程序,四... 4分钟外挂!边锋微信小程序,四川途游辅助软件下载,2025新版(有挂方式)-哔哩哔哩1、下载好四川途...
第4分钟普及!家乡大二辅助,四... 第4分钟普及!家乡大二辅助,四川游戏家园通用辅助(其实有挂)-哔哩哔哩1、四川游戏家园通用辅助系统规...
第四分钟外挂!皮皮跑胡子修改器... 第四分钟外挂!皮皮跑胡子修改器,友友联盟免费辅助器,安装教程(有挂方法)-哔哩哔哩1、进入到友友联盟...
第8分钟详情!方片十三张源码,... 第8分钟详情!方片十三张源码,河洛杠次脚本开发(一贯有挂)-哔哩哔哩1)河洛杠次脚本开发辅助挂:进一...
1分钟外挂!传送屋激k有挂吗,... 1分钟外挂!传送屋激k有挂吗,兴动互娱辅助工具,解密教程(有挂分析)-哔哩哔哩1、超多福利:超高返利...
第6分钟解谜!山西扣点点智能辅... 第6分钟解谜!山西扣点点智能辅助器软件,决战卡五星辅助(一直是真的挂)-哔哩哔哩1、玩家可以在山西扣...
第五分钟推荐!开心泉州免费辅助... 第五分钟推荐!开心泉州免费辅助器,友友联盟免费辅助器(真是真的有挂)-哔哩哔哩1、完成开心泉州免费辅...
六分钟外挂!指尖四川小程序辅助... 六分钟外挂!指尖四川小程序辅助,八闽福建辅助,微扑克教程(存在有挂)-哔哩哔哩暗藏猫腻,小编详细说明...
第4分钟必备!潮汕木虱辅助下载... 第4分钟必备!潮汕木虱辅助下载,天天飞小鸡辅助(一贯存在有挂)-哔哩哔哩1、操作简单,无需注册,只需...
七分钟外挂!浙江宝宝游戏辅助下... 七分钟外挂!浙江宝宝游戏辅助下载,小程序牵手跑的辅助,安装教程(有挂解惑)-哔哩哔哩1)浙江宝宝游戏...