流缓存
流缓存
虽然流类型(如 StreamSource、InputStream 和 Reader)出于性能考虑在消息传递中被广泛使用,但它们也有一个重要的缺点:只能被读取一次。为了能够多次处理消息内容,就需要对流进行缓存。
流会被缓存到内存中。不过,对于大型流消息,你可以设置 spoolEnabled=true,这样超过 128 KB 的大消息就会被缓存到临时文件中。一旦缓存的流不再需要,Camel 会自行负责删除该临时文件。
为什么我的消息是空的?
在 Camel 中,消息体可以是任何类型。有些类型可以安全地被多次读取,因此不会遭遇变成空的问题。
因此,当你的消息体突然变为空时,通常是因为使用了不可重复读取的消息类型;换句话说,消息体只能被读取一次,后续的读取得到的就是空内容。基于流的类型(例如 java.util.InputStream 等)就会出现这种情况。
许多 Camel 组件开箱即用地支持并使用流式类型,例如 HTTP 相关组件、CXF 等。
Camel 提供了流缓存功能,会对流进行缓存,使其可以被重复读取。
StreamCache —— 会影响消息负载
StreamCache 会影响你的负载对象,因为它会将 Stream 负载替换为 org.apache.camel.StreamCache 对象。这个 StreamCache 可以被重复读取,因此在 Camel 中使用重新投递或基于内容的路由器等方式进行路由时会更加可靠。
为了判断消息负载是否需要缓存,Camel 使用类型转换器功能,来确定消息负载的类型是否可以转换为 org.apache.camel.StreamCache 实例。
Camel 发行版中所有实现了 org.apache.camel.StreamCache 的类都不是供最终用户创建实例使用的,它们是 Camel 流缓存功能的一部分。 |
|---|
配置流缓存
流缓存通过 org.apache.camel.spi.StreamCachingStrategy 进行配置。
该策略具有以下选项:
| 选项 | 默认值 | 说明 |
|---|---|---|
allowClasses | 用于过滤给定的允许/拒绝类集合的流缓存。默认情况下,所有属于 java.io.InputStream 的类都被允许。多个类名可用逗号分隔。 | |
anySpoolRules | false | 决定是否需要所有 SpoolRule 返回 true(还是任意一个返回 true)来确定流是否应被写入临时文件。这相当于对所有规则应用 AND/OR 二进制逻辑。默认基于 AND。 |
bufferSize | 4096 | 设置分配内存流缓存所用的内存缓冲区时的缓冲区大小。 |
denyClasses | 用于过滤给定的允许/拒绝类集合的流缓存。默认情况下,所有属于 java.io.InputStream 的类都被允许。多个类名可用逗号分隔。 | |
enabled | true | 是否启用流缓存 |
removeSpoolDirectoryWhenStopping | true | 在停止 CamelContext 时是否删除临时文件目录。 |
spoolCipher | null | 如果设置,临时文件将使用指定的密码变换进行加密(即有效的流密码或 8 位密码名称,例如 "RC4"、"AES/CTR/NoPadding"。空名称 "" 视为 null)。 |
spoolDirectory | ${java.io.tmpdir}/camel/camel-tmp-#uuid# | 存放写入临时文件的流所对应的临时文件的基础目录。此选项支持如下所述的命名模式。 |
spoolEnabled | false | 是否启用写入磁盘 |
spoolThreshold | 128 KB | 当流的大小达到该字节数时,应将其写入磁盘而不是保留在内存中。使用 0 或负值可完全禁用此功能,这样无论流的大小如何都始终保留在内存中。 |
spoolUsedHeapMemoryLimit | Max | 如果 spoolUsedHeapMemoryThreshold 正在使用,则用于确定已使用堆内存的上限是 Max(最大值)还是 Committed(已提交值)。 |
spoolUsedHeapMemoryThreshold | 0 | 以当前已使用堆内存的百分比(1 到 99)作为将流写入磁盘的阈值。上界基于堆的已提交内存(JVM 可以保证获取的内存)。当内存不足时,可据此将流写入磁盘。 |
statisticsEnabled | false | 是否启用利用率统计。启用后,你可以通过 JMX 等方式查看这些统计数据。 |
SpoolDirectory 命名模式
支持以下模式:
#uuid#= 一个随机的 UUID#camelId#= CamelContext 的 id(例如名称)#name#= 与#camelId#相同#counter#= 一个递增的计数器#bundleId#= OSGi bundle id(仅适用于 OSGi 环境)#symbolicName#= OSGi 符号名称(仅适用于 OSGi 环境)#version#= OSGi bundle 版本(仅适用于 OSGi 环境)${env:key}= 键对应的环境变量${key}= 键对应的 JVM 系统属性
几个示例:
要存储在 Java 临时目录下的一个以 CamelContext 名称命名的子目录中:
context.getStreamCachingStrategy().setSpoolDirectory"${java.io.tmpdir}#name#/");存储到 KARAF_HOME/tmp/bundleId 目录中:
context.getStreamCachingStrategy().setSpoolDirectory"${env:KARAF_HOME}/tmp/bundle#bundleId#");在 Java 中配置 StreamCachingStrategy
你可以在 Java 中按照如下方式配置 StreamCachingStrategy:
context.getStreamCachingStrategy().setSpoolEnabled(true);
context.getStreamCachingStrategy().setSpoolDirectory("/tmp/cachedir");
context.getStreamCachingStrategy().setSpoolThreshold(64 * 1024);
context.getStreamCachingStrategy().setBufferSize(16 * 1024);
// to enable encryption using RC4
// context.getStreamCachingStrategy().setSpoolCipher("RC4");并且记得在 CamelContext 上启用 Stream caching:
context.setStreamCaching(true);或应用于路由:
- Java
- XML
- YAML
from("file:inbox")
.streamCache(true)
.to("bean:foo");<route streamCache="true">
<from uri="file:inbox"/>
<to uri="bean:foo"/>
</route>- route:
streamCache: "true"
from:
uri: file:inbox
steps:
- to:
uri: bean:foo使用应用属性配置流缓存
使用 Spring Boot、Quarkus 或 Camel 独立版时,建议在 application.properties 配置文件中配置流缓存:
应用属性
camel.main.streamCachingSpoolEnabled = true
camel.main.streamCachingSpoolDirectory = /tmp/cachedir
camel.main.streamCachingSpoolThreshold = 65536
camel.main.streamCachingBufferSize = 16384你可以在命令行中运行 Camel CLI 命令 camel doc main --filter=stream 来查看所有选项。 |
|---|
在 Spring XML 中配置 StreamCachingStrategy
在 Spring XML 中,你可以先在 <camelContext> 上启用流缓存,然后在 streamCaching 元素中进行配置:
<camelContext streamCache="true">
<streamCaching id="myCacheConfig" bufferSize="16384" spoolEnabled="true" spoolDirectory="/tmp/cachedir" spoolThreshold="65536"/>
<route>
<from uri="direct:c"/>
<to uri="mock:c"/>
</route>
</camelContext>使用 spoolUsedHeapMemoryThreshold
默认情况下,流缓存只会将较大的负载(128 KB 及以上)写入磁盘。不过,你也可以设置 spoolUsedHeapMemoryThreshold 选项,它表示已使用堆内存的百分比。当内存即将不足时,可以利用该选项将数据一并写入磁盘。
例如,可以使用:
- Spring XML
- Application Properties
<streamCaching id="myCacheConfig" spoolEnabled="true" spoolDirectory="/tmp/cachedir" spoolUsedHeapMemoryThreshold="70"/>camel.main.streamCachingSpoolEnabled = true
camel.main.streamCachingSpoolDirectory = /tmp/cachedir
camel.main.streamCachingSpoolUsedHeapMemoryThreshold = 70然后要注意,spoolThreshold 默认已启用且为 128 KB,因此两个阈值(spoolThreshold 和 spoolUsedHeapMemoryThreshold)都在生效。在这个示例中,只有当负载大于 128 KB 并且已用堆内存大于 70% 时,我们才会将数据溢写到磁盘。原因是选项 anySpoolRules 默认为 false,这意味着两条规则必须同时为 true(即 AND 关系)。
如果我们希望只要满足其中任意一条规则就溢写到磁盘(即 OR 关系),则可以这样配置:
- Spring XML
- Application Properties
<streamCaching id="myCacheConfig" spoolEnabled="true" spoolDirectory="/tmp/cachedir" spoolUsedHeapMemoryThreshold="70" anySpoolRules="true"/>camel.main.streamCachingSpoolEnabled = true
camel.main.streamCachingSpoolDirectory = /tmp/cachedir
camel.main.streamCachingSpoolUsedHeapMemoryThreshold = 70
camel.main.streamCachingAnySpoolRules = true如果我们只希望在内存不足时才将数据写入磁盘,可以设置:
- Spring XML
- 应用程序属性
<streamCaching id="myCacheConfig" spoolEnabled="true" spoolDirectory="/tmp/cachedir" spoolThreshold="-1" spoolUsedHeapMemoryThreshold="70"/>camel.main.streamCachingSpoolEnabled = true
camel.main.streamCachingSpoolDirectory = /tmp/cachedir
camel.main.streamCachingSpoolThreshold = -1
camel.main.streamCachingSpoolUsedHeapMemoryThreshold = 70那么我们将不使用 spoolThreshold 规则,只使用基于堆内存的规则。
默认情况下,已用堆内存的上限是基于最大堆内存大小的。不过你也可以配置为使用已提交堆内存大小作为上限,这可以通过 spoolUsedHeapMemoryLimit 选项来完成,如下所示:
- Spring XML
- 应用属性(Application Properties)
<streamCaching id="myCacheConfig" spoolEnabled="true" spoolDirectory="/tmp/cachedir" spoolUsedHeapMemoryThreshold="70" spoolUsedHeapMemoryLimit="Committed"/>camel.main.streamCachingSpoolEnabled = true
camel.main.streamCachingSpoolDirectory = /tmp/cachedir
camel.main.streamCachingSpoolUsedHeapMemoryThreshold = 70
camel.main.streamCachingSpoolUsedHeapMemoryLimit = Committed使用自定义的 SpoolRule 实现(高级)
你可以实现自己的规则,以决定是否将流假脱机(spool)到磁盘。这可以通过实现 org.apache.camel.spi.StreamCachingStrategy.SpoolRule 接口来完成,该接口只有一个方法:
boolean shouldSpoolCache(long length);length 表示流的长度。要使用该规则,请按如下方式将其添加到 StreamCachingStrategy 中:
- Java
- Spring XML
- 应用属性
SpoolRule mySpoolRule = ...
context.getStreamCachingStrategy().addSpoolRule(mySpoolRule);而在 Spring XML 中,你需要为自定义规则定义一个 <bean>:
<bean id="mySpoolRule" class="com.foo.MySpoolRule"/>
<streamCaching id="myCacheConfig" spoolEnabled="true" spoolDirectory="/tmp/cachedir" spoolRules="mySpoolRule"/>在 <streamCaching> 上使用 spoolRules 属性。如果有多个规则,则用逗号分隔。
<streamCaching id="myCacheConfig" spoolEnabled="true" spoolDirectory="/tmp/cachedir" spoolRules="mySpoolRule,myOtherSpoolRule"/>使用 Spring Boot 或 Camel 独立运行模式时,你也可以在 application.properties 中进行配置:
# refers to the bean id of the rule object
camel.main.streamCachingSpoolRules=mySpoolRule
# you can also specify the class via
camel.main.streamCachingSpoolRules=#class:com.foo.MySpoolRule为每个 Exchange 使用自定义 spool 目录(高级)
默认情况下,所有 spool 的流都会写入同一个 spool 目录。如果你想控制每个 Exchange 的 spool 文件写入到何处(例如,按路由隔离 spool 数据),可以实现一个自定义的 StreamCachingStrategy,并重写 resolveSpoolDirectory 方法。
当某个流即将被 spool 到磁盘时,系统会为每个 Exchange 调用一次 resolveSpoolDirectory 方法。默认实现返回已配置的 spoolDirectory。自定义实现可以根据 Exchange 返回不同的目录,例如使用路由 ID 作为子目录:
DefaultStreamCachingStrategy strategy = new DefaultStreamCachingStrategy() {
@Override
public File resolveSpoolDirectory(Exchange exchange) {
String routeId = exchange.getFromRouteId();
if (routeId != null) {
return new File(getSpoolDirectory(), routeId);
}
return getSpoolDirectory();
}
};
context.setStreamCachingStrategy(strategy);如果解析出的目录不存在,将自动创建。
使用 StreamCachingProcessor
从 Camel 4.11 开始,可以使用该处理器将当前消息体转换为 StreamCache。这样消息体就可以被多次重新读取,并且可以放置在 Camel 路由中的任意位置。
from("direct:start")
.process(new StreamCachingProcessor())
.to("log:cached");如何在 Camel 中进行调试日志记录时启用流
当 Camel 以 DEBUG 级别运行时,它会不时地记录消息及其内容。由于某些消息可能包含流,而流往往无法被多次读取,因此 Camel 默认不会记录这些类型。
以下是默认不会被记录的典型类型:
java.xml.transform.StreamSourcejava.io.InputStreamjava.io.OutputStreamjava.io.Readerjava.io.Writer
你将在日志中看到如下内容:
DEBUG ProducerCache - >>>> Endpoint[direct:start] Exchange[Message: [Body is instance of java.xml.transform.StreamSource]]这里的消息是基于 XML 流的。你可以自行定制 Camel 是否仍然记录该负载。
你可以通过 Java 将其作为 CamelContext 上的全局选项启用:
context.getGlobalOptions().put(Exchange.LOG_DEBUG_BODY_STREAMS, "true");在 application.properties 中也可以按如下方式配置:
camel.main.globalOptions[CamelLogDebugBodyStreams] = true注意,默认值为 false。
评论
登录后参与评论
KnowForge