如何让Spring Integration处理返回Flux的MessageHandler并传递解析对象
问题解决:让Spring Integration正确处理Flux并传递单个元素
你当前的问题在于,第一个handle方法返回的Flux被直接当作完整payload传给了下一个处理器,所以输出的是FluxFlatMap对象,而非Flux内部的Account或Order元素。Spring Integration默认不会自动订阅Flux并拆分元素,需要用专门的反应式处理方式来解决。
两种可行解决办法:
1. 用handleReactive替代普通handle
将返回Flux的handle替换为handleReactive,它会自动订阅Flux,并把每个元素作为独立消息发送给后续处理器:
@Bean public IntegrationFlow flow1() { Function<Map,Flux<Account>> handler = (map)->service.getAccounts(map) .take(20) .log() .flatMap(account -> orderService.getOrdersForAcct(account).take(3)); return IntegrationFlow.from("accountMessageChannel") .transform(Transformers.fromJson(Map.class)) .handleReactive((payload,headers)-> handler.apply(payload)) // 替换为handleReactive .handle((payload,headers)->{log.info("got something");log.info(payload.toString());return null;}) .get(); }
2. 在handle后添加split()操作符
如果不想替换handle方法,可以在返回Flux的handle步骤后追加.split(),它会把Flux拆分为单个元素逐个传递:
@Bean public IntegrationFlow flow1() { Function<Map,Flux<Account>> handler = (map)->service.getAccounts(map) .take(20) .log() .flatMap(account -> orderService.getOrdersForAcct(account).take(3)); return IntegrationFlow.from("accountMessageChannel") .transform(Transformers.fromJson(Map.class)) .handle((payload,headers)-> handler.apply(payload)) .split() // 新增split拆分Flux元素 .handle((payload,headers)->{log.info("got something");log.info(payload.toString());return null;}) .get(); }
至于你尝试的FluxMessageChannel无效,是因为仅传递Flux到通道,并没有告诉框架要拆分它。上面两种方法的核心都是让Spring Integration订阅Flux并逐个发送元素,这样后续处理器就能拿到具体的Account或Order对象了。
内容的提问来源于stack exchange,提问作者VPN236
相关产品推荐
相关产品推荐

