跳到主要内容

Spring 集成交互

ChatGPT-4o-mini 中英对照 Spring Integration Spring Integration Interaction

Spring Integration Framework 扩展了 Spring 编程模型,以支持众所周知的企业集成模式。它在基于 Spring 的应用程序中实现轻量级消息传递,并通过声明性适配器支持与外部系统的集成。它还提供了一个高级 DSL,用于将各种操作(端点)组合成一个逻辑集成流程。通过这种 DSL 配置的 lambda 风格,Spring Integration 已经很好地采用了 java.util.function 接口。@MessagingGateway 代理接口也可以作为 FunctionConsumer,根据 Spring Cloud Function 环境可以注册到函数目录中。有关其对函数支持的更多信息,请参见 Spring Integration ReferenceManual

另一方面,从版本 4.0.3 开始,Spring Cloud Function 引入了一个 spring-cloud-function-integration 模块,该模块提供了更深入、更具云特性和基于自动配置的 API,用于从 Spring Integration DSL 视角与 FunctionCatalog 进行交互。FunctionFlowBuilder 被自动配置并自动注入 FunctionCatalog,并且代表了针对目标 IntegrationFlow 实例的函数特定 DSL 的入口点。除了标准的 IntegrationFlow.from() 工厂(为了方便),FunctionFlowBuilder 还暴露了一个 fromSupplier(String supplierDefinition) 工厂,用于在提供的 FunctionCatalog 中查找目标 Supplier。然后,这个 FunctionFlowBuilder 引导到 FunctionFlowDefinition。这个 FunctionFlowDefinitionIntegrationFlowExtension 的一个实现,并暴露了 apply(String functionDefinition)accept(String consumerDefinition) 操作符,分别用于从 FunctionCatalog 中查找 FunctionConsumer。有关更多信息,请参见它们的 Javadocs。

以下示例演示了 FunctionFlowBuilder 的实际应用,以及 IntegrationFlow API 其余部分的强大功能:

@Configuration
public class IntegrationConfiguration {

@Bean
Supplier<byte[]> simpleByteArraySupplier() {
return "simple test data"::getBytes;
}

@Bean
Function<String, String> upperCaseFunction() {
return String::toUpperCase;
}

@Bean
BlockingQueue<String> results() {
return new LinkedBlockingQueue<>();
}

@Bean
Consumer<String> simpleStringConsumer(BlockingQueue<String> results) {
return results::add;
}

@Bean
QueueChannel wireTapChannel() {
return new QueueChannel();
}

@Bean
IntegrationFlow someFunctionFlow(FunctionFlowBuilder functionFlowBuilder) {
return functionFlowBuilder
.fromSupplier("simpleByteArraySupplier")
.wireTap("wireTapChannel")
.apply("upperCaseFunction")
.log(LoggingHandler.Level.WARN)
.accept("simpleStringConsumer");
}

}
java

由于 FunctionCatalog.lookup() 功能不仅限于简单的函数名称,还可以在提到的 apply()accept() 操作符中使用函数组合特性:

@Bean
IntegrationFlow functionCompositionFlow(FunctionFlowBuilder functionFlowBuilder) {
return functionFlowBuilder
.from("functionCompositionInput")
.accept("upperCaseFunction|simpleStringConsumer");
}
java

当我们在 Spring Cloud 应用程序中添加预定义功能的自动配置依赖项时,这个 API 变得更加相关。例如 Stream Applications 项目,除了应用程序镜像外,还提供用于各种集成用例的功能工件,例如 debezium-supplierelasticsearch-consumeraggregator-function 等。

以下配置分别基于 http-supplierspel-functionfile-consumer

@Bean
IntegrationFlow someFunctionFlow(FunctionFlowBuilder functionFlowBuilder) {
return functionFlowBuilder
.fromSupplier("httpSupplier", e -> e.poller(Pollers.trigger(new OnlyOnceTrigger())))
.<Flux<?>>handle((fluxPayload, headers) -> fluxPayload, e -> e.async(true))
.channel(c -> c.flux())
.apply("spelFunction")
.<String, String>transform(String::toUpperCase)
.accept("fileConsumer");
}
java

我们需要做的就是将他们的配置添加到 application.properties 中(如果有必要):

http.path-pattern=/testPath
spel.function.expression=new String(payload)
file.consumer.name=test-data.txt
properties