本节摘要:非分区主题是单一队列,一笔账本从头写到尾;分区主题把一条河道拆成若干并行支流,每条支流独立存、独立消费。本节讲清两种形态的取舍、消息键的路由规则,以及面对"流量要涨五倍"时如何定分区数。
一条消息进主题,主题在大促流量下必须并行。Pulsar 的答案是分区(partition):把一个逻辑主题拆成 N 个物理分区,每个分区是一条独立的队列、独立记账本、独立被消费。生产者发消息时按规则选定某个分区,消费者按分区并行接货。非分区主题(non-partitioned)则可以理解为分区数固定为一的特例——单队列,简单但吞吐有限。
与 Kafka 不同的是,Pulsar 的分区不是部署时的物理绑定。正因为第 2 章里账本不认机器,分区在这里更像一个"逻辑编号":加 Broker 不需要动分区,扩分区也不需要搬数据——分区多了只是账本更多,照样条带摊在货仓上。这让"分区数"从沉重的基础设施决策变成了相对轻的应用层决策。
分区主题上,生产者要为每条消息选分区。Pulsar 客户端默认提供三种路由方式:按消息键哈希(同键必同分区,顺序保证的来源)、轮询分发(无键消息均匀摊开)、自定义路由(自己写选择逻辑)。看一段 Java 代码把三种方式跑起来:
Producer<String> producer = client.newProducer(Schema.STRING) .topic("persistent://trade-order/transaction/order-events") // 单分区主题也可显式声明路由,这里演示键哈希(默认带键消息走哈希) .routingMode(RoutingMode.SinglePartition) .create(); // 带键消息:同一订单号永远进同一分区 -> 分区内严格有序 producer.newMessage().key("order-1024").value("CREATED").send(); producer.newMessage().key("order-1024").value("PAID").send(); // 无键消息:按轮询策略摊到各分区,只保证吞吐不保证全局顺序 producer.send("heartbeat-snapshot"); // 看看消息实际落在了哪个分区 producer.newMessage().key("order-1024").value("SHIPPED") .sendAsync() .thenAccept(id -> System.out.println("落在分区: " + id.getPartitionIndex()));
最后打印的 getPartitionIndex 就是这条消息漂进的支流编号。工程含义要钉死:顺序保证的最小单位是"分区内的同一消息键"。想让同一订单的状态变更有序,就把订单号当键;想让同一用户的行为有序,就把用户 ID 当键。全局有序只有一条路——单分区(或非分区主题),代价是吞吐上限。
分区数的决策框架,比任何"建议值"都耐用。分三步想:
第一步算目标并行度。大促峰值吞吐除以单分区可稳定承载的吞吐(实测得出,通常受单分区消费者处理能力约束),向上取整再留余量。比如峰值每秒五万条、实测单分区稳定八千条,则至少七到八个分区。
第二步想消费侧形状。按键共享订阅下,分区数决定上游并行,消费端并行由键的分布决定;共享订阅下分区数影响写入并行度,消费并行独立扩展。若有"分区内顺序"需求,分区数同时约束了消费并行上限。
第三步留增长余量但要克制。Pulsar 加分区是轻操作(新分区新账本,旧数据不动),但加分区会打断"键到分区"的旧映射——同一键的新消息会落到新分区,跨分区的顺序窗口会短暂错乱。所以高频写入强顺序的主题,宁可一次给足余量;吞吐型无键主题则可以少给、随需扩。
⚠️ 常见坑:为了"以后省事"一口气建上千分区。每个分区都有元数据、游标与账本管理成本,主题数乘以分区数上到十万级后,Broker 的内存与名册压力显著上升。分区的正确姿势是按需增长,不是预防性囤积。
既然分区好处多,非分区主题何时用?三类场景:控制信令类小流量(配置下发、服务开关),单队列的简单性是优点;要求严格全局有序且吞吐极低的场景(单实例的任务调度指令);以及客户端生态尚不支持分区查找的老组件对接。判断标准只有一对词:吞吐与有序的相对稀缺性——哪个更稀缺,决定了你选哪扇门。
分区知识最好落在 MessageId 上收尾。前文说过,MessageId 由账本号与条目号组成,完整形态还带着分区信息——把一条已发消息的坐标打出来看:
MessageId id = producer.newMessage().key("order-2048") .value("PAID").send(); System.out.println("分区号: " + id.getPartitionIndex()); System.out.println("完整坐标: " + id); // 典型输出: // 分区号: 3 // 完整坐标: persistent://trade-order/transaction/order-events-partition-3/ledger-18/entry-42
三行输出对应三层理解:分区号说明键哈希把它送进了第四条支流;坐标里的分区主题名(主题名加分区后缀)说明分区在存储层就是独立的物理主题;账本与条目号说明它在这条支流里的确切床位。反过来,运维拿着报障信息里的坐标,可以立刻判断消息落在哪个分区、哪段账本——排查时先读坐标再动手,能省掉大量误操作。
分区不是越早越多,但也别等雪崩才扩。两个信号提示该扩:订阅积压随流量同比增长(消费端已无余力),且分区写入吞吐逼近实测基线(写入端接近天花板)。执行时的顺序值得照做:先扩消费端能力确认(新分区会有新消费者接管)、再改分区数、然后观察键到分区的新映射——强顺序业务的消费端要能容忍短暂的跨分区窗口错乱,通常做法是切换期把消费并行降一档,让每条支流的处理追平后再恢复。
还有个容易忽略的执行细节:扩分区后立即生效的是"新消息的路由",老消息还躺在原分区的账本里。依赖"全量数据在固定分区"的工具(比如按分区号并行的数仓作业)要改成按键感知的读取,否则会漏掉新分区的数据。把这句写进扩分区的操作手册,能避免一次诡异的数据缺失排障。
下一站跟着一条消息走完全程:从 send 调用到消费者手中,每一跳都值得停下来看。