如何在Quarkus中创建指定分区与副本数的Kafka Topic?
在Quarkus中创建指定分区和副本数的Kafka Topic
Spring中的实现方式
在Spring生态中,只需创建NewTopic类型的Bean即可自动完成Kafka Topic的创建,示例代码如下:
@Configuration @RequiredArgsConstructor @EnableKafka @EnableConfigurationProperties(KafkaProperties.class) public class KafkaConfig { private final KafkaProperties kafkaProperties; @Bean public NewTopic requestTopic() { Map<String, String> configs = new HashMap<>(); configs.put("retention.ms", kafkaProperties.getStartReport().getRequestReplyTimeOut().toString()); return new NewTopic(kafkaProperties.getStartReport().getTopicName(), kafkaProperties.getStartReport().getTopicPartitions(), kafkaProperties.getStartReport().getReplicationFactor()).configs(configs); } }
问题场景
想要在Quarkus中实现相同逻辑(创建指定分区数、副本数的Kafka Topic),但以下代码无法生效:
@Slf4j @Startup @ApplicationScoped public class KafkaConfig { @Inject @ConfigProperty(name = "kafka.bootstrap.servers") String kafkaHost; @Inject KafkaOutgoingProperties kafkaOutgoingProperties; @Produces public NewTopic requestTopic() { Map<String, String> configs = new HashMap<>(); return new NewTopic(kafkaOutgoingProperties.finishReportChannel().topic(), kafkaOutgoingProperties.finishReportChannel().partitions(), kafkaOutgoingProperties.finishReportChannel().replicationFactor()).configs(configs); } }
解决方案
Quarkus没有像Spring那样自动识别NewTopic Bean并创建Topic的机制,需通过以下两种方式实现:
方式一:配置文件自动创建(推荐)
通过application.properties直接配置要创建的Topic信息,Quarkus的Kafka客户端扩展会自动完成创建:
# Kafka Broker地址 kafka.bootstrap.servers=your-kafka-broker:9092 # 配置目标Topic kafka.topic.finish-report.name=finish-report-topic kafka.topic.finish-report.partitions=3 kafka.topic.finish-report.replication-factor=1 # 可选:添加Topic额外配置项 kafka.topic.finish-report.config.retention.ms=3600000
方式二:代码主动创建(动态场景)
如果需要动态逻辑控制Topic创建,需使用AdminClient主动执行创建操作:
@Slf4j @Startup @ApplicationScoped public class KafkaTopicCreator { @Inject @ConfigProperty(name = "kafka.bootstrap.servers") String bootstrapServers; @Inject KafkaOutgoingProperties kafkaOutgoingProperties; void onStart(@Observes StartupEvent event) { Map<String, Object> adminConfig = new HashMap<>(); adminConfig.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); try (AdminClient adminClient = AdminClient.create(adminConfig)) { String topicName = kafkaOutgoingProperties.finishReportChannel().topic(); int partitions = kafkaOutgoingProperties.finishReportChannel().partitions(); short replicationFactor = kafkaOutgoingProperties.finishReportChannel().replicationFactor(); // 检查Topic是否已存在 Set<String> existingTopics = adminClient.listTopics().names().get(); if (!existingTopics.contains(topicName)) { NewTopic newTopic = new NewTopic(topicName, partitions, replicationFactor); // 设置Topic额外配置 Map<String, String> topicConfigs = new HashMap<>(); topicConfigs.put("retention.ms", "3600000"); newTopic.configs(topicConfigs); // 执行创建操作 adminClient.createTopics(Collections.singleton(newTopic)).all().get(); log.info("成功创建Topic: {},分区数: {},副本数: {}", topicName, partitions, replicationFactor); } else { log.info("Topic {} 已存在,跳过创建", topicName); } } catch (InterruptedException | ExecutionException e) { log.error("创建Topic失败", e); Thread.currentThread().interrupt(); } } }
原代码失效原因
Quarkus并未内置扫描NewTopic实例并自动触发创建的逻辑,仅通过@Produces生成NewTopic对象不会执行实际的Topic创建操作,必须通过AdminClient主动调用创建接口。
内容的提问来源于stack exchange,提问作者Andrew Samoilov
相关产品推荐
相关产品推荐

