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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 17:45:29