【发布时间】:2022-08-14 01:32:11
【问题描述】:
在这里,我使用分散聚集模式。如果我想在aggregate() 之后和to() 之前调用另一个IntegrationFlow,我该怎么做?我可以在这里使用接收者流,以便我也可以使该流成为有条件的吗?
@Bean
public IntegrationFlow flow() {
return flow ->
flow.handle(validatorService, \"validateRequest\")
.split()
.channel(c -> c.executor(Executors.newCachedThreadPool()))
.scatterGather(
scatterer ->
scatterer
.applySequence(true)
.recipientFlow(flow1())
.recipientFlow(flow2())
.recipientFlow(flow3()),
gatherer ->
gatherer
.releaseLockBeforeSend(true)
.releaseStrategy(group -> group.size() == 2))
.aggregate(lionService.someMethod())
// here I want to call other Integration flows
.gateway(someFlow())
.to(someFlow2());
}
@Bean
public IntegrationFlow flow1() {
return flow ->
flow.channel(c -> c.executor(Executors.newCachedThreadPool()))
.enrichHeaders(h -> h.errorChannel(\"flow1ErrorChannel\", true))
.handle(cdRequestService, \"prepareCDRequestFromLoanRequest\");
}
//same way I have flow2 and flow3, and I have set an custom error channel header for all the flows
@Bean
public IntegrationFlow someFlow() {
return flow ->
flow.filter(\"headers.sourceSystemCode.equals(\"001\")\").channel(c -> c.executor(Executors.newCachedThreadPool()))
.enrichHeaders(h -> h.errorChannel(\"someFlow1ErrorChannel\", true))
.handle( Http.outboundGateway(\"http://localhost:4444/test2\")
.httpMethod(HttpMethod.POST)
.expectedResponseType(String.class)).bridge();
}
到目前为止,只要在任何流程中发生任何错误,它都会通过已分配给它们的自定义错误通道,然后我会处理错误,但是当我在 .gateway(someFlow()) 中使用 someFlow1() 时,该流程中发生的错误不是进入指定的错误通道。如何解决?
在 errorhandler 类中,我正在做类似下面的事情——
//errorhandlerclass
@ServiceActivator(inputChannel = \"flow1ErrorChannel\")
public Message<?> processDBError(MessagingException payload) {
logger.atSevere().withStackTrace(StackSize.FULL).withCause(payload).log(
Objects.requireNonNull(payload.getFailedMessage()).toString());
MessageHeaders messageHeaders = Objects.requireNonNull(payload.getFailedMessage()).getHeaders();
return MessageBuilder.withPayload(
new LionException(ErrorCode.DATABASE_ERROR.getErrorData()))
.setHeader(MessageHeaders.REPLY_CHANNEL, messageHeaders.get(\"originalErrorChannel\"))
.build();
}
@ServiceActivator(inputChannel = \"someFlow1ErrorChannel\")
public Message<?> processDBError(MessagingException payload) {
logger.atSevere().withStackTrace(StackSize.FULL).withCause(payload).log(
Objects.requireNonNull(payload.getFailedMessage()).toString());
MessageHeaders messageHeaders = Objects.requireNonNull(payload.getFailedMessage()).getHeaders();
return MessageBuilder.withPayload(
new LionException(ErrorCode.CUSTOM_ERROR.getErrorData()))
.setHeader(MessageHeaders.REPLY_CHANNEL, messageHeaders.get(\"originalErrorChannel\"))
.build();
}
同样,如果someFlow() 中有任何错误,则会显示错误,但我希望它转到我正在根据我的要求处理错误的方法。
另外,您可以看到我在someFlow() 中使用了过滤器,所以当过滤器表达式评估为真时没有问题,但当它变为假时,它会抛出错误,但我希望它转义并转到下一个,即@ 987654327@。我使用.bridge() 认为它会返回到以前的上下文,但那没有发生。我知道我的理解有些差距。请帮助解决以上两个问题。
标签: java spring-boot spring-integration spring-integration-dsl spring-integration-http