Java 中的 Hazelcast:面向 Spring Boot 应用的分布式内存数据网格
Java 中 Hazelcast 的实操、代码优先之旅——这是一张内存数据网格,同时充当 Spring Boot 应用的缓存、协调层与分布式计算引擎。覆盖 IMap、EntryProcessor、MapStore、Near Cache,以及何时该使用它。
大多数缓存是你通过网络对话的一个盒子。Hazelcast 就是那张网络。这种区别就是全部要点——也是这篇文章的全部主张。
Hazelcast 究竟是什么
大多数团队接触缓存,是因为数据库读流量太慢。他们装上 Redis,给几个方法标上 @Cacheable,仪表盘就变绿了。在你需要的不只是缓存之前,这套很好用。在你需要一个两个服务必须达成共识的计数器之前。一条节点重启时不能丢的队列。一种当一项长任务在改某条记录时跨集群对它加锁的方法。需要这些任何一项的那一刻,简单的缓存就不够了。
Hazelcast 就是当你拿起"缓存"二字然后问"如果缓存就是集群呢?"得到的答案。它是一个内存数据网格(IMDG)——一组对等的 JVM,共享数据、协同执行代码,并悄无声息地处理分区、复制和故障转移。你得到分布式 Map、分布式 Queue、分布式 Lock、一个把任务跑在数据已经存在的节点上的 ExecutorService,以及让大部分功能看起来像普通缓存的 Spring Boot 集成。
这是一次小旅行。我们将从"把它嵌入 Spring Boot 应用并当缓存用",一直走到"把它当作分布式系统的协调层"。先上代码,意见在它配得上的位置出现。
嵌入式 vs 客户端-服务端——尽早做出选择
Hazelcast 有两种部署形态,选择会染色后续一切。
嵌入式模式。Hazelcast 跑在你应用的 JVM 里。你服务的每个实例同时也是一个 Hazelcast 成员,它们通过网络自动组成集群。共驻意味着数据查找通常是本地内存访问——没有网络跳。代价:你不能独立于应用扩展数据网格。服务从 3 节点扩到 30 节点,网格也跟着扩。
客户端-服务端模式。Hazelcast 作为一组专属 JVM 的集群运行。你的应用通过一个轻客户端与它通信。这更接近团队使用 Redis 的方式。你可以单独调整数据层规模,重启应用节点而不丢数据,并让一个多语言栈共用一个网格。
对大多数发布单一 Java 服务的 Spring Boot 团队来说,嵌入式是最简单的入口。对多语言栈,或者运维希望把缓存作为独立层来管理的场景,客户端-服务端才是合理选择。
在 Spring Boot 中搭建 Hazelcast
Spring Boot 内建支持——加上依赖就能在 classpath 上得到一个嵌入式集群。
<!-- pom.xml -->
<dependency>
<groupId>com.hazelcast</groupId>
<artifactId>hazelcast-spring</artifactId>
<version>5.4.0</version>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-cache</artifactId>
</dependency>
一个最小配置——你可以在 classpath 放一个 hazelcast.yaml,或者用 Java 构造配置。
// src/main/resources/hazelcast.yaml
hazelcast:
cluster-name: orders-cluster
network:
port:
auto-increment: true
port: 5701
join:
multicast:
enabled: false
tcp-ip:
enabled: true
member-list:
- 10.0.0.10
- 10.0.0.11
- 10.0.0.12
map:
products:
time-to-live-seconds: 300
max-size:
policy: PER_NODE
max-size: 10000
eviction:
eviction-policy: LRU
或者用 Java 配置:
@Configuration
public class HazelcastConfig {
@Bean
public Config hazelcastConfig() {
Config config = new Config();
config.setClusterName("orders-cluster");
MapConfig productsMap = new MapConfig("products")
.setTimeToLiveSeconds(300)
.setEvictionConfig(new EvictionConfig()
.setEvictionPolicy(EvictionPolicy.LRU)
.setMaxSizePolicy(MaxSizePolicy.PER_NODE)
.setSize(10_000));
config.addMapConfig(productsMap);
return config;
}
}
在你的应用类上加 @EnableCaching,Spring 就会自动用 Hazelcast 来支撑 @Cacheable。
@SpringBootApplication
@EnableCaching
public class OrdersApplication {
public static void main(String[] args) {
SpringApplication.run(OrdersApplication.class, args);
}
}
你拿到的分布式数据结构
Hazelcast 的价值来自这些积木。它们看上去像熟悉的 Java API,但状态在集群之间共享,按规模分区,按容错复制。下面是你常会用到的。
IMap——分布式 HashMap
IMap 是主力。它看起来像 ConcurrentHashMap,但它的键在集群间分区,每个分区被复制到一个或多个备份节点。
@Service
public class ProductCatalog {
private final HazelcastInstance hazelcast;
public ProductCatalog(HazelcastInstance hazelcast) {
this.hazelcast = hazelcast;
}
public Product get(String productId) {
IMap<String, Product> map = hazelcast.getMap("products");
return map.get(productId);
}
public void put(String productId, Product product) {
IMap<String, Product> map = hazelcast.getMap("products");
map.put(productId, product);
}
public Product getOrLoad(String productId) {
IMap<String, Product> map = hazelcast.getMap("products");
return map.computeIfAbsent(productId, this::loadFromDatabase);
}
private Product loadFromDatabase(String productId) {
// hits the DB
return null;
}
}
IMap 还免费给你谓词、索引和条目监听器。你可以 map.values(Predicates.equal("category", "books")),Hazelcast 会把谓词并行下推到每个分区。
IQueue, ITopic——集群级消息
需要一个节点重启后还在、由任何活着的成员消费的队列?IQueue 是一个分布式 BlockingQueue。
IQueue<OrderEvent> queue = hazelcast.getQueue("order-events");
queue.put(new OrderEvent(orderId, "PLACED")); // producer
OrderEvent next = queue.take(); // consumer (any node)
需要广播——每个节点都收到每条消息——用 ITopic:
ITopic<CacheInvalidate> topic = hazelcast.getTopic("cache-invalidations");
topic.addMessageListener(message -> {
String key = message.getMessageObject().getKey();
localCache.remove(key);
});
topic.publish(new CacheInvalidate("product:123"));
IAtomicLong, FencedLock——协调原语
两个服务需要就下一个序列号达成一致。或者其中一个需要在长时间操作期间对一条记录上独占锁。Hazelcast 给你和单 JVM 一样的原语,但作用范围是整个集群。
// Cluster-wide counter
IAtomicLong sequence = hazelcast.getCPSubsystem().getAtomicLong("invoice-seq");
long next = sequence.incrementAndGet();
// Cluster-wide lock (with timeout — never use the unbounded version in production)
FencedLock lock = hazelcast.getCPSubsystem().getLock("order:" + orderId);
if (lock.tryLock(2, TimeUnit.SECONDS)) {
try {
// do the thing only one cluster member should be doing
} finally {
lock.unlock();
}
}
CPSubsystem 使用 Raft 共识算法——这些原语在网络分区下也是正确的,代价是需要法定人数。把它们用在正确性要紧的地方;同名 API 的 AP 版本更快但最终一致。
MultiMap 和 ReplicatedMap
MultiMap 在一个键下保存多个值——做标签索引或一对多查询很有用。ReplicatedMap 把整个数据集放到每个节点上——非常适合极小、极频繁读取的参考数据,因为每次读都是一次本地内存访问。
MultiMap<String, String> tagsByPost = hazelcast.getMultiMap("post-tags");
tagsByPost.put("post-1", "java");
tagsByPost.put("post-1", "hazelcast");
Collection<String> tags = tagsByPost.get("post-1"); // [java, hazelcast]
ReplicatedMap<String, String> countryByCode = hazelcast.getReplicatedMap("countries");
countryByCode.put("US", "United States"); // copied to every node
由 Hazelcast 支撑的 Spring @Cacheable
大多数团队最先用上的,也是最简单的:给方法加注解,让 Spring 干活,把 Hazelcast 当成可即插即用的缓存。
@Service
public class ProductService {
private final ProductRepository repository;
public ProductService(ProductRepository repository) {
this.repository = repository;
}
@Cacheable(value = "products", key = "#productId")
public Product getProduct(String productId) {
log.info("Loading product {} from DB", productId);
return repository.findById(productId)
.orElseThrow(() -> new ProductNotFoundException(productId));
}
@CachePut(value = "products", key = "#product.id")
public Product updateProduct(Product product) {
return repository.save(product);
}
@CacheEvict(value = "products", key = "#productId")
public void deleteProduct(String productId) {
repository.deleteById(productId);
}
@CacheEvict(value = "products", allEntries = true)
public void clearAll() {
// wipes the products cache cluster-wide
}
}
和 Redis 支撑的缓存有什么不同:当一个节点更新了商品,其他每个节点立刻看到,因为底层的 IMap 就是同一个分布式对象。"跨服务缓存失效"这件事根本不需要解决——缓存本身就是共享的。
Near Cache——读为主流量时
即便是一个内存网格,当数据落在另一个分区时也要付一点小代价。对极热的键——一张配置表、一张特性开关表——网络跳累积起来就显著了。Near Cache 用本地保留最近访问条目的副本来解决这件事,并由网格在源变更时推送失效。
NearCacheConfig nearCacheConfig = new NearCacheConfig()
.setName("default")
.setInMemoryFormat(InMemoryFormat.OBJECT)
.setMaxIdleSeconds(60)
.setEvictionConfig(new EvictionConfig()
.setSize(1000)
.setEvictionPolicy(EvictionPolicy.LRU));
MapConfig featureFlags = new MapConfig("feature-flags")
.setNearCacheConfig(nearCacheConfig);
第一次命中之后,读直接来自本地内存。陈旧数据风险被失效推送和 max-idle 限定。"多旧才算太旧"是产品问题,不是技术问题——根据数据如何运转选择,而不是框架默认。
MapStore——对数据库的读穿透与写穿透
有时你希望网格就是访问层。应用调用 map.get(id),如果条目不在内存里,Hazelcast 从数据库取回,返回,然后缓存。map.put(id, value) 时 Hazelcast 替你写库。这就是 MapStore。
public class ProductMapStore implements MapStore<String, Product> {
private final JdbcTemplate jdbc;
public ProductMapStore(JdbcTemplate jdbc) {
this.jdbc = jdbc;
}
@Override
public Product load(String key) {
return jdbc.queryForObject(
"SELECT id, name, price FROM products WHERE id = ?",
new ProductRowMapper(),
key);
}
@Override
public Map<String, Product> loadAll(Collection<String> keys) {
// batch load
return jdbc.query(
"SELECT id, name, price FROM products WHERE id IN (...)",
new ProductRowMapper())
.stream()
.collect(Collectors.toMap(Product::getId, p -> p));
}
@Override
public void store(String key, Product value) {
jdbc.update(
"INSERT INTO products(id, name, price) VALUES(?, ?, ?) " +
"ON CONFLICT (id) DO UPDATE SET name = ?, price = ?",
value.getId(), value.getName(), value.getPrice(),
value.getName(), value.getPrice());
}
@Override
public void delete(String key) {
jdbc.update("DELETE FROM products WHERE id = ?", key);
}
}
在配置里接上:
MapStoreConfig mapStoreConfig = new MapStoreConfig()
.setEnabled(true)
.setImplementation(productMapStore)
.setWriteDelaySeconds(0); // 0 = synchronous write-through; >0 = write-behind batched
config.getMapConfig("products").setMapStoreConfig(mapStoreConfig);
写后批量(writeDelaySeconds > 0)为吞吐量批量写。它也意味着节点崩溃可能丢掉尚未刷盘的写——明确的取舍,把它写在团队看得见的地方。
EntryProcessor——更新就发生在数据所在处
更新条目的自然方式是"取、改、放"。在分布式系统里这是两次网络往返,加上一个两个客户端可能互相覆盖更改的窗口。EntryProcessor 把代码送到拥有数据的分区,并在那里原子地执行。
public class IncrementStockProcessor implements EntryProcessor<String, Product, Integer> {
private final int delta;
public IncrementStockProcessor(int delta) { this.delta = delta; }
@Override
public Integer process(Map.Entry<String, Product> entry) {
Product product = entry.getValue();
if (product == null) return 0;
product.setStock(product.getStock() + delta);
entry.setValue(product);
return product.getStock();
}
}
// Usage
IMap<String, Product> products = hazelcast.getMap("products");
Integer newStock = (Integer) products.executeOnKey("sku-42", new IncrementStockProcessor(-1));
一次往返,没有竞态,没有手动加锁。要做聚合——"对每个匹配此谓词的条目把这个计数器自增"——用 executeOnEntries(processor, predicate),Hazelcast 会在所有分区上并行运行处理器。
用 IExecutorService 做分布式计算
网格不仅存数据——还跑代码。IExecutorService 是一个集群感知的 ExecutorService,能把任务提交给某个特定成员、所有成员,或者最有用的,"拥有这个键的成员",让任务跑在数据所在处,而不是穿过网络。
IExecutorService executor = hazelcast.getExecutorService("default");
// Run on a specific key's owner — task arrives where the data already lives
Future<Integer> future = executor.submitToKeyOwner(new InventoryCheck("sku-42"), "sku-42");
// Or run on every member and aggregate
Map<Member, Future<Long>> results = executor.submitToAllMembers(new RowCountTask());
long total = results.values().stream()
.mapToLong(this::getQuietly)
.sum();
任务类需要在每个成员的 classpath 上。对于需要扫描大量数据的临时分析或定时工作,这种模式比把数据全部拉回到一个节点再集中处理要好。
生产注意事项
集群发现。多播在笔记本上能用,在生产里几乎哪都不能用。用带显式成员列表的 TCP/IP,云上部署用 Kubernetes / AWS / Azure 发现插件,让成员通过平台自身的服务发现彼此找到。
备份。默认下 IMap 把每个分区的一个同步备份保存在另一个成员上。对于不能丢的数据,把 backup-count 提到 2——每个分区便落在三个节点上,可以丢两个不丢数据。代价是更多内存和稍慢一点的写。
脑裂。如果网络分区把集群切成两半,两边都会继续接受写。连接恢复后 Hazelcast 会检测到分裂,跑合并策略来对账——明确选一个(LATEST_UPDATE、HIGHER_HITS 或自定义),不要靠默认。对于在分区间必须一致的数据,使用 CP subsystem。
内存管理。堆外(HD-Memory)让你把数据存在 Java 堆之外,让 GC 没什么可扫。对多 GB 的缓存来说,这就是"灵敏的集群"和"每分钟暂停一秒的集群"的差别。
可观测性。Hazelcast 直接对外发布 JMX 指标,Management Center 给你一个集群健康的 UI。把两者都接进监控;不要从客户那里第一次听说慢分区。
什么时候伸手去拿 Hazelcast
当问题不只是"让数据库更快"时,Hazelcast 在工具箱里赚到位置。下列任一情况可考虑:
- 一个同时承担短暂状态(会话、领导选举状态、在途工作流状态)真相来源的缓存。
- 跨集群协调——分布式锁、原子计数器、消息广播——而不必另起一套 ZooKeeper 或 etcd。
- 对大到不便搬运的数据进行计算——把函数送到数据,而不是把数据搬到函数。
- 每个应用实例同时是网格成员的嵌入式网格,并且你希望亚毫秒级读取,不必再运维一个独立缓存层。
如果你只需要键值缓存,且团队已经在用 Redis,Hazelcast 可能比问题需要的网格大太多。反过来也成立:多年前选了 Redis 的团队,如今为了弥补 Redis 的不足而到处洒上分布式锁、队列和 Lua 脚本——这种团队往往面对一个 Hazelcast 形状的问题,正在用难走的路解。
最后一句
选与问题形状相配的工具。第一次用 Hazelcast,是因为官方教程友好,嵌入式模式让我少运维一件东西。第二次用,是因为有三个服务在一个共享工作流上协调,我已经不想再在 HTTP 上发明协议了。这是一个诚实的范式:一个长成协调层的缓存,正好就是 IMDG 当初被设计出来要做的事。
常见问题
Hazelcast 只是 Redis 的对手吗?
在缓存上有重叠,但它属于另一类工具。Redis 是一个有丰富数据类型、单线程内核的远程键值服务器。Hazelcast 是一个分布式内存网格,作为你 JVM 的一部分(或独立集群)运行,支持并行计算,并在集群范围内提供 Java 风格的并发原语。仅做缓存它们看起来相似;做分布式协调、与数据共驻的计算,以及嵌入式场景,它们并不等价。
我应该用嵌入式还是客户端-服务端模式?
嵌入式在运维上更简单——你的服务和网格一起扩、一起部署、一起失败。客户端-服务端在以下情况更合适:运维想单独调整数据层、应用重启不能等同缓存重启、网格被多语言服务共用。多数单 Java 服务团队从嵌入式开始,只有当上述约束出现时才转客户端-服务端。
Hazelcast 与 JCache 或 Caffeine 有什么不同?
Caffeine 是一个快的本地缓存——单 JVM、无复制、无集群感知。JCache(JSR 107)是一个标准缓存 API,多个产品(包括 Hazelcast 和 Caffeine)都实现它。Hazelcast 是底下的分布式系统,带成员管理、分区、复制等等。问题如果在一个 JVM 内能解决,Caffeine 更快更简单。如果数据需要跨实例共享,Hazelcast 是更大问题的更大答案。
用于流处理的 Hazelcast Jet 呢?
Jet 已经被并入核心产品。承载 IMap 的同一批节点也能跑一条流处理流水线,从 Kafka 主题读取、与 map 做 join、写入另一张 IMap——全部不出网格。对低延迟状态查询的事件驱动富化场景,Jet 是 JVM 生态里更优雅的答案之一。
Hazelcast 能熬过我整个集群的重启吗?
开箱即用,数据在内存里,整集群重启会丢。要熬过整集群重启,请配置 Persistence(原 Hot Restart Store)——条目镜像到每个成员的本地磁盘,这样集群回来后成员重新载入分区并带着数据加入。对完全不能丢数据的负载,再叠加一个把变更写穿到数据库的 MapStore——本地持久化用于快速恢复,数据库作为持久的真相来源。