Spring Cloud Stream下RabbitMQ消费者@StreamListener错误触发问题排查
问题根源
- RabbitMQ队列绑定错误:你配置的3个RabbitMQ绑定使用了相同的
destination(对应RabbitMQ的Exchange,值为icebank)和相同的group,Spring Cloud Stream Rabbit Binder的队列命名规则为{destination}.{group},相同的destination+group会生成同一个队列,该队列会绑定你3个配置中所有的bindingRoutingKey,因此队列会收到3个路由键的全部消息,而非每个路由键对应独立队列。 - @StreamListener缺少匹配规则:未给
@StreamListener添加过滤条件时,Spring Cloud Stream会将所有收到的消息依次尝试匹配所有监听方法的入参类型,类型匹配成功则执行方法,匹配失败就抛出异常继续匹配下一个方法,这就是为什么单个消息会依次触发4个监听方法,甚至跨Binder触发Kafka对应的监听方法。
修复方案
- 修改RabbitMQ绑定的group配置:给每个RabbitMQ绑定分配独立的group,确保每个绑定生成独立的队列,每个队列仅绑定对应自己的路由键:
spring: cloud: stream: bindings: loanContractWaitingForActivationInput: binder: rabbit destination: icebank group: LH_LoanEventKeeper_testing2_activation # 独立group contentType: application/json loanContractComponentsChangedInput: binder: rabbit destination: icebank group: LH_LoanEventKeeper_testing2_components # 独立group contentType: application/json loanOfferPresentedInput: binder: rabbit destination: icebank group: LH_LoanEventKeeper_testing2_offer # 独立group contentType: application/json
- 给@StreamListener添加过滤条件:通过
condition属性指定监听方法的匹配规则,避免非目标消息触发方法:
// RabbitMQ监听方法添加路由键匹配规则 @StreamListener( target = MultiInputChannelsRabbit.LOAN_CONTRACT_WAITING_FOR_ACTIVATION_INPUT, condition = "headers['amqp_receivedRoutingKey'] == 'loans.contracts.waitingForActivation'" ) void receiveLoanContractWaitingForActivation(@Payload LoanContractWaitingForActivationRabbit message, @Headers MessageHeaders headers) { messageBusConsumerService.receive(message, headers); } @StreamListener( target = MultiInputChannelsRabbit.LOAN_CONTRACT_COMPONENTS_CHANGED_INPUT, condition = "headers['amqp_receivedRoutingKey'] == 'loans.contracts.components'" ) void receiveLoanContractComponentsChanged(@Payload LoanContractComponentsChangedRabbit message, @Headers MessageHeaders headers) { messageBusConsumerService.receive(message, headers); } @StreamListener( target = MultiInputChannelsRabbit.LOAN_OFFER_PRESENTED_INPUT, condition = "headers['amqp_receivedRoutingKey'] == 'loans.offers.presented'" ) void receiveLoanOfferPresented(@Payload LoanOfferPresentedRabbit message, @Headers MessageHeaders headers) { messageBusConsumerService.receive(message, headers); } // Kafka监听方法添加专属header匹配规则,避免匹配RabbitMQ消息 @StreamListener( target = MultiInputChannelsKafka.POSTING_ENTRY_REQUEST_INPUT, condition = "headers.containsKey('kafka_receivedTopic')" ) void receivePostingEntryRequest(@Payload PostingEntryRequest message, @Headers MessageHeaders headers) { messageBusConsumerService.receive(message, headers); }
- (可选)配置无效消息直接丢弃:如果不需要类型匹配失败的消息继续流转,可以给每个绑定添加如下配置,避免匹配失败抛出异常:
spring: cloud: stream: bindings: <对应通道名>: consumer: dropInvalid: true
内容的提问来源于stack exchange,提问作者greengold
相关产品推荐
相关产品推荐

