架构

KeyValueRepository

师成师成· 更新于 2026-09-28· 阅读 22 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

KeyValueRepository

自 Camel 4.23 起

KeyValueRepository SPI(org.apache.camel.spi.KeyValueRepository)提供了一个统一的键值抽象,可由任意存储技术作为后端。无需为每种存储后端分别实现仓储接口,单一的 KeyValueRepository 实现即可作为多种 Camel 模式的后备存储:

这意味着你只需配置一个后端(例如 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
内存(默认)MemoryKeyValueRepositorycamel-support惰性逐出支持(computeIfPresent)
JDBCJdbcKeyValueRepositorycamel-sqlexpires_at 列不支持(默认)
JPAJpaKeyValueRepositorycamel-jpaexpiresAt 字段不支持(默认)
CassandraCassandraKeyValueRepositorycamel-cassandraql原生 USING TTL不支持(默认)
KafkaKafkaKeyValueRepositorycamel-kafka消息封装时间戳不支持(默认)
CaffeineCaffeineKeyValueRepositorycamel-caffeine原生按条目 Expiry支持(computeIfPresent)
EhcacheEhcacheKeyValueRepositorycamel-ehcache惰性逐出不支持(默认)
JCache(JSR-107)JCacheKeyValueRepositorycamel-jcache惰性逐出不支持(默认)
HazelcastHazelcastKeyValueRepositorycamel-hazelcast原生按条目支持(IMap.replace)
RedisRedisKeyValueRepositorycamel-redis原生键过期部分支持(CAS 支持,CAS+TTL 非原子)
InfinispanInfinispanRemoteKeyValueRepositorycamel-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");  // optional

TTL 使用 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),以获得更好的并发性能。

评论

登录后参与评论

正在加载评论…