Webflux流中返回业务对象并按需发送消息总线消息的问题
解决方法
核心思路是保留generateBusinessObject()生成的业务对象,同时在它执行成功后触发消息发送操作,可根据消息发送的容错需求选择不同方案:
方案1:消息发送不影响业务对象返回
如果消息发送的成败不影响接口返回结果(即便消息发送失败,仍返回正常生成的业务对象),可以这样实现:
return Mono.just(id) .flatMap(this::checkAccess) // 权限校验不通过会触发错误流,终止后续操作 .flatMap(ignore -> generateBusinessObject()) // 校验通过后生成业务对象 .flatMap(businessObj -> // 发送消息完成后,返回原业务对象 postMessageOnBus() .thenReturn(businessObj) // 可选:消息发送失败时吞掉错误,继续返回业务对象 .onErrorReturn(businessObj) ) .toFuture();
方案2:消息发送失败时整体返回错误
如果要求消息必须发送成功,否则接口返回错误,可通过Mono.zip合并操作:
return Mono.just(id) .flatMap(this::checkAccess) .flatMap(ignore -> generateBusinessObject()) .flatMap(businessObj -> Mono.zip( Mono.just(businessObj), postMessageOnBus() ).map(tuple -> tuple.getT1()) // 仅返回业务对象 ) .toFuture();
关键说明
checkAccess方法需保证:校验不通过时抛出异常或返回空Mono,才能触发错误流,阻止后续操作执行。- 原代码中
onErrorResume(throwable -> return Mono.error())属于错误写法,正确形式为onErrorResume(throwable -> Mono.error(throwable)),但该逻辑可省略——Reactor默认会向上传播错误,仅在需要自定义错误转换时才需配置。
内容的提问来源于stack exchange,提问作者Cristian
相关产品推荐
相关产品推荐

