本节摘要:消费者代码只有三段式:订阅、取消息、签收——但每一段都有分岔口。本节把监听模式与循环拉取两种写法的适用场景摆开,讲清确认时机的两条铁律,并把重试、死信与按键共享的客户端配置一次配齐。读完这一节,第 3 章的全部投递语义都能落进你的代码。
先看骨架。任何 Pulsar 消费者都是这三步,复杂度都在细节里:
Consumer<String> consumer = client.newConsumer(Schema.STRING) .topic("persistent://trade-order/transaction/order-events") .subscriptionName("stock-deduct") // 订阅名:岗位名 .subscriptionType(SubscriptionType.Key_Shared) // 模式:键内有序并行 .subscribe(); while (running) { Message<String> msg = consumer.receive(5, TimeUnit.SECONDS); // 取货 if (msg == null) continue; // 超时轮询 try { handleOrderEvent(msg); // 业务处理 consumer.acknowledge(msg); // 签收 } catch (Exception e) { consumer.negativeAcknowledge(msg); // 打回重投 } }
三段里最容易写错的是最后一段的位置:确认必须在业务处理成功之后调用,处理与确认之间的顺序就是"不丢"与"可能重"的分界线。确认放在处理前,处理一崩这条消息就算漂完了——丢了;确认放在成功后,崩在确认前则会被重投——可能重但不丢。第 3 章的结论在这里落成一行代码的位置问题:先处理后签收,配幂等。
循环拉取(上面骨架的写法)把节奏握在自己手里:用 receive 带超时轮询,循环里可以穿插限流、优雅停机检查、批量聚合等自定义逻辑。监听模式(messageListener)则是回调驱动:注册一个监听器,客户端线程自动拉消息并调用你,写法短但线程模型归了库管。
选型的经验法则:需要与线程池、限流器、事务边界协作的消费逻辑用循环拉取,代码直白好调试;简单的转发型处理(收下来丢给下游或落库)用监听模式省样板代码。注意监听模式里如果回调抛异常而不做 nack,默认行为是等到确认超时才重投——延迟远大于显式 nack,所以监听器内部同样要写 try-catch。

// 监听模式等价写法:短,但线程模型在客户端库手里 consumer = client.newConsumer(Schema.STRING) .topic(TOPIC) .subscriptionName("stock-deduct") .subscriptionType(SubscriptionType.Key_Shared) .messageListener((c, msg) -> { try { handleOrderEvent(msg); c.acknowledge(msg); } catch (Exception e) { c.negativeAcknowledge(msg); // 显式打回,别等超时 } }) .subscribe();
receiverQueueSize 决定消费者本地预取多少条消息。预取大,单实例吞吐高(省网络往返),但某实例宕机时它缓冲里未处理的消息要等确认超时才重投,尾延迟与积压转移都变差;预取小,分发均匀、故障收敛快,吞吐让步。批量消费场景(一次取一批落库)配大预取配合批量确认是正确姿势;均衡敏感的共享订阅把预取调小(比如五十)能显著改善负载均衡的公平性。
💡 关键直觉:预取是"把消息先押在你这"。押得多干得快,但押金丢了要赔——每个消费者实例的宕机代价与预取深度成正比。
一段把死信、重试、按键共享、批量确认全部配置齐的代码,可直接作为模板:
Consumer<String> consumer = client.newConsumer(Schema.STRING) .topic(TOPIC) .subscriptionName("stock-deduct") .subscriptionType(SubscriptionType.Key_Shared) .receiverQueueSize(200) .ackTimeout(30, TimeUnit.SECONDS) // 处理超时自动重投 .ackTimeoutTickTime(5, TimeUnit.SECONDS) // 超时检查精度 .negativeAckRedeliveryDelay(10, TimeUnit.SECONDS) .deadLetterPolicy(DeadLetterPolicy.builder() .maxRedeliverCount(5) .deadLetterTopic(TOPIC + "-STOCK-DLQ") // 耗尽即入死信 .initialSubscriptionName("dlq-inspect")) // 死信自动建订阅便于检查 .subscribe();
运行验证流程建议照做:正常消息签收后,在管理端看该订阅的 msgBacklog 归零;人为抛异常的消息,观察它按十秒间隔重投、五次后出现在死信主题;再用一个临时消费者订阅死信主题确认内容完好。这条验证链跑通,说明第 3 章的语义在你的代码里真正闭环了。
消费代码的收尾功夫在停机。进程收到下线信号时,正确顺序是:停止取新消息、把手头的消息处理完并签收、调用 close 释放订阅连接。跳过收尾直接 kill 的代价是"在途消息全变未确认",要等确认超时后重投——短的重投延迟窗口可以接受,但如果停机是滚动发布,每个实例都甩一遍在途消息,重投风暴会拖慢整个发布的收敛。把优雅停机写进消费骨架的模板代码里,是性价比极高的一行工程投资。
与停机相关的是再均衡期间的流量预热。按键共享与共享模式下,新实例加入时分货比例会重新分配,新实例冷启动(缓存空、连接池未热)却立刻接到满额货,常见现象是发布后短暂延迟升高。给新实例留预热期(启动后先小量接收,逐步放开)或把发布安排在低峰,都是简单有效的缓冲手段——这些细节不属于消息系统本身,却决定着上线体验的顺滑度。
下一站认识两位特殊的工具:不排队只回放的 Reader,以及 Java 之外的多语言客户端。