KeyValueRepository
KeyValueRepository
自 Camel 4.23 起
KeyValueRepository SPI(org.apache.camel.spi.KeyValueRepository)提供了一个统一的键值抽象,可由任意存储技术作为后端。无需为每种存储后端分别实现仓储接口,单一的 KeyValueRepository 实现即可作为多种 Camel 模式的后备存储:
- 幂等消费者 —— 通过
KeyValueIdempotentRepository - 聚合器 —— 通过
KeyValueAggregationRepository - Cache EIP —— 直接将
KeyValueRepository用作缓存后端 - State Store 组件 —— 在路由中进行键值操作
这意味着你只需配置一个后端(例如 Redis),即可在所有四种模式中复用,无需重复配置或管理多个仓储实例。
API
KeyValueRepository 接口继承自 Service,并定义了以下操作:
| 方法 | 说明 |
|---|---|
get(key) | 获取指定键的值(不存在或已过期时返回 null) |
put(key, value, ttl) | 存储值并可选指定 TTL;返回被替换的旧值 |
delete(key) | 删除并返回指定键的值 |
contains(key) | 检查是否存在未过期的条目 |
keys() | 返回所有未过期的键 |
clear() | 删除所有条目 |
size() | 统计未过期的条目数量 |
putIfAbsent(key, value, ttl) | 仅在键不存在时存储;返回已存在的值或 null |
replace(key, expected, new, ttl) | 比较并交换:当当前值与期望值匹配时进行替换 |
delete(key, expected) | 比较并交换:仅当当前值与期望值匹配时才删除 |
TTL 通过 java.time.Duration 指定。null、零或负值时长表示永不过期。
putIfAbsent、replace 和 delete(key, expected) 方法带有默认实现,这些实现并非原子操作。支持原生原子操作的后端会覆盖这些方法,以获得更好的并发保证——参见下表中的原子性一列。
可用后端
Camel 开箱即用地提供了以下 KeyValueRepository 实现:
| 后端 | 类 | 模块 | TTL | 原子 CAS |
|---|---|---|---|---|
| 内存(默认) | MemoryKeyValueRepository | camel-support | 惰性逐出 | 支持(computeIfPresent) |
| JDBC | JdbcKeyValueRepository | camel-sql | expires_at 列 | 不支持(默认) |
| JPA | JpaKeyValueRepository | camel-jpa | expiresAt 字段 | 不支持(默认) |
| Cassandra | CassandraKeyValueRepository | camel-cassandraql | 原生 USING TTL | 不支持(默认) |
| Kafka | KafkaKeyValueRepository | camel-kafka | 消息封装时间戳 | 不支持(默认) |
| Caffeine | CaffeineKeyValueRepository | camel-caffeine | 原生按条目 Expiry | 支持(computeIfPresent) |
| Ehcache | EhcacheKeyValueRepository | camel-ehcache | 惰性逐出 | 不支持(默认) |
| JCache(JSR-107) | JCacheKeyValueRepository | camel-jcache | 惰性逐出 | 不支持(默认) |
| Hazelcast | HazelcastKeyValueRepository | camel-hazelcast | 原生按条目 | 支持(IMap.replace) |
| Redis | RedisKeyValueRepository | camel-redis | 原生键过期 | 部分支持(CAS 支持,CAS+TTL 非原子) |
| Infinispan | InfinispanRemoteKeyValueRepository | camel-infinispan | 原生存活时间 | 支持(完全原子) |
内存(默认)
默认后端将条目存储在 ConcurrentHashMap 中,并采用惰性 TTL 逐出。无需额外依赖。
KeyValueRepository repo = new MemoryKeyValueRepository();JDBC
通过原生 JDBC 将条目存储在关系数据库表中。
JdbcKeyValueRepository repo = new JdbcKeyValueRepository();
repo.setDataSource(myDataSource);
repo.setTableName("camel_kvr"); // optional, defaults to "camel_kvr"该表在启动时若不存在则会自动创建。TTL(生存时间)以 expires_at epoch-millis 列的形式存储;过期条目在读取时会被惰性清理。
JPA
使用带有 EntityManager 的 JPA 来存储条目。
JpaKeyValueRepository repo = new JpaKeyValueRepository();
repo.setEntityManagerFactory(myEmf);使用映射到 camel_kvr 表的 KeyValueEntry 实体。请将该实体添加到 persistence.xml 持久化单元中。
Cassandra
使用 DataStax 驱动将条目存储到 Apache Cassandra 中。
CassandraKeyValueRepository repo = new CassandraKeyValueRepository();
repo.setSession(cassandraSession);
repo.setTableName("camel_kvr"); // optionalTTL 使用 Cassandra 原生的 USING TTL(由 Duration 转换为秒)。该表会在启动时自动创建。
Kafka
使用 Kafka 的压实(compacted)主题作为持久化键值存储。
KafkaKeyValueRepository repo = new KafkaKeyValueRepository();
repo.setBootstrapServers("localhost:9092");
repo.setTopic("camel-kvr"); // optional条目在启动时从主题回放至内存中的 ConcurrentHashMap 缓存。TTL 以 epoch 毫秒时间戳的形式存储在值信封中。
Caffeine
基于进程内的缓存,通过 Caffeine 的 Expiry API 原生支持按条目设置 TTL。
CaffeineKeyValueRepository repo = new CaffeineKeyValueRepository();
repo.setMaximumSize(10000); // optional, defaults to 10,000通过 cache.asMap().computeIfPresent() 实现原子 CAS 操作。
Ehcache
使用 Ehcache 3 的 CacheManager。
EhcacheKeyValueRepository repo = new EhcacheKeyValueRepository();
repo.setCacheManager(myCacheManager);
repo.setCacheName("camel-kvr");TTL 通过一个包装器进行管理,采用读取时的延迟清除策略(Ehcache 3 本身不支持逐条目的 TTL)。
JCache(JSR-107)
可与任意 JCache 提供方(Hazelcast、Ehcache 等)配合使用。
JCacheKeyValueRepository repo = new JCacheKeyValueRepository();
repo.setCachingProvider(myProvider);
repo.setCacheName("camel-kvr");Hazelcast
使用 Hazelcast IMap 的分布式键值存储。
HazelcastKeyValueRepository repo = new HazelcastKeyValueRepository();
repo.setHazelcastInstance(myHzInstance);
repo.setMapName("camel-kvr"); // optional支持原生的按条目 TTL,以及通过 IMap.replace(K, V, V) 实现的原子 CAS 操作。
Redis
基于 Redisson 客户端实现。
RedisKeyValueRepository repo = new RedisKeyValueRepository("localhost:6379");
repo.setKeyPrefix("camel-kvr:"); // optional, defaults to "camel-kvr:"通过 RBucket.compareAndSet() 实现原子 CAS。注意:使用带 TTL 的 replace() 时,CAS 和 TTL 会作为两次独立的 Redis 调用执行——其间会有一个短暂的窗口,新值已存在但尚未设置 TTL。
Infinispan
由 Infinispan HotRod RemoteCacheManager 提供支持。
InfinispanRemoteKeyValueRepository repo = new InfinispanRemoteKeyValueRepository();
repo.setHost("localhost:11222");
repo.setCacheName("camel-kvr"); // optional完全原子的 CAS 操作,包括在单次原生调用中带 TTL 的 replace。
序列化
持久化后端(JDBC、JPA、Cassandra、Kafka、Hazelcast、Redis 和 Infinispan)通过共享的 KeyValueRepositoryHelper 以普通 Java 序列化产生的字节形式存储值。因此,值必须实现 java.io.Serializable。
由于这些后端背后的字节存储是共享基础设施——一张数据库表、一个 Kafka 主题、一个 Redis 实例、一个数据网格——读取回来的条目可能是 Camel 并未写入的数据的反序列化。因此,每次读取都会安装一个 JEP-290 的 java.io.ObjectInputFilter,与聚合仓库(aggregation repositories)的做法完全一致。
默认情况下,该过滤器拒绝 java.net.,其余的 java.、javax. 和 org.apache.camel. 则予以放行,同时应用 JEP-290 的图结构限制 maxdepth、maxrefs 和 maxbytes 作为纵深防御。如果设置了 JVM 全局的 jdk.serialFilter 系统属性,它将优先于 Camel 的默认值。
这意味着存储你自己的类的实例时,需要通过 deserializationFilter 选项放宽过滤器,而每个持久化后端都提供该选项:
RedisKeyValueRepository repo = new RedisKeyValueRepository("localhost:6379");
repo.setDeserializationFilter("com.mycompany.model.**;java.**;javax.**;org.apache.camel.**;!*");该模式使用的语法与 jdk.serialFilter 相同。请尽可能保持其范围狭窄:像 * 这样宽松的模式会让仓库重新对任意 gadget 类开放。
当通过 KeyValueAggregationRepository 使用持久化后端时,这一点同样间接适用:聚合后的 Exchange 会以 DefaultExchangeHolder 形式存储,因此消息体以及任何可序列化的头部类型也必须被过滤器覆盖。
内存后端(Memory、Caffeine、Ehcache、JCache)保存的是对象引用,不会进行序列化,因此它们没有 deserializationFilter 选项。
与 Camel 模式配合使用
单个后端,多个模式
在 Camel 注册表中注册一个 KeyValueRepository,所有模式都会自动发现它:
@BindToRegistry("kvRepo")
public KeyValueRepository kvRepo() {
return new RedisKeyValueRepository("localhost:6379");
}只需在注册表中放入这一个 bean:
- State Store 组件会自动发现并将其用于键值操作。
- Idempotent Consumer EIP 会自动发现并将其包装为
KeyValueIdempotentRepository。 - Aggregator EIP 会自动发现并将其包装为
KeyValueAggregationRepository。 - Cache EIP 会将其用作缓存后端。
无需任何显式装配。
显式装配
当你需要更强的控制力时,也可以显式地装配该后端:
// Idempotent Consumer
KeyValueRepository store = new JdbcKeyValueRepository(myDataSource);
IdempotentRepository idempotent = new KeyValueIdempotentRepository(store);
// Aggregation Repository
AggregationRepository aggregation = new KeyValueAggregationRepository(store);XML / YAML DSL
<bean name="kvStore" type="org.apache.camel.component.redis.RedisKeyValueRepository">
<property name="endpoint" value="localhost:6379"/>
</bean>
<!-- Idempotent Consumer backed by Redis -->
<bean name="idempotentRepo" type="org.apache.camel.support.KeyValueIdempotentRepository">
<constructors>
<constructor value="#kvStore"/>
</constructors>
</bean>编写自定义后端
实现 KeyValueRepository 接口即可创建你自己的后端:
public class MyCustomRepository extends ServiceSupport implements KeyValueRepository {
@Override
public Object get(String key) { /* ... */ }
@Override
public Object put(String key, Object value, Duration ttl) { /* ... */ }
@Override
public Object delete(String key) { /* ... */ }
@Override
public boolean contains(String key) { /* ... */ }
@Override
public Set<String> keys() { /* ... */ }
@Override
public void clear() { /* ... */ }
}如果存储支持原生原子操作,请重写 putIfAbsent、replace 和 delete(key, expected),以获得更好的并发性能。
评论
登录后参与评论
KnowForge