Spring Cloud Stream使用StreamBridge动态路由时Service3无法调用问题

问题现象
有3个服务,其中service-1会根据传入负载中id字段的值动态路由到另外两个服务。但目前payment-processor始终路由到后缀为-out-0的主题spring.cloud.stream.function.bindings.processor-out-0,service-3始终无法被调用。尝试过已废弃的配置项spring.cloud.stream.sendto.destination,但依旧无法正常路由到service-3。
问题根源
- processor函数返回值冗余触发额外消息发送:你将processor定义为带返回值的
Function,Spring Cloud Stream会自动将函数返回的所有内容发送到默认的processor-out-0绑定,导致无论分支逻辑怎么走,都会有一条消息发到out-0对应的e-payment主题。 - 多输出绑定未声明数量:processor有2个输出绑定(out-0、out-1),未显式声明输出数量,Spring Cloud Stream默认只会初始化第一个out-0绑定,out-1绑定不存在,导致
streamBridge.send("processor-out-1",xxx)调用无效。 - 冗余配置干扰:service1的
spring.cloud.stream.source=processorCheck;processorEPayment是service2、service3的函数配置,service1不需要该配置,属于冗余干扰项。 - 代码语法错误:service1代码中
payload变量未声明类型,编译阶段就会报错。
修复步骤
方案1:使用StreamBridge实现动态路由(推荐)
1. 修改service-1代码
将不需要返回值的Function改为Consumer,删除冗余返回逻辑,补全变量声明:
@Bean public Consumer<Flux<String>> processor() { return messageFlux -> messageFlux.flatMap( stringMessage -> { Map payload = null; try { payload = jsonMapper.readValue(stringMessage, Map.class); } catch (JsonProcessingException e) { e.printStackTrace(); return Mono.empty(); } logger.info("Payment payload :" + payload.get("id")); // 直接发送到目标主题,避免绑定匹配问题 if(payload.get("id").equalsIgnoreCase("e-payment")) streamBridge.send("e-payment",stringMessage); else streamBridge.send("check-payment",stringMessage); return Mono.empty(); }).subscribe(); } @Bean public Consumer<String> consumer() { return data -> logger.info("Payment Response :" + data); }
2. 修正service-1 application.properties配置
server.port = 9000 spring.application.name=payment-process # 无用到supplier可删除 spring.cloud.function.definition=processor;consumer # 绑定配置 spring.cloud.stream.function.bindings.processor-in-0=payment-process spring.cloud.stream.function.bindings.consumer-in-0=payment-response # 分组配置 spring.cloud.stream.bindings.processor-in-0.group=check-service
方案2:保留Function多输出绑定写法
如果不想用StreamBridge,可以通过消息Header指定路由目标,对应修改如下:
1. 修改service-1代码
@Bean public Function<Flux<String>, Flux<Message<String>>> processor() { return messageFlux -> messageFlux.map(stringMessage -> { Map payload = jsonMapper.readValue(stringMessage, Map.class); String targetDestination = payload.get("id").equalsIgnoreCase("e-payment") ? "e-payment" : "check-payment"; return MessageBuilder.withPayload(stringMessage) .setHeader("spring.cloud.stream.sendto.destination", targetDestination) .build(); }); }
2. 修正service-1 application.properties配置
server.port = 9000 spring.application.name=payment-process spring.cloud.function.definition=processor;consumer # 声明processor有2个输出绑定 spring.cloud.stream.function.bindings.processor.output-count=2 spring.cloud.stream.function.bindings.processor-out-0=e-payment spring.cloud.stream.function.bindings.processor-out-1=check-payment spring.cloud.stream.function.bindings.processor-in-0=payment-process spring.cloud.stream.function.bindings.consumer-in-0=payment-response spring.cloud.stream.bindings.processor-in-0.group=check-service
其他服务无需修改
service-2、service-3的代码和配置均符合要求,不需要调整。
内容的提问来源于stack exchange,提问作者Vishal Mamidi
相关产品推荐
相关产品推荐

