Spring Cloud Stream基于PubSubReactiveFactory的背压流控实现咨询
问题说明
尝试使用Spring Cloud Stream响应式客户端实现Pub/Sub拉取场景下的流控能力,当前实现代码如下:
@Bean ApplicationRunner reactiveSubscriber(PubSubReactiveFactory reactiveFactory, PubSubMessageConverter converter) { return (args) -> { reactiveFactory.poll("orders-subscription", 250L) // Convert a JSON payload into an object .map(msg -> converter.fromPubSubMessage(msg.getPubsubMessage(), Order.class)) .doOnNext(order -> proccessOrder(order)) // Mannually acknowledge the message .doOnNext(AcknowledgeablePubsubMessage::ack); .subscribe(); }; }
经测试当前实现无法通过背压机制限制系统内流转的消息数量,单条消息处理耗时可达数秒,面对每日百万级的消息流入量存在内存占用过高风险。预期效果是Flux管道内恒定维持100条消息并发处理,需要确认:
- 该需求是否有可落地的成熟实现方案
- 替换为响应式RabbitMQ是否可以满足该流控需求
解答
现有Pub/Sub响应式栈的实现方案
当前代码背压失效的核心原因有两点:
- 默认的
poll方法拉取逻辑没有和下游消费速率强绑定,底层客户端预拉取的消息量不会随下游实际处理能力动态调整 - 代码中使用的
map、doOnNext都是同步串行算子,既没有做并发控制,也没有把消息处理的异步完成信号回传给上游,自然无法触发背压调节。
不需要替换中间件,对代码做如下调整即可实现固定100并发的效果:
@Bean ApplicationRunner reactiveSubscriber(PubSubReactiveFactory reactiveFactory, PubSubMessageConverter converter) { return (args) -> { reactiveFactory.poll("orders-subscription", 250L) .flatMap(msg -> Mono.fromRunnable(() -> { Order order = converter.fromPubSubMessage(msg.getPubsubMessage(), Order.class); proccessOrder(order); msg.ack(); }).subscribeOn(Schedulers.boundedElastic()), 100 // 该参数直接控制管道内最大并发处理数 ) .subscribe(); }; }
额外配置要求:需要同步在配置文件中调整Pub/Sub订阅端的流控参数,将最大未确认消息数、并行拉取数、消费线程数设置为和100的并发值匹配,避免客户端底层预拉取过量消息堆积在内存中。
响应式RabbitMQ的适配性说明
响应式RabbitMQ客户端原生做了背压和预取机制(prefetch count)的打通:下游每处理完成一条消息,才会向Broker请求拉取新的消息,只要把prefetch值设置为100,天然就能实现管道内最多100条待处理消息的流控效果,不需要额外做算子层面的并发限制。但从成本角度考虑,现有Pub/Sub技术栈只需要调整代码和配置就可以完全满足流控需求,专门为了流控能力替换消息中间件的迁移成本远高于现有方案改造成本,没有必要。
内容的提问来源于stack exchange,提问作者Jonathan Chevalier
相关产品推荐
相关产品推荐

