【问题标题】:incorrectly invoked stream listeners on rabbitmq message consumer with spring-cloud-stream使用 spring-cloud-stream 在 rabbitmq 消息使用者上错误调用流侦听器
【发布时间】:2021-11-20 11:00:38
【问题描述】:

我有一个项目是spring-cloud-starter-parent:Hoxton.SR9。 这里的主要目标是消息传递,因此声明了消息(kafka/rabbit)生产者和消费者:

spring:
  cloud:
    stream:
      default-binder: kafka
      contentType: application/*+avro
      bindings:
        loanContractWaitingForActivationInput:
          binder: rabbit
          destination: icebank
          group: LH_LoanEventKeeper_testing2
          contentType: application/json
        loanContractComponentsChangedInput:
          binder: rabbit
          destination: icebank
          group: LH_LoanEventKeeper_testing2
          contentType: application/json
        loanOfferPresentedInput:
          binder: rabbit
          destination: icebank
          group: LH_LoanEventKeeper_testing2
          contentType: application/json
        postingEntryRequestInput:
          destination: PAY_PaymentServices_PostingEntryRequest
          group: LH_LoanEventKeeper
          contentType: application/*+avro
          consumer:
            maxAttempts: 1
            autoCommitOnError: true
            headerMode: embeddedHeaders

      kafka:
        bindings:
          postingEntryRequestInput:
            consumer:
              autoCommitOnError: true
        binder:
          configuration:
            security:
              protocol: ${KAFKA_SECURITY_PROTOCOL:SSL}
            ssl:
              truststore:
                location: path
                password: ****
                type: JKS

      rabbit:
        binder:
          useSSL: ${AMQP_SSL_ENABLED:true}
        bindings:
          loanContractWaitingForActivationInput:
            consumer:
              bindingRoutingKey: loans.contracts.waitingForActivation
              exchange-auto-delete: false
          loanContractComponentsChangedInput:
            consumer:
              bindingRoutingKey: loans.contracts.components
              exchange-auto-delete: false
          loanOfferPresentedInput:
            consumer:
              bindingRoutingKey: loans.offers.presented
              exchange-auto-delete: false

连同输入的@StreamListeners 提供:

    @StreamListener(target = MultiInputChannelsRabbit.LOAN_CONTRACT_WAITING_FOR_ACTIVATION_INPUT)
    void receiveLoanContractWaitingForActivation(@Payload LoanContractWaitingForActivationRabbit message, @Headers MessageHeaders headers) {
        messageBusConsumerService.receive(message, headers);
    }

    @StreamListener(target = MultiInputChannelsRabbit.LOAN_CONTRACT_COMPONENTS_CHANGED_INPUT)
    void receiveLoanContractComponentsChanged(@Payload LoanContractComponentsChangedRabbit message, @Headers MessageHeaders headers) {
        messageBusConsumerService.receive(message, headers);
    }

    @StreamListener(target = MultiInputChannelsRabbit.LOAN_OFFER_PRESENTED_INPUT)
    void receiveLoanOfferPresented(@Payload LoanOfferPresentedRabbit message, @Headers MessageHeaders headers) {
        messageBusConsumerService.receive(message, headers);
    }

    @StreamListener(target = MultiInputChannelsKafka.POSTING_ENTRY_REQUEST_INPUT)
    void receivePostingEntryRequest(@Payload PostingEntryRequest message, @Headers MessageHeaders headers) {
        messageBusConsumerService.receive(message, headers);
    }

配置了rabbitmq交换:

话虽这么说,但问题是使用此配置程序在所有 4 个 streamListener 消费者端点上接收具有特定路由键的消息。反过来。所以首先 receiveLoanContractWaitingForActivation(..) 被调用,消息类型与声明的有效负载不匹配,因为第一个参数,然后 receiveLoanContractComponentsChanged(..) 与第一个相同的故事,然后是正确的(根据消息的路由键:receiveLoanOfferPresented(..)),消息在调用正确的侦听器时得到正确处理,最后一个不正确的receivePostingEntryRequest(..) 最终与第一个和第二个相同。

所以基本上我似乎无法将@StreamListener 绑定到我认为将由路由键指定的交换中的正确队列。

请您指出这里缺少/不正确的配置吗?

谢谢!

【问题讨论】:

    标签: java spring-boot rabbitmq spring-cloud-stream spring-rabbit


    【解决方案1】:

    您不能对所有 3 个绑定使用相同的目标/组;你需要一个单独的队列。

    路由键只是队列和交换机之间的绑定;使用您当前的设置,您有一个包含三个绑定的队列。

    更改组,使其独一无二。

    【讨论】:

    • 是的,可以解决这个问题。不幸的是group对于兔子和卡夫卡消费者来说意味着不同的东西,而对于卡夫卡你可以拥有相同的......
    • 即使使用 Kafka,最好使用不同的组 - 否则对一个组进行重新平衡会导致对其他组进行不必要的重新平衡。
    猜你喜欢
    • 1970-01-01
    • 2020-03-06
    • 2021-05-07
    • 1970-01-01
    • 1970-01-01
    • 2021-03-25
    • 1970-01-01
    • 1970-01-01
    • 2018-01-15
    相关资源
    最近更新 更多