【问题标题】:Using Spring Integration executorChannel with Spring Cloud Function将 Spring Integration executorChannel 与 Spring Cloud 功能一起使用
【发布时间】:2020-04-13 13:07:47
【问题描述】:

我正在使用 Spring Cloud 功能来处理来自带有 Flux 的 kafka 的数据。默认情况下,它在消费者线程(消费消息的地方)中处理数据。我将为并行数据处理和节流实现线程池,Spring Cloud Integration 中有一个很棒的实现,称为 executorChannel (https://docs.spring.io/spring-integration/api/org/springframework/integration/channel/ExecutorChannel.html)

功能实现示例:

public static class FN1 implements Function<Flux<String>, Flux<String>> {
  public Flux<String> apply(Flux<String> data) {
    return data
      .map(f ->  doSomething() )      
  }
}

所以我发现没有简单的方法来连接通过 executorChannel 实现的函数。

M.b.有没有办法定义 inputChannel 类型?

UPD:在 Oleg 的回答下阅读 cmets。它们非常有用。

【问题讨论】:

    标签: java spring-integration spring-cloud-function


    【解决方案1】:

    你的意思是这样的?

    @SpringBootApplication
    public class SampleFunctoinAppApplication  {
    
        public static void main(String[] args) throws Exception {
    
            ApplicationContext context = SpringApplication.run(SampleFunctoinAppApplication.class, args);
            SubscribableChannel output = context.getBean("output", SubscribableChannel.class);
            output.subscribe(System.out::println);
    
            MessageChannel channel = context.getBean("executorChannel", MessageChannel.class);
            channel.send(new GenericMessage<String>("hello"));
        }
    
        @Bean
        public IntegrationFlow flow() {
            return IntegrationFlows
                    .from("executorChannel")
                    .transform(echo())
                    .channel("output")
                    .get();
        }
    
        @Bean
        public ExecutorChannel executorChannel() {
            return new ExecutorChannel(Executors.newCachedThreadPool());
        }
    
        public Function<String, String> echo() {
            return v -> v;
        }
    }
    

    “定义 inputChannel 类型”是什么意思?

    【讨论】:

    • 是的,但 SCF 的主要功能之一是您可以仅实现一个或多个功能,并使用来自 swping 云流和 function.definition 的绑定器与外部配置相互绑定。在您的示例中,此定义是硬编码的。
    • "define inputChannel type" like this ``` @Bean @Transformer(inputChannel = "headerEnricherChannel", outputChannel = "transformChannel") public HeaderEnricher headerEnricher() { return new HeaderEnricher(Collections.singletonMap(SimpMessageHeaderAccessor .SESSION_ID_HEADER, new ExpressionEvaluatingHeaderValueMessageProcessor("payload", null))); } ```
    • 那么您说的是 spring-cloud-stream 和功能绑定 - cloud.spring.io/spring-cloud-static/spring-cloud-stream/…,对吗?
    • 好的。 1.我有几个实现的函数Function, Flux<...>>。 2. 我可以使用 Spring Cloud Streams binders 将此函数绑定到消息代理,只需配置它们。所以实现很简单。一切都很棒 3. 我想做节流代理 -> (direct channel) -> FN1 -> (executorChannel) -> FN2 我想保持 spring-cloud-functions 的灵活性并重用已经实施者 executorChannel 进行节流,因为它很棒解决方案。只需编写所有代码 - 好的,但是 m.b.还有更优雅的解决方案吗?
    • 你坚持拥有 executorChannel 却没有意识到框架已经处理了消费者并发 - cloud.spring.io/spring-cloud-static/spring-cloud-stream/…。请浏览文档。此外,您也不需要执行者通道。这只是众多实现细节之一。您有一个我认为需要并发的业务案例,因此请考虑说明您的案例,以便我们可以更好地为您提供春季投资组合中正确路径的建议。
    猜你喜欢
    • 1970-01-01
    • 2014-10-23
    • 2018-08-28
    • 2020-10-24
    • 1970-01-01
    • 2016-09-27
    • 1970-01-01
    • 2023-02-06
    • 2017-09-13
    相关资源
    最近更新 更多