【发布时间】: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);
}
话虽这么说,但问题是使用此配置程序在所有 4 个 streamListener 消费者端点上接收具有特定路由键的消息。反过来。所以首先
receiveLoanContractWaitingForActivation(..) 被调用,消息类型与声明的有效负载不匹配,因为第一个参数,然后 receiveLoanContractComponentsChanged(..) 与第一个相同的故事,然后是正确的(根据消息的路由键:receiveLoanOfferPresented(..)),消息在调用正确的侦听器时得到正确处理,最后一个不正确的receivePostingEntryRequest(..) 最终与第一个和第二个相同。
所以基本上我似乎无法将@StreamListener 绑定到我认为将由路由键指定的交换中的正确队列。
请您指出这里缺少/不正确的配置吗?
谢谢!
【问题讨论】:
标签: java spring-boot rabbitmq spring-cloud-stream spring-rabbit