【问题标题】:Main processing flow programmatic approach when using Spring Integration with Project Reactor使用 Spring Integration with Project Reactor 时的主要处理流程编程方法
【发布时间】:2022-01-21 15:39:26
【问题描述】:

我想定义一个使用 Reactor Kafka 消耗 kafka 并写入 MongoDB 的流,并且只有在成功时才会将 ID 写入 Kafka。我正在使用带有 Spring Integration JavaDSL 的 Project Reactor,并且我希望有一个 FlowBuilder 类来定义我的高级管道。我目前有以下方向:

public IntegrationFlow buildFlow() {
   return IntegrationFlows.from(reactiveKafkaConsumerTemplate)
      .publishSubscribeChannel(c -> c
                        .subscribe(sf -> sf
                                .handle(MongoDb.reactiveOutboundChannelAdapter())) 
      .handle(writeToKafka)
      .get();
}

我见过docs that there is a support for a different approach, that also works with Project Reactor。这种方法不包括使用IntegrationFlows。这看起来像这样:

@MessagingGateway
public static interface TestGateway {

    @Gateway(requestChannel = "promiseChannel")
    Mono<Integer> multiply(Integer value);

    }

        ...

    @ServiceActivator(inputChannel = "promiseChannel")
    public Integer multiply(Integer value) {
            return value * 2;
    }

        ...

    Flux.just("1", "2", "3", "4", "5")
            .map(Integer::parseInt)
            .flatMap(this.testGateway::multiply)
            .collectList()
            .subscribe(integers -> ...);

我想知道在使用这两个库时更推荐的处理方式。我想知道如何在第二个示例中使用 Reactive MongoDB 适配器。如果没有 IntegrationFlows 包装器,我不确定第二种方法是否可行。

【问题讨论】:

    标签: spring spring-boot spring-integration project-reactor reactor-kafka


    【解决方案1】:

    @MessagingGateway 是为高级最终用户 API 设计的,以尽可能地将消息隐藏在下面。因此,当您开发目标服务的逻辑时,目标服务没有任何消息传递抽象。

    可以使用IntegrationFlow 中的这样一个接口适配器,您应该将其视为常规服务激活器,因此它看起来像这样:

    .handle("testGateway", "multiply", e -> e.async(true))
    

    async(true) 使这个服务激活器订阅返回的Mono。您可以省略它,然后您自己在下游订阅它,因为这个 Mono 将成为流中下一条消息的 payload

    如果你想要相反的东西:从Flux 调用IntegrationFlow,就像flatMap(),然后考虑使用流定义中的toReactivePublisher() 运算符返回Publisher&lt;?&gt; 并声明它作为一颗豆子。在这种情况下,最好不要使用 MongoDb.reactiveOutboundChannelAdapter(),而只使用 ReactiveMongoDbStoringMessageHandler 让它返回的 Mono 传播到 Publisher

    另一方面,如果您希望 @MessagingGatewayMono 返回,但仍从中调用 ReactiveMongoDbStoringMessageHandler,则将其声明为 bean 并用 @ServiceActivator 标记它。

    我们还有一个ExpressionEvaluatingRequestHandlerAdvice 来捕获特定端点上的错误(或成功)并分别处理它们:https://docs.spring.io/spring-integration/docs/current/reference/html/messaging-endpoints.html#expression-advice

    我想你要找的是这样的:

    public IntegrationFlow buildFlow() {
       return IntegrationFlows.from(reactiveKafkaConsumerTemplate)
          .handle(reactiveMongoDbStoringMessageHandler, "handleMessage")
          .handle(writeToKafka)
          .get();
    }
    

    注意.handle(reactiveMongoDbStoringMessageHandler) - 这与MongoDb.reactiveOutboundChannelAdapter() 无关。因为这个将ReactiveMessageHandler 包装成ReactiveMessageHandlerAdapter 用于自动订阅。您需要的是看起来更像是您希望将 Mono&lt;Void&gt; 返回到您自己的控制中,因此您可以将其用作您的 writeToKafka 服务的输入,并自己在那里订阅并按照您的解释处理成功或错误。关键是,使用 Reactive Stream,我们无法提供命令式错误处理。该方法与任何异步 API 使用相同。因此,对于 Reactive Streams,我们也将错误发送到 errorChannel

    我们可能可以使用returnMono(true/false) 之类的东西改进MongoDb.reactiveOutboundChannelAdapter(),让像您这样的用例开箱即用。

    【讨论】:

    • 非常感谢。如果我理解正确,如果我想自己控制反应式订阅,我目前没有理由使用实现 ReactiveMessageHandler 的适配器。所以是的,如果有一个returnMono 选项可以改善这一点。
    • 对于一个案例,我想将reactiveKafkaConsumerTemplate 移动到不同的类中。推荐的方法是实现一个ReactorKafkaInboundChannelAdapter?
    • 是的。这是可能的。请参阅 MongoDbChangeStreamMessageProducer 作为示例,了解如何为响应式源实现消息驱动的通道适配器。
    • 您能解释一下是什么使它成为通道适配器吗?我在阅读文档后认为通道适配器必须实现SourcePollingChannelAdapter 或由@InboundChannelAdapter 注释。我只看到 MongoDbChangeStreamMessageProducer 实现了 MessageProducerSupport 所以我对 documentation 很困惑。
    • 正确。这就是为什么我指出 MongoDbChangeStreamMessageProducer 具有类似的反应特性。因此,您可以完全借用它为准备Flux 所做的一切以及它最终如何调用subscribeToPublisher(changeStreamFlux);。与ReactiveKafkaMessageProducer 相同。当您对解决方案感到满意时,请随时将其回馈给 Spring Integration:github.com/spring-projects/spring-integration/issues
    猜你喜欢
    • 2018-09-04
    • 2017-05-18
    • 1970-01-01
    • 2017-07-29
    • 2019-02-27
    • 2013-11-13
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多