本文目录导读:

在分布式系统中,单调一致性(Monotonic Consistency)通常指单调读一致性(Monotonic Read Consistency),它的核心规则是:一旦一个进程读取到某个数据的某个版本(或更新后的值),它后续的所有读操作都不应再看到该数据更旧的值(即时间上的回退)。
换句话说,如果你先读到了 v2(新值),系统保证你后续绝对不会再读到 v1(旧值),这解决了“时光倒流”的问题,但不需要所有节点立刻一致(弱于强一致性)。
下面详细解释如何实现单调一致性,以及它在 Java 分布式系统中的典型实现方式。
实现单调一致性的核心思路
单调一致性的实现本质上是记录并维护一个“版本锚点”,客户端在每次读操作后,记住当前读到数据的最新版本号(或时间戳),后续读请求中,客户端将这个版本号作为参数传递给服务端,要求服务端只返回大于等于该版本的数据。
常见做法:
- 客户端携带读时间戳:客户端在请求中携带
last_read_timestamp,服务端确保返回的数据的版本 >= 该时间戳。 - 服务端重定向或等待:如果服务端节点当前的数据版本低于客户端的要求(说明该节点较旧),它必须等待其他节点同步(如等待主节点的新数据),或把请求转发给拥有更新数据的节点(如主节点或最新的副本)。
Java 分布式系统实现举例
方案 A:基于数据库的版本号 + 客户端锚点
假设数据存储在 MySQL 或 TiDB 中,每行数据有一个乐观锁字段(如 version)。
客户端代码逻辑:
public class MonotonicReaderService {
private final DataRepository repository;
// 客户端本地记录上次读取的版本号(线程私有或存储在请求上下文)
private volatile long lastReadVersion = 0;
// 读取数据,保证单调性
public DataRecord readWithMonotonicity(Long id) {
long clientVersion = lastReadVersion;
// 查询时,要求返回的版本号 >= 客户端已知的最新版本
DataRecord record = repository.findByIdAndVersionGreaterThanOrEqual(id, clientVersion);
if (record != null) {
// 更新本地的版本号
lastReadVersion = record.getVersion();
} else {
// 如果数据库还没有大于等于客户端版本的数据,说明读的是旧节点或数据未同步
// 需要重试或等待(例如等待主节点同步)
throw new RetryLaterException("Data version too old, retry...");
}
return record;
}
// 写入数据(由其他服务完成)
public void writeData(Long id, String value) {
// 插入时 version = currentTimestamp
repository.saveWithTimestamp(id, value);
}
}
方案 B:基于 Cassandra 或 DynamoDB 的 Read-after-write
在 Cassandra 中,可以通过设置 Read Consistency Level 和 Write Consistency Level 实现单调读。
- Write 采用
QUORUM(写多数节点) - Read 采用
QUORUM(读多数节点)
QUORUM 无法完全保证单调读(因为读可能命中不同的节点集合),Cassandra 提供了 SERIAL 级别(需要用到线性一致性),但开销大。
更好的方案是客户端使用 monotonic clock:
import com.datastax.oss.driver.api.core.CqlSession;
import com.datastax.oss.driver.api.core.cql.*;
import java.time.Instant;
import java.util.concurrent.atomic.AtomicLong;
public class CassandraMonotonicReader {
private final CqlSession session;
private final AtomicLong lastReadTimestamp = new AtomicLong(0);
public Row readMonotonic(String key) {
long clientTs = lastReadTimestamp.get();
// 查询时显式指定时间戳(Cassandra 支持 USING TIMESTAMP)
// 注意:Cassandra 的 TIMESTAMP 用于冲突解决,而不是直接过滤
// 更实际的做法:查询返回后自己判断版本
ResultSet rs = session.execute(
SimpleStatement.builder("SELECT * FROM my_table WHERE key = ?")
.addPositionalValue(key)
.build()
);
Row row = rs.one();
if (row != null) {
long serverVersion = row.getLong("version");
if (serverVersion < clientTs) {
// 数据回退,需要重试或读取其他副本
throw new StaleDataException("Caught stale data, retry on different node");
}
// 成功读取,更新客户端状态
lastReadTimestamp.set(serverVersion);
}
return row;
}
}
更高级的实现:使用分布式共识算法
如果想在最严格的意义上实现单调读(如线性一致性下的单调读),通常需要引入分布式共识。
- ZooKeeper:每次读操作可以从 Follower 读,但必须带上
zxid,Follower 的zxid小于客户端的lastZxid,则请求必须转发给 Leader 或等待 Follower 追上。 - Etcd:客户端维护一个
cluster_id和revision,Etcd 保证每个节点的revision单调递增,客户端请求时携带revision,Etcd 会确保返回的响应revision>= 客户端传入的。
Etcd 客户端 Java 代码示例:
import io.etcd.jetcd.*;
public class EtcdMonotonicReader {
private final Client client;
private long lastRevision = 0;
public EtcdMonotonicReader(String endpoints) {
this.client = Client.builder().endpoints(endpoints).build();
}
public ByteSequence read(String key) throws Exception {
// 创建读取请求,携带客户端的已知最新版本
Kv kvClient = client.getKVClient();
CompletableFuture<GetResponse> future = kvClient.get(
ByteSequence.from(key, UTF_8),
GetOption.newBuilder()
.withRevision(lastRevision) // 关键:只返回 >= 该 revision 的数据
.build()
);
GetResponse response = future.get();
if (response.getCount() > 0) {
long kvRevision = response.getKvs().get(0).getModRevision();
// 更新本地已知版本
lastRevision = kvRevision;
return response.getKvs().get(0).getValue();
}
return null;
}
}
分布式缓存场景(Redis Cluster + Redisson)
Redis 本身不提供单调性保证(因为主从异步复制),但可以结合 Redisson 的 RedLock 或 RAtomicLong 实现单调版本控制。
思路:
- 每次写操作在 Redis 中设置一个 version(或时间戳)
- 每次读操作从 Redis 读取,并检查版本是否 >= 本地缓存版本
- 如果小于,则强制从主节点(或多数节点)重新读
public class RedisMonotonicCache {
private final RedissonClient redisson;
private final ConcurrentHashMap<String, AtomicLong> localVersions = new ConcurrentHashMap<>();
public String getWithMonotonicity(String key) {
AtomicLong localVersion = localVersions.computeIfAbsent(key, k -> new AtomicLong(0));
long expectedVersion = localVersion.get();
// 读取 Redis(默认可能读到旧副本)
RBucket<String> bucket = redisson.getBucket(key);
String value = bucket.get();
// Redis 返回的版本低于期望,说明读到旧副本,需要重试或读主节点
long redisVersion = extractVersion(value);
if (redisVersion < expectedVersion) {
// 方式1:强制从主节点读(需要配置读写分离策略)
// 方式2:等待后重试
value = readFromMaster(key);
redisVersion = extractVersion(value);
}
localVersion.set(Math.max(redisVersion, expectedVersion));
return value;
}
}
实现单调一致性的常用“配方”
| 方法 | 适用场景 | 实现复杂度 | 性能影响 |
|---|---|---|---|
| 客户端携带版本号+服务端过滤 | 简单数据库 | 低 | 低 |
| Cassandra一致性级别+重试 | NoSQL | 中 | 中 |
| ZooKeeper/Etcd 的 zxid/revision | 共识系统 | 高 | 中 |
| Redis 主从 + 版本检查 | 缓存 | 中 | 低 |
关键实践:
- 客户端必须是有状态的:需要记录上次读到的最新版本。
- 服务端必须返回版本信息:无论是时间戳、逻辑时钟、还是递增 ID。
- 处理“版本回退”:当读到旧版本时,不能返回给用户,必须等待或重试。
如果你有具体的技术栈(Spring Cloud + Redis / TiDB / Cassandra),我可以给出更贴合的 Java 代码示例。