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

Java环境下如何对已创建的Apache Pulsar Topic进行分区?

已创建的Pulsar Topic分区处理方案

首先明确:Pulsar中已创建的非分区Topic无法直接转换为分区Topic,需通过以下步骤完成迁移:

1. 创建新的分区Topic

先创建对应名称的分区Topic,以原Topic名称test为例,可通过两种方式操作:

  • 命令行方式:
bin/pulsar-admin topics create-partitioned-topic persistent://public/default/test --partitions 3

(3为分区数量,可根据业务需求调整)

  • Java Admin API方式:
PulsarAdmin admin = PulsarAdmin.builder().serviceHttpUrl("http://pulsar-broker:8080").build();
admin.topics().createPartitionedTopic("persistent://public/default/test", 3);

2. 迁移Producer到新分区Topic

修改现有Producer Bean,确保指向刚创建的分区Topic(Pulsar中分区Topic与非分区Topic为不同实体,同名情况下优先识别已创建的分区版本):

@Bean(name = "test")
public Producer<Test> testProducer(PulsarClient pulsarClient){
    return pulsarClient.newProducer()
            .topic("test")
            .create();
}

若担心命名混淆,可直接指定完整Topic路径:.topic("persistent://public/default/test")

3. 迁移Consumer到新分区Topic

若存在对应Consumer,同样修改配置指向分区Topic,Pulsar Consumer默认会订阅分区Topic的所有分区:

Consumer<Test> consumer = pulsarClient.newConsumer()
        .topic("test")
        .subscriptionName("test-sub")
        .subscribe();

4. 清理原非分区Topic

待所有流量完全迁移到新分区Topic后,删除原非分区Topic:

bin/pulsar-admin topics delete persistent://public/default/test

内容的提问来源于stack exchange,提问作者Dovile Barkauskaite

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 03:01:24