编程模型
函数目录与灵活的函数签名
Spring Cloud Function 的主要特性之一是适应并支持用户定义函数的一系列类型签名,同时提供一致的执行模型。这就是为什么所有用户定义的函数都会通过 FunctionCatalog 转换为规范表示形式。
虽然用户通常不需要关心 FunctionCatalog,但了解用户代码中支持哪些类型的函数是很有用的。
同样重要的是,Spring Cloud Function 对由 Project Reactor 提供的响应式 API 提供了原生支持。这使得响应式原语如 Mono 和 Flux 可以作为用户定义函数中的类型使用,从而在为实现函数选择编程模型时提供了更大的灵活性。响应式编程模型还支持一些特性的功能性实现,这些特性在使用命令式编程风格时可能难以或无法实现。更多相关内容,请阅读 函数元数 部分。
Java 8 函数支持
Spring Cloud Function 接受并构建在自 Java 8 以来定义的 3 个核心函数式接口之上。
-
供应商<O>
-
函数<I, O>
-
消费者<I>
为了避免频繁提到 Supplier、Function 和 Consumer,在本手册的其余部分,我们将在适当的地方将它们统称为功能 Bean。
简而言之,ApplicationContext 中的任何 Functional bean 都将被懒加载注册到 FunctionCatalog 中。这意味着它可以受益于本参考手册中描述的所有附加功能。
在最简单的应用中,你只需要在应用配置中声明一个类型为 Supplier、Function 或 Consumer 的 @Bean。然后,你可以使用 FunctionCatalog 根据名称查找特定的函数。
例如:
@Bean
public Function<String, String> uppercase() {
return value -> value.toUpperCase();
}
// . . .
FunctionCatalog catalog = applicationContext.getBean(FunctionCatalog.class);
Function uppercase = catalog.lookup(“uppercase”);
重要的是要理解,既然 uppercase 是一个 bean,你当然可以直接从 ApplicationContext 中获取它,但你得到的只是你声明的 bean,而不包含 SCF 提供的任何额外功能。当你通过 FunctionCatalog 查找一个函数时,你接收到的实例是被包装(instrumented)过的,带有本手册中描述的额外功能(例如类型转换、组合等)。
此外,重要的是要理解,典型用户并不直接使用 Spring Cloud Function。相反,典型用户会实现一个 Java Function、Supplier 或 Consumer,并期望在不同的执行环境中使用它,而无需额外的工作。
例如,同一个 Java 函数可以通过 Spring Cloud Function 提供的适配器以及其他使用 Spring Cloud Function 作为核心编程模型的框架(如 Spring Cloud Stream)表示为 REST 端点、流式消息处理器 或 AWS Lambda,甚至更多形式。
总的来说,Spring Cloud Function 为 Java 函数提供了额外的功能,使其能够在多种执行环境中使用。
函数定义
前面的示例展示了如何以编程方式在 FunctionCatalog 中查找函数,但在典型的集成场景中,当 Spring Cloud Function 被另一个框架(如 Spring Cloud Stream)用作编程模型时,你可以通过 spring.cloud.function.definition 属性声明要使用的函数。了解并理解在 FunctionCatalog 中发现函数的默认行为非常重要。
例如,如果你的 ApplicationContext 中只有一个 Functional bean,通常不需要设置 spring.cloud.function.definition 属性,因为 FunctionCatalog 中的单个函数可以通过空名称或任意名称进行查找。例如,假设 uppercase 是你目录中唯一的函数,它可以通过 catalog.lookup(null)、catalog.lookup(“”) 或 catalog.lookup(“foo”) 进行查找。
也就是说,对于使用诸如 Spring Cloud Stream 这样的框架的情况,该框架使用 spring.cloud.function.definition,建议始终使用 spring.cloud.function.definition 属性。
例如,
spring.cloud.function.definition=uppercase
过滤不符合条件的函数
一个典型的 ApplicationContext 可能包含一些有效的 Java 函数,但它们并不打算作为注册到 FunctionCatalog 的候选者。这些 bean 可能来自其他项目的自动配置,或者是任何其他符合 Java 函数资格的 bean。
该框架提供了默认过滤机制,排除不应注册到 FunctionCatalog 的已知 bean。您还可以通过使用 spring.cloud.function.ineligible-definitions 属性提供一个以逗号分隔的 bean 定义名称列表,将其他 bean 添加到此列表中。
例如,
spring.cloud.function.ineligible-definitions=foo,bar
供应商
供应商可以是 响应式 的 - Supplier<Flux<T>> 或 命令式 的 - Supplier<T>。从调用的角度来看,这对于实现此类 Supplier 的人来说应该没有区别。
然而,当在框架(例如 Spring Cloud Stream)中使用时,Supplier,尤其是响应式的,通常用于表示流的源头。因此,它们会被调用一次以获取流(例如 Flux),消费者可以订阅该流。换句话说,这样的 Supplier 相当于一个无限流。
尽管相同的反应式供应商也可以表示一个有限的流(例如轮询 JDBC 数据的结果集)。在这些情况下,这些反应式供应商必须连接到底层框架的某种轮询机制。
为了协助解决这个问题,Spring Cloud Function 提供了一个标记注解 org.springframework.cloud.function.context.PollableBean,用于指示此类供应商生成的是有限流,并且可能需要再次轮询。然而,重要的是要理解,Spring Cloud Function 本身并不为这个注解提供任何行为。
此外,PollableBean 注解还暴露了一个 splittable 属性,用于指示生成的流需要进行拆分(参见 Splitter EIP)。
这是一个示例:
@PollableBean(splittable = true)
public Supplier<Flux<String>> someSupplier() {
return () -> {
String v1 = String.valueOf(System.nanoTime());
String v2 = String.valueOf(System.nanoTime());
String v3 = String.valueOf(System.nanoTime());
return Flux.just(v1, v2, v3);
};
}
函数
函数也可以以命令式或响应式的方式编写。然而,与 Supplier 和 Consumer 不同,除了理解在框架(如 Spring Cloud Stream)中使用时,响应式函数只会被调用一次以传递对流的引用(即 Flux 或 Mono),而命令式函数会在每个事件中被调用一次之外,实现者无需特别考虑其他因素。
public Function<String, String> uppercase() {
. . . .
}
BiFunction
如果您需要接收一些额外的数据(元数据)与有效负载一起,您可以始终声明您的函数签名以接收包含带有额外信息的头部映射的 Message。
public Function<Message<String>, String> uppercase() {
. . . .
}
为了使你的函数签名更轻量且更符合 POJO 风格,还有另一种方法。你可以使用 BiFunction。
public BiFunction<String, Map, String> uppercase() {
. . . .
}
由于 Message 只包含两个属性(payload 和 headers),而 BiFunction 需要两个输入参数,框架将自动识别此签名并从 Message 中提取 payload,将其作为第一个参数传递,并将 headers 的 Map 作为第二个参数传递。因此,你的函数不会与 Spring 的消息 API 耦合。请记住,BiFunction 需要严格的签名,其中第二个参数必须是 Map。同样的规则适用于 BiConsumer。
消费者
Consumer 有点特殊,因为它有一个 void 返回类型,这意味着它可能会阻塞,至少是有这种可能性。大多数情况下,你不需要编写 Consumer<Flux<?>>,但如果你确实需要这样做,请记得订阅输入的 Flux。
函数组合
函数组合(Function Composition)是一个允许将多个函数组合成一个的特性。其核心支持基于 Java 8 提供的 Function.andThen(..) 函数组合特性。然而,Spring Cloud Function 在此基础上提供了一些额外的功能。
声明式函数组合
此功能允许您在设置 spring.cloud.function.definition 属性时,使用 |(管道符)或 ,(逗号)分隔符以声明式的方式提供组合指令。
例如:
--spring.cloud.function.definition=uppercase|reverse
在这里,我们有效地提供了一个单一函数的定义,该函数本身是 uppercase 和 reverse 两个函数的组合。实际上,这也是属性名称为 definition 而非 name 的原因之一,因为一个函数的定义可以是多个命名函数的组合。如前所述,你可以使用 , 代替 |,例如 …definition=uppercase,reverse。
组合非函数
Spring Cloud Function 还支持将 Supplier 与 Consumer 或 Function 组合,以及将 Function 与 Consumer 组合。重要的是要理解这些定义的最终产物。将 Supplier 与 Function 组合仍然会得到 Supplier,而将 Supplier 与 Consumer 组合实际上会生成 Runnable。按照同样的逻辑,将 Function 与 Consumer 组合会得到 Consumer。
当然,你也不能组合那些不可组合的对象,比如 Consumer 和 Function、Consumer 和 Supplier 等。
函数路由与过滤
自 2.2 版本起,Spring Cloud Function 提供了一个路由功能,允许你调用一个单一的函数,该函数充当路由器,指向你希望调用的实际函数。这一功能在某些 FAAS 环境中非常有用,在这些环境中,维护多个函数的配置可能很繁琐,或者无法暴露多个函数。
RoutingFunction 在 FunctionCatalog 中注册的名称为 functionRouter。为了简单和一致性,你也可以引用 RoutingFunction.FUNCTION_NAME 常量。
该函数具有以下签名:
public class RoutingFunction implements Function<Object, Object> {
// . . .
}
路由指令可以通过多种方式进行传达。我们支持通过消息头、系统属性以及可插拔策略来提供指令。让我们来看一些细节。
消息路由回调
MessageRoutingCallback 是一种策略,用于协助确定路由到函数定义的名称。
public interface MessageRoutingCallback {
default String routingResult(Message<?> message) {
return (String) message.getHeaders().get(FunctionProperties.FUNCTION_DEFINITION);
}
}
你需要做的就是实现并注册一个 MessageRoutingCallback 作为 bean,以便被 RoutingFunction 接收。例如:
@Bean
public MessageRoutingCallback customRouter() {
return new MessageRoutingCallback() {
@Override
public String routingResult(Message<?> message) {
return (String) message.getHeaders().get(FunctionProperties.FUNCTION_DEFINITION);
}
};
}
在前面的示例中,你可以看到一个非常简单的 MessageRoutingCallback 实现,它从传入 Message 的 FunctionProperties.FUNCTION_DEFINITION Message 头中确定函数定义,并返回一个表示要调用的函数定义的 String 实例。
消息头
如果输入参数的类型是 Message<?>,你可以通过设置 spring.cloud.function.definition 或 spring.cloud.function.routing-expression 这两个 Message 头来传递路由指令。正如属性名称所示,spring.cloud.function.routing-expression 依赖于 Spring 表达式语言 (SpEL)。对于更静态的情况,你可以使用 spring.cloud.function.definition 头,它允许你提供一个单一函数的名称(例如 …definition=foo)或一个组合指令(例如 …definition=foo|bar|baz)。对于更动态的情况,你可以使用 spring.cloud.function.routing-expression 头,并提供一个 SpEL 表达式,该表达式应解析为函数的定义(如上所述)。
SpEL 评估上下文的根对象是实际的输入参数,因此在 Message<?> 的情况下,您可以构造一个既可以访问 payload 又可以访问 headers 的表达式(例如 spring.cloud.function.routing-expression=headers.function_name)。
SpEL 允许用户提供要执行的 Java 代码的字符串表示形式。鉴于 spring.cloud.function.routing-expression 可以通过消息头提供,这意味着设置此类表达式的功能可能会暴露给最终用户(例如在使用 web 模块时的 HTTP 头),这可能会导致一些问题(例如恶意代码)。为了管理这一点,所有通过消息头传递的表达式将仅在 SimpleEvaluationContext 中进行评估,该上下文功能有限,设计用于仅评估上下文对象(在我们的情况下是消息)。另一方面,通过属性或系统环境变量设置的所有表达式将在 StandardEvaluationContext 中进行评估,允许充分利用 Java 语言的灵活性。虽然通过系统/应用程序属性或环境变量设置表达式通常被认为是安全的,因为在正常情况下不会暴露给最终用户,但在某些情况下,通过 Spring Boot Actuator 端点(由其他 Spring 项目、第三方或最终用户创建的自定义实现提供)确实会暴露给最终用户可见性和更新系统、应用程序和环境变量的能力。此类端点必须使用行业标准的 Web 安全实践进行保护。Spring Cloud Function 不会暴露任何此类端点。
在特定的执行环境/模型中,适配器负责通过 Message 头来转换和传递 spring.cloud.function.definition 和/或 spring.cloud.function.routing-expression。例如,在使用 spring-cloud-function-web 时,你可以将 spring.cloud.function.definition 作为 HTTP 头提供,框架会将其与其他 HTTP 头一起作为消息头进行传播。
应用程序属性
路由指令也可以通过 spring.cloud.function.definition 或 spring.cloud.function.routing-expression 作为应用程序属性进行传递。上一节中描述的规则在这里同样适用。唯一的区别是你需要将这些指令作为应用程序属性提供(例如,--spring.cloud.function.definition=foo)。
需要注意的是,将 spring.cloud.function.definition 或 spring.cloud.function.routing-expression 作为消息头提供时,仅适用于命令式函数(例如 Function<Foo, Bar>)。也就是说,我们只能通过命令式函数进行每条消息的路由。对于响应式函数,我们无法进行每条消息的路由。因此,你只能通过应用程序属性提供路由指令。这完全取决于工作单元。在命令式函数中,工作单元是消息,因此我们可以基于这样的工作单元进行路由。而在响应式函数中,工作单元是整个流,因此我们只能根据通过应用程序属性提供的指令进行操作,并路由整个流。
路由指令的优先级顺序
鉴于我们有多种提供路由指令的机制,理解在同时使用多种机制时冲突解决的优先级是非常重要的。以下是优先级顺序:
-
MessageRoutingCallback(当函数为命令式时优先使用,无论是否定义了其他内容) -
消息头(如果函数为命令式且未提供
MessageRoutingCallback) -
应用程序属性(适用于任何函数)
不可路由的消息
如果目录中没有可用的 route-to 函数,你将收到一个异常提示。
在某些情况下,这种行为可能不是我们所期望的,你可能希望有一个“全能”类型的函数能够处理这类消息。为了实现这一点,框架提供了 org.springframework.cloud.function.context.DefaultMessageRoutingHandler 策略。你只需要将其注册为一个 bean。它的默认实现会简单地记录消息不可路由的事实,但会允许消息流继续而不会抛出异常,从而有效地丢弃不可路由的消息。如果你需要更复杂的处理,你只需要提供自己的策略实现并将其注册为一个 bean 即可。
@Bean
public DefaultMessageRoutingHandler defaultRoutingHandler() {
return new DefaultMessageRoutingHandler() {
@Override
public void accept(Message<?> message) {
// do something really cool
}
};
}
函数过滤
过滤是一种路由类型,其中只有两条路径 - “通过”或“丢弃”。在函数方面,这意味着您只希望在某个条件返回“true”时调用某个函数,否则您希望丢弃输入。
然而,在丢弃输入方面,关于它在您的应用程序上下文中可能意味着什么,有许多不同的解释。例如,您可能希望记录它,或者您可能希望维护一个被丢弃消息的计数器。您也可能选择什么都不做。
由于这些不同的路径,我们并没有提供一个通用的配置选项来处理被丢弃的消息。相反,我们建议简单地定义一个 Consumer,用来表示“丢弃”路径:
@Bean
public Consumer<?> devNull() {
// log, count, or whatever
}
现在你可以有一个路由表达式,它实际上只有两条路径,有效地成为一个过滤器。例如:
--spring.cloud.function.routing-expression=headers.contentType.toString().equals('text/plain') ? 'echo' : 'devNull'
所有不符合传递给 echo 函数条件的消息都将传递到 devNull,在那里你可以简单地不做任何处理。签名 Consumer<?> 还将确保不会尝试进行类型转换,从而几乎不会产生执行开销。
在处理反应式输入(例如 Publisher)时,路由指令必须仅通过函数属性提供。这是由于反应式函数的特性,它们仅被调用一次以传递一个 Publisher,其余部分由反应器处理,因此我们无法访问和/或依赖于通过单个值(例如 Message)传递的路由指令。
多路由器
默认情况下,框架将始终配置一个路由函数,如前面章节所述。然而,有时您可能需要多个路由函数。在这种情况下,除了现有的路由函数外,您还可以创建自己的 RoutingFunction bean 实例,只要您为其指定一个不同于 functionRouter 的名称即可。
你可以将 spring.cloud.function.routing-expression 或 spring.cloud.function.definition 作为键值对传递给 RoutingFunction。
这是一个简单的示例:
@Configuration
protected static class MultipleRouterConfiguration {
@Bean
RoutingFunction mySpecialRouter(FunctionCatalog functionCatalog, BeanFactory beanFactory, @Nullable MessageRoutingCallback routingCallback) {
Map<String, String> propertiesMap = new HashMap<>();
propertiesMap.put(FunctionProperties.PREFIX + ".routing-expression", "'reverse'");
return new RoutingFunction(functionCatalog, propertiesMap, new BeanFactoryResolver(beanFactory), routingCallback);
}
@Bean
public Function<String, String> reverse() {
return v -> new StringBuilder(v).reverse().toString();
}
@Bean
public Function<String, String> uppercase() {
return String::toUpperCase;
}
}
这是一个演示其工作原理的测试:
@Test
public void testMultipleRouters() {
System.setProperty(FunctionProperties.PREFIX + ".routing-expression", "'uppercase'");
FunctionCatalog functionCatalog = this.configureCatalog(MultipleRouterConfiguration.class);
Function function = functionCatalog.lookup(RoutingFunction.FUNCTION_NAME);
assertThat(function).isNotNull();
Message<String> message = MessageBuilder.withPayload("hello").build();
assertThat(function.apply(message)).isEqualTo("HELLO");
function = functionCatalog.lookup("mySpecialRouter");
assertThat(function).isNotNull();
message = MessageBuilder.withPayload("hello").build();
assertThat(function.apply(message)).isEqualTo("olleh");
}
输入/输出增强
在某些情况下,你可能需要修改或调整传入或传出的消息,并且希望保持代码的整洁,避免与非功能性关注点混杂在一起。你不希望这些操作出现在业务逻辑内部。
你总是可以通过函数组合来实现它。这种方法有几个好处:
-
它允许你将这个非功能性关注点隔离到一个单独的函数中,你可以将其与业务函数组合为一个函数定义。
-
它为你提供了完全的自由(和危险),可以在传入消息到达实际业务函数之前对其进行修改。
@Bean
public Function<Message<?>, Message<?>> enrich() {
return message -> MessageBuilder.fromMessage(message).setHeader("foo", "bar").build();
}
@Bean
public Function<Message<?>, Message<?>> myBusinessFunction() {
// do whatever
}
然后,通过提供以下函数定义来组合你的函数:enrich|myBusinessFunction。
虽然所描述的方法最为灵活,但也最为复杂。它要求您编写一些代码,然后将其作为一个 bean,或者手动注册为一个函数,之后才能像前面的示例中那样将其与业务函数组合使用。
但是,如果像前面的例子那样,您尝试进行的修改(丰富)是微不足道的呢?是否有更简单、更动态和可配置的机制来实现相同的效果?
从 3.1.3 版本开始,框架允许你提供 SpEL 表达式来丰富进入函数和从函数输出的单个消息头。让我们以一个测试为例来看一下。
@Test
public void testMixedInputOutputHeaderMapping() throws Exception {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
SampleFunctionConfiguration.class).web(WebApplicationType.NONE).run(
"--logging.level.org.springframework.cloud.function=DEBUG",
"--spring.main.lazy-initialization=true",
"--spring.cloud.function.configuration.split.output-header-mapping-expression.keyOut1='hello1'",
"--spring.cloud.function.configuration.split.output-header-mapping-expression.keyOut2=headers.contentType",
"--spring.cloud.function.configuration.split.input-header-mapping-expression.key1=headers.path.split('/')[0]",
"--spring.cloud.function.configuration.split.input-header-mapping-expression.key2=headers.path.split('/')[1]",
"--spring.cloud.function.configuration.split.input-header-mapping-expression.key3=headers.path")) {
FunctionCatalog functionCatalog = context.getBean(FunctionCatalog.class);
FunctionInvocationWrapper function = functionCatalog.lookup("split");
Message<byte[]> result = (Message<byte[]>) function.apply(MessageBuilder.withPayload("hello")
.setHeader(MessageHeaders.CONTENT_TYPE, "application/json")
.setHeader("path", "foo/bar/baz")
.build());
assertThat(result.getHeaders()).containsKey("keyOut1"));
assertThat(result.getHeaders().get("keyOut1")).isEqualTo("hello1");
assertThat(result.getHeaders()).containsKey("keyOut2"));
assertThat(result.getHeaders().get("keyOut2")).isEqualTo("application/json");
}
}
在这里,你可以看到名为 input-header-mapping-expression 和 output-header-mapping-expression 的属性,它们前面是函数的名称(即 split),后面是你想要设置的消息头键的名称以及作为 SpEL 表达式的值。第一个表达式(针对 keyOut1)是一个用单引号括起来的字面量 SpEL 表达式,实际上将 keyOut1 设置为值 hello1。keyOut2 则被设置为现有 contentType 头的值。
你也可以在输入头映射中观察到一些有趣的功能,我们实际上是在拆分现有头 path 的值,将 key1 和 key2 的单独值设置为基于索引的拆分元素的值。
如果由于某种原因提供的表达式评估失败,函数的执行将继续进行,就好像什么都没发生过一样。但是,您会在日志中看到 WARN 消息,通知您此情况。
o.s.c.f.context.catalog.InputEnricher : Failed while evaluating expression "hello1" on incoming message. . .
在你处理具有多个输入的函数时(下一节),你可以紧接在 input-header-mapping-expression 之后使用索引:
--spring.cloud.function.configuration.echo.input-header-mapping-expression[0].key1=‘hello1'
--spring.cloud.function.configuration.echo.input-header-mapping-expression[1].key2='hello2'
函数元数
有时候,数据流需要被分类和组织。例如,考虑一个经典的大数据用例,处理包含“订单”和“发票”的无组织数据,你希望每个数据都进入一个单独的数据存储。这时,函数元数(具有多个输入和输出的函数)支持就派上用场了。
让我们来看一个这样的函数示例:MessageRoutingCallback。
完整的实现细节可以在这里找到。
@Bean
public Function<Flux<Integer>, Tuple2<Flux<String>, Flux<String>>> organise() {
return flux -> ...;
}
鉴于 Project Reactor 是 SCF 的核心依赖项,我们正在使用其 Tuple 库。元组为我们提供了一个独特的优势,即能够传达基数和类型信息。这两者在 SCSt 的上下文中都极为重要。基数让我们知道需要创建多少个输入和输出绑定,并将它们绑定到函数的相应输入和输出上。对类型信息的了解确保了正确的类型转换。
此外,这也是绑定名称命名规则中“索引”部分的由来,因为在这个函数中,两个输出绑定的名称分别是 organise-out-0 和 organise-out-1。
目前,函数参数数量(arity)仅支持用于以复杂事件处理为中心的响应式函数(Function<TupleN<Flux<?>…>, TupleN<Flux<?>…>>),其中对事件汇合处的评估和计算通常需要查看事件流而非单个事件。
输入头传播
在典型场景中,输入消息头不会传播到输出,这是合理的,因为函数的输出可能是其他需要自己一组消息头的输入。然而,有时这种传播可能是必要的,因此 Spring Cloud Function 提供了几种机制来实现这一点。
首先,你可以始终手动复制消息头。例如,如果你有一个函数签名接收 Message 并返回 Message(即 Function<Message, Message>),你可以简单且有选择性地自行复制消息头。请记住,如果你的函数返回 Message,框架除了正确转换其有效负载外,不会对其进行任何操作。然而,这种方法可能显得有些繁琐,尤其是在你只想复制所有消息头的情况下。为了帮助处理这种情况,我们提供了一个简单的属性,允许你在希望传播输入消息头的函数上设置一个布尔标志。该属性是 copy-input-headers。
例如,假设你有以下配置:
@EnableAutoConfiguration
@Configuration
protected static class InputHeaderPropagationConfiguration {
@Bean
public Function<String, String> uppercase() {
return x -> x.toUpperCase();
}
}
如你所知,你仍然可以通过向它发送消息来调用此函数(框架将负责类型转换和有效负载提取)
只需将 spring.cloud.function.configuration.uppercase.copy-input-headers 设置为 true,以下断言也将成立
Function<Message<String>, Message<byte[]>> uppercase = catalog.lookup("uppercase", "application/json");
Message<byte[]> result = uppercase.apply(MessageBuilder.withPayload("bob").setHeader("foo", "bar").build());
assertThat(result.getHeaders()).containsKey("foo");
类型转换(内容类型协商)
内容类型协商是 Spring Cloud Function 的核心功能之一,它不仅允许将传入的数据转换为函数签名声明的类型,还可以在函数组合期间执行相同的转换,从而使原本无法通过类型组合的函数变得可组合。
为了更好地理解内容协商机制及其必要性,我们通过以下函数作为一个非常简单的用例来探讨:
@Bean
public Function<Person, String> personFunction {..}
前面示例中显示的函数期望一个 Person 对象作为参数,并生成一个 String 类型的输出。如果用 Person 类型调用此函数,则一切正常。但是,通常情况下,函数扮演着处理传入数据的角色,这些数据通常以原始格式(如 byte[]、JSON 字符串 等)传入。为了使框架能够成功地将传入数据作为参数传递给此函数,它必须以某种方式将传入数据转换为 Person 类型。
Spring Cloud Function 依赖于两个 Spring 原生的机制来实现这一点。
-
MessageConverter - 用于将传入的 Message 数据转换为函数声明的类型。
-
ConversionService - 用于将传入的非 Message 数据转换为函数声明的类型。
这意味着根据原始数据的类型(Message 或非 Message),Spring Cloud Function 将应用不同的机制。
在大多数情况下,当处理作为某些其他请求(例如 HTTP、消息传递等)的一部分而被调用的函数时,框架依赖于 MessageConverters,因为这些请求已经被转换为 Spring 的 Message。换句话说,框架会定位并应用适当的 MessageConverter。为了实现这一点,框架需要用户提供一些指示。其中一项指示已经由函数本身的签名(Person 类型)提供。因此,理论上这应该(并且在某些情况下确实)足够了。然而,对于大多数用例,为了选择合适的 MessageConverter,框架还需要一个额外的信息。这个缺失的信息就是 contentType 头。
此类标头通常作为消息的一部分出现,最初由创建该消息的相应适配器注入。例如,HTTP POST 请求的 content-type HTTP 标头会被复制到消息的 contentType 标头中。
对于不存在此类标头的情况,框架依赖于默认的内容类型 application/json。
内容类型与参数类型
如前所述,框架要选择合适的 MessageConverter,需要参数类型,以及可选的内容类型信息。选择合适 MessageConverter 的逻辑由参数解析器负责,这些解析器会在用户定义函数调用之前触发(此时框架已知实际的参数类型)。如果参数类型与当前负载的类型不匹配,框架会委托给预先配置的 MessageConverter 堆栈,以查看是否有任何一个转换器能够转换负载。
contentType 和参数类型的组合是框架通过定位适当的 MessageConverter 来确定消息是否可以转换为目标类型的机制。如果未找到合适的 MessageConverter,则会抛出异常,您可以通过添加自定义的 MessageConverter 来处理此异常(请参阅 [用户定义的消息转换器](#user-defined-message-converters))。
不要期望 Message 仅基于 contentType 被转换为其他类型。请记住,contentType 是对目标类型的补充。它是一个提示,MessageConverter 可能会考虑,也可能不会考虑。
消息转换器
MessageConverters 定义了两个方法:
Object fromMessage(Message<?> message, Class<?> targetClass);
Message<?> toMessage(Object payload, @Nullable MessageHeaders headers);
理解这些方法的契约及其使用方式非常重要,特别是在 Spring Cloud Stream 的上下文中。
fromMessage 方法将传入的 Message 转换为参数类型。Message 的有效载荷可以是任何类型,MessageConverter 的实际实现需要支持多种类型。
提供的消息转换器
如前所述,框架已经提供了一组 MessageConverters 来处理大多数常见用例。以下列表按优先级顺序描述了提供的 MessageConverters(使用第一个可用的 MessageConverter):
-
JsonMessageConverter:支持在contentType为application/json时,使用 Jackson(默认)或 Gson 库将Message的有效负载与 POJO 相互转换。该消息转换器还支持type参数(例如,application/json;type=foo.bar.Person)。这在开发函数时可能不知道类型的情况下非常有用,因此函数签名可能类似于Function<?, ?>或Function或Function<Object, Object>。换句话说,对于类型转换,我们通常从函数签名中推导类型。拥有 mime-type 参数允许你以更动态的方式传递类型信息。 -
ByteArrayMessageConverter:支持在contentType为application/octet-stream时,将Message的有效负载从byte[]转换为byte[]。它本质上是一个直通的转换器,主要用于向后兼容。 -
StringMessageConverter:支持在contentType为text/plain时,将任何类型转换为String。
当没有找到合适的转换器时,框架会抛出异常。当这种情况发生时,你应该检查你的代码和配置,确保你没有遗漏任何内容(即确保你通过绑定或头部提供了 contentType)。然而,最有可能的是,你遇到了一些不常见的情况(例如自定义的 contentType),而当前提供的 MessageConverters 栈不知道如何转换。如果是这种情况,你可以添加自定义的 MessageConverter。请参阅用户自定义的 Message Converters。
用户自定义的消息转换器
Spring Cloud Function 提供了一种机制来定义和注册额外的 MessageConverter。要使用它,实现 org.springframework.messaging.converter.MessageConverter,并将其配置为 @Bean。然后它会附加到现有的 MessageConverter 栈中。
需要注意的是,自定义的 MessageConverter 实现会被添加到现有栈的头部。因此,自定义的 MessageConverter 实现会优先于现有的转换器,这让你既可以覆盖现有的转换器,也可以添加新的转换器。
以下示例展示了如何创建一个消息转换器 bean 以支持名为 application/bar 的新内容类型:
@SpringBootApplication
public static class SinkApplication {
...
@Bean
public MessageConverter customMessageConverter() {
return new MyCustomMessageConverter();
}
}
public class MyCustomMessageConverter extends AbstractMessageConverter {
public MyCustomMessageConverter() {
super(new MimeType("application", "bar"));
}
@Override
protected boolean supports(Class<?> clazz) {
return (Bar.class.equals(clazz));
}
@Override
protected Object convertFromInternal(Message<?> message, Class<?> targetClass, Object conversionHint) {
Object payload = message.getPayload();
return (payload instanceof Bar ? payload : new Bar((byte[]) payload));
}
}
关于 JSON 选项的说明
在 Spring Cloud Function 中,我们支持 Jackson 和 Gson 机制来处理 JSON。为了方便您使用,我们将其抽象为 org.springframework.cloud.function.json.JsonMapper,它本身知道这两种机制,并将使用您选择的机制或遵循默认规则。默认规则如下:
-
无论哪个库在类路径上,都将使用该库作为机制。因此,如果你在类路径上有
com.fasterxml.jackson.*,则将使用 Jackson;如果你有com.google.code.gson,则将使用 Gson。 -
如果你同时拥有这两个库,那么 Gson 将是默认的选择,或者你可以设置
spring.cloud.function.preferred-json-mapper属性为gson或jackson中的一个值。
也就是说,类型转换通常对开发者是透明的。然而,鉴于 org.springframework.cloud.function.json.JsonMapper 也被注册为一个 bean,你可以在需要时轻松地将其注入到你的代码中。
Kotlin Lambda 支持
我们还提供了对 Kotlin lambda 的支持(自 v2.0 起)。考虑以下示例:
@Bean
open fun kotlinSupplier(): () -> String {
return { "Hello from Kotlin" }
}
@Bean
open fun kotlinFunction(): (String) -> String {
return { it.toUpperCase() }
}
@Bean
open fun kotlinConsumer(): (String) -> Unit {
return { println(it) }
}
上述内容展示了将 Kotlin lambda 配置为 Spring bean 的方式。每个 lambda 的签名映射到 Java 中对应的 Supplier、Function 和 Consumer,因此这些签名被框架支持/识别。尽管 Kotlin 到 Java 的映射机制超出了本文档的范围,但重要的是要理解,这里同样适用于“Java 8 函数支持”部分中概述的签名转换规则。
要启用 Kotlin 支持,您只需在类路径中添加 Kotlin SDK 库,这将触发适当的自动配置和支持类。
函数组件扫描
Spring Cloud Function 如果存在一个名为 functions 的包,将会扫描该包中 Function、Consumer 和 Supplier 的实现。使用此功能,你可以编写不依赖于 Spring 的函数——甚至不需要 @Component 注解。如果你想使用不同的包,可以设置 spring.cloud.function.scan.packages。你也可以使用 spring.cloud.function.scan.enabled=false 来完全关闭扫描功能。
数据脱敏
一个典型的应用程序通常包含多个级别的日志记录。某些云/无服务器平台可能会在记录的数据包中包含敏感信息,这些信息可能会被所有人看到。虽然检查正在记录的数据是每个开发者的责任,但由于日志记录来自框架本身,从 4.1 版本开始,我们引入了 JsonMasker 来初步帮助屏蔽 AWS Lambda 负载中的敏感数据。然而,JsonMasker 是通用的,适用于任何模块。目前,它仅适用于结构化数据,如 JSON。你只需要指定你想要屏蔽的键,它就会处理其余的事情。键应该在 META-INF/mask.keys 文件中指定。文件的格式非常简单,你可以通过逗号、换行符或两者来分隔多个键。
以下是此类文件内容的示例:
eventSourceARN
asdf1, SS
在这里,您可以看到定义了三个键。一旦存在这样的文件,JsonMasker 将使用它来屏蔽指定键的值。
以下是展示用法的示例代码:
private final static JsonMasker masker = JsonMasker.INSTANCE();
// . . .
logger.info("Received: " + masker.mask(new String(payload, StandardCharsets.UTF_8)));