Cassandra写后读出现陈旧读取问题,求Spring Boot及Spring Data Cassandra解决方案
Cassandra写后读场景陈旧读问题解决方案
一、陈旧数据识别核心方法
- 利用Cassandra原生
writetime内置函数:每行数据写入时会自动生成微秒级时间戳,读取时获取该字段与写入操作返回的时间戳比对,读取结果的writetime小于写入时间戳即可判定为陈旧数据 - 业务层新增版本号字段:每次写入操作版本号自增1,读取时将返回的版本号与业务侧记录的最新版本号比对,版本号偏小即为陈旧数据
- 轻量级事务(LWT)校验:写入时使用条件更新操作,成功后返回最新版本标识,读取时直接比对标识即可识别陈旧数据
二、避免返回陈旧数据的可选方案
- 调整读写一致性级别实现强一致:Cassandra实现强读一致性不存在高难度是常见误区,只要满足「读一致性级别 + 写一致性级别 > 副本数即可实现强一致。通常副本数设为3时,读写均采用
QUORUM级别即可满足强一致,多数据中心场景可使用LOCAL_QUORUM读写保证单中心强一致 - 业务侧缓存最新写入标记:写入成功后将最新版本号/写入时间戳存入本地缓存,缓存过期时间设置为大于等于Cassandra的最大同步超时时间即可,读取时先做校验
- 读请求触发手动读修复:识别到陈旧数据时主动触发读修复,等待同步完成后再返回结果
三、Spring Boot + Spring Data Cassandra 实现示例
基础配置(application.yml)
spring: cassandra: keyspace-name: your_keyspace contact-points: 127.0.0.1 port: 9042 request: # 普通读写一致性级别 consistency: QUORUM # 轻量级事务一致性级别,多中心用LOCAL_SERIAL serial-consistency: SERIAL
实体类定义
import org.springframework.data.cassandra.core.mapping.Column; import org.springframework.data.cassandra.core.mapping.PrimaryKey; import org.springframework.data.cassandra.core.mapping.Table; @Table("user_info") public class UserInfo { @PrimaryKey private String userId; @Column private String userName; // 业务版本号字段 @Column private Long version; // 省略getter、setter方法 }
DAO层接口
import org.springframework.data.cassandra.repository.CassandraRepository; import org.springframework.data.cassandra.repository.Query; import org.springframework.data.repository.query.Param; public interface UserInfoRepository extends CassandraRepository<UserInfo, String> { // 条件更新,版本号自增,返回写入是否成功 @Query("UPDATE user_info SET userName = :userName, version = :newVersion WHERE userId = :userId IF version = :oldVersion") Boolean updateUserInfo(@Param("userId") String userId, @Param("userName") String userName, @Param("oldVersion") Long oldVersion, @Param("newVersion") Long newVersion); // 查询数据同时返回版本号 @Query("SELECT userId, userName, version FROM user_info WHERE userId = :userId") UserInfo getUserInfo(@Param("userId") String userId); }
业务层逻辑实现
import com.google.common.cache.CacheBuilder; import com.google.common.cache.LoadingCache; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; import java.util.concurrent.TimeUnit; @Service public class UserInfoService { @Autowired private UserInfoRepository userInfoRepository; // 缓存最近写入的最新版本号,过期时间匹配Cassandra最大同步超时时间 private final LoadingCache<String, Long> latestVersionCache = CacheBuilder.newBuilder() .expireAfterWrite(3, TimeUnit.HOURS) .maximumSize(10000) .build(); // 写入逻辑 public void writeUserInfo(String userId, String userName) { Long oldVersion = latestVersionCache.getIfPresent(userId); oldVersion = oldVersion == null ? 0L : oldVersion; Long newVersion = oldVersion + 1; Boolean writeSuccess = userInfoRepository.updateUserInfo(userId, userName, oldVersion, newVersion); if (writeSuccess) { latestVersionCache.put(userId, newVersion); return; } throw new RuntimeException("写入冲突,请重试"); } // 读取逻辑 public UserInfo readUserInfo(String userId) { UserInfo result = userInfoRepository.getUserInfo(userId); Long cachedLatestVersion = latestVersionCache.getIfPresent(userId); // 版本比对识别陈旧数据 if (cachedLatestVersion != null && result.getVersion() < cachedLatestVersion) { // 可根据业务需要选择重试读取、触发读修复或者直接返回提示 throw new RuntimeException("数据同步中,请稍后重试"); } return result; } }
内容的提问来源于stack exchange,提问作者JUser
相关产品推荐
相关产品推荐

