跳转到内容

缓存一致性问题

草稿难度:中级#分布式系统

原文:CSDN(历史文章导入,当前状态为草稿)

缓存一致性问题的本质是在分布式系统中,同一数据在不同存储介质(缓存和数据库)之间存在的状态差异,以及这种差异如何影响系统的正确性和可用性。

缓存一致性问题主要源于以下四个原因.
数据冗余存储:同一数据同时存在于缓存和数据库中,形成了数据冗余
更新时序问题:在高并发环境下,对缓存和数据库的操作顺序难以保证
分布式环境的复杂性:网络延迟、部分失败等分布式特性使数据同步变得困难
缓存策略的权衡:缓存的失效策略、更新策略都会影响一致性

缓存系统之所以存在,是因为我们需要在速度较慢但容量大、持久的存储(数据库)之外,增加一层速度快但容量有限、易失的存储(缓存)。这种冗余设计本身就包含了潜在的一致性问题。

当同一数据需要同时写入两个不同的系统时,就面临着原子性问题:

  • 无法保证对缓存和数据库的操作是原子的
  • 任何一方的写入失败都会导致数据不一致
  • 即使都成功,由于执行时间差异,也会出现短暂的不一致窗口

在读写分离架构中,这个问题更加复杂:

  • 写操作通常路由到主库
  • 读操作可能路由到从库
  • 主从复制延迟进一步扩大了不一致窗口

考虑以下并发更新场景:
线程A读取数据库值X=1
线程B读取数据库值X=1
线程A计算X+1=2,更新数据库
线程A删除缓存
线程B计算X+1=2,更新数据库(实际上应该是X+2=3)
线程B删除缓存
新请求来临,缓存未命中,从数据库读取X=2并写入缓存
最终结果是X=2,而非正确的X=3,出现了数据不一致。

不同的缓存更新策略面临不同的时序挑战:

  • 如果删除缓存失败,会导致缓存中保留旧数据
  • 在删除缓存到重建缓存之间,可能有其他线程读取到旧数据
  • 如果删除缓存后,更新数据库前,有读请求,会用旧数据重建缓存
  • 更新数据库后,缓存中仍是旧数据

延迟双删策略试图解决这些时序问题:

  1. 先删除缓存
  2. 更新数据库
  3. 等待一段时间(大于读操作重建缓存的时间)
  4. 再次删除缓存
    这样设计是为了确保在数据库更新完成后,任何可能使用旧数据重建的缓存都会被清除。

在分布式系统中,网络分区是不可避免的:

  • 缓存服务器可能暂时无法访问
  • 数据库可能暂时不可用
  • 不同节点之间的通信可能延迟或中断
    这些情况都会导致缓存和数据库之间的同步失败,产生不一致。

分布式系统中的部分失败是常态:

  • 更新数据库成功但删除缓存失败
  • 删除缓存成功但更新数据库失败
  • 主从复制部分成功部分失败
    每种失败情况都需要特定的恢复机制,增加了系统复杂性。

为了应对这些挑战,常见的最终一致性实现机制包括:

  • 基于消息队列的异步通知
  • 定时任务扫描和修复不一致数据
  • CDC(变更数据捕获)技术监控数据库变更并同步到缓存
  • 读取:先查缓存,缓存没有则查数据库,然后将结果放入缓存
  • 更新:先更新数据库,然后删除缓存(而非更新缓存)
  • 优点:实现简单,适合读多写少场景
  • 缺点:存在短暂不一致窗口
  • 应用只与缓存交互,由缓存层负责与数据库的交互
  • 写入时,先更新缓存,缓存同步更新数据库
  • 优点:对应用透明,一致性较好
  • 缺点:增加了写延迟,缓存层复杂度高
  • 写入时只更新缓存,缓存异步批量更新数据库
  • 优点:写性能高,可合并多次更新
  • 缺点:数据库与缓存一致性最差,可能丢失数据

缓存过期策略直接影响一致性:

  • 较短的过期时间:减少不一致窗口,但增加缓存miss率和数据库压力
  • 较长的过期时间:提高缓存命中率,但可能长时间保持不一致状态
  • 不过期策略:需要显式的失效机制,否则可能永远不一致

缓存预热是提前加载热点数据到缓存:

  • 优点:避免系统启动时大量缓存miss
  • 缺点:如果预热数据与实际访问不匹配,会浪费缓存空间
  • 一致性挑战:预热数据可能很快过期或变得不一致

现代系统通常有多级缓存:

  • 本地缓存(应用内存)
  • 分布式缓存(如Redis)
  • 数据库缓存
  • CDN缓存
    每增加一层缓存,就增加了一层一致性挑战。

基于事件的缓存更新机制:

  1. 数据变更时发布事件到消息队列
  2. 缓存服务订阅这些事件并更新缓存
  3. 保证消息至少被处理一次(可能需要幂等处理)
    这种机制可以解耦数据库操作和缓存操作,提高系统弹性。

使用版本号或时间戳标记数据版本:

  • 缓存中存储数据时同时存储版本信息
  • 读取时比较版本,发现过时则更新
  • 写入时检查版本,避免写入旧数据

对于关键业务数据,设计特定的补偿机制:

  • 定期对账:比对缓存和数据库数据,修复不一致
  • 业务重试:关键操作失败时自动重试
  • 人工干预:对无法自动修复的问题提供人工处理接口

业务场景:电商商品信息更新
在电商平台中,商品信息(如价格、库存)经常变动,需要确保数据库和缓存的一致性。

// 1. 定义商品更新事件
public class ProductUpdateEvent {
private Long productId;
private String field;
private Object newValue;
private long timestamp;
// 构造器、getter和setter省略
}
// 2. 商品服务实现
@Service
public class ProductServiceImpl implements ProductService {
@Autowired
private ProductRepository productRepository;
@Autowired
private KafkaTemplate<String, ProductUpdateEvent> kafkaTemplate;
@Override
@Transactional
public void updateProductPrice(Long productId, BigDecimal newPrice) {
// 更新数据库
Product product = productRepository.findById(productId)
.orElseThrow(() -> new ProductNotFoundException(productId));
product.setPrice(newPrice);
product.setUpdateTime(new Date());
productRepository.save(product);
// 发布商品更新事件
ProductUpdateEvent event = new ProductUpdateEvent();
event.setProductId(productId);
event.setField("price");
event.setNewValue(newPrice);
event.setTimestamp(System.currentTimeMillis());
kafkaTemplate.send("product-updates", event);
log.info("Product price updated and event published: {}", productId);
}
}
// 3. 缓存更新服务
@Service
public class ProductCacheService {
@Autowired
private RedisTemplate<String, Product> redisTemplate;
@Autowired
private ProductRepository productRepository;
@KafkaListener(topics = "product-updates")
public void handleProductUpdate(ProductUpdateEvent event) {
String cacheKey = "product:" + event.getProductId();
try {
// 获取当前缓存
Product cachedProduct = redisTemplate.opsForValue().get(cacheKey);
if (cachedProduct != null) {
// 如果缓存存在,更新特定字段
switch (event.getField()) {
case "price":
cachedProduct.setPrice((BigDecimal) event.getNewValue());
break;
case "stock":
cachedProduct.setStock((Integer) event.getNewValue());
break;
// 其他字段更新...
default:
// 对于复杂更新,直接删除缓存
redisTemplate.delete(cacheKey);
return;
}
// 更新缓存,设置过期时间
redisTemplate.opsForValue().set(cacheKey, cachedProduct, 1, TimeUnit.HOURS);
log.info("Cache updated for product: {}", event.getProductId());
}
} catch (Exception e) {
// 出现异常,删除缓存,下次查询时重建
redisTemplate.delete(cacheKey);
log.error("Error updating cache, cache deleted: {}", event.getProductId(), e);
}
}
// 缓存重建方法
public Product getProduct(Long productId) {
String cacheKey = "product:" + productId;
// 尝试从缓存获取
Product product = redisTemplate.opsForValue().get(cacheKey);
if (product == null) {
// 缓存未命中,从数据库获取
product = productRepository.findById(productId)
.orElseThrow(() -> new ProductNotFoundException(productId));
// 放入缓存
redisTemplate.opsForValue().set(cacheKey, product, 1, TimeUnit.HOURS);
log.info("Cache rebuilt for product: {}", productId);
}
return product;
}
}

优势分析

  • 解耦数据库和缓存操作:通过消息队列分离写操作和缓存更新
  • 提高系统弹性:即使缓存服务暂时不可用,消息会在队列中等待处理
  • 支持多种缓存更新策略:可以选择更新缓存或删除缓存
  • 可追踪性:事件包含时间戳和变更信息,便于问题排查

业务场景:用户资料更新
用户资料可能在多个服务中被访问和修改,需要确保缓存中的数据是最新的。

// 1. 带版本号的用户实体
@Entity
@Table(name = "users")
public class User implements Serializable {
@Id
private Long id;
private String username;
private String email;
private String phone;
@Version
private Long version;
@Column(name = "update_time")
private Timestamp updateTime;
// 构造器、getter和setter省略
}
// 2. 用户服务实现
@Service
public class UserServiceImpl implements UserService {
@Autowired
private UserRepository userRepository;
@Autowired
private RedisTemplate<String, User> redisTemplate;
@Override
@Transactional
public User updateUserProfile(Long userId, UserProfileDTO profileDTO) {
User user = userRepository.findById(userId)
.orElseThrow(() -> new UserNotFoundException(userId));
// 更新用户信息
user.setEmail(profileDTO.getEmail());
user.setPhone(profileDTO.getPhone());
// 其他字段更新...
// 更新时间戳(@Version会自动更新版本号)
user.setUpdateTime(new Timestamp(System.currentTimeMillis()));
// 保存到数据库
User updatedUser = userRepository.save(user);
// 更新缓存,包含版本信息
String cacheKey = "user:" + userId;
redisTemplate.opsForValue().set(cacheKey, updatedUser, 30, TimeUnit.MINUTES);
return updatedUser;
}
@Override
public User getUserProfile(Long userId) {
String cacheKey = "user:" + userId;
// 尝试从缓存获取
User cachedUser = redisTemplate.opsForValue().get(cacheKey);
if (cachedUser != null) {
// 检查缓存数据是否过时
User latestVersion = userRepository.findVersionAndUpdateTimeById(userId);
if (latestVersion != null &&
(cachedUser.getVersion() < latestVersion.getVersion() ||
cachedUser.getUpdateTime().before(latestVersion.getUpdateTime()))) {
// 缓存版本过时,从数据库重新获取完整数据
User freshUser = userRepository.findById(userId)
.orElseThrow(() -> new UserNotFoundException(userId));
// 更新缓存
redisTemplate.opsForValue().set(cacheKey, freshUser, 30, TimeUnit.MINUTES);
return freshUser;
}
return cachedUser;
}
// 缓存未命中,从数据库获取
User user = userRepository.findById(userId)
.orElseThrow(() -> new UserNotFoundException(userId));
// 放入缓存
redisTemplate.opsForValue().set(cacheKey, user, 30, TimeUnit.MINUTES);
return user;
}
}
// 3. 用户仓库接口
public interface UserRepository extends JpaRepository<User, Long> {
// 只查询版本和更新时间,减少数据传输
@Query("SELECT new User(u.id, u.version, u.updateTime) FROM User u WHERE u.id = :userId")
User findVersionAndUpdateTimeById(@Param("userId") Long userId);
}

优势分析

  • 版本检测:通过版本号或时间戳快速判断缓存是否过时
  • 按需更新:只有当发现缓存过时时才从数据库获取完整数据
  • 减少数据库负载:版本检查查询很轻量,不需要获取所有字段
  • 乐观锁支持:@Version注解自动提供乐观锁功能,防止并发更新问题

业务场景:订单支付状态同步
支付系统和订单系统之间的状态同步,需要确保即使在系统部分失败的情况下,最终也能达到一致状态。

// 1. 订单状态检查和补偿服务
@Service
public class OrderConsistencyService {
@Autowired
private OrderRepository orderRepository;
@Autowired
private PaymentRepository paymentRepository;
@Autowired
private RedisTemplate<String, Order> redisTemplate;
@Autowired
private NotificationService notificationService;
// 定时任务,每5分钟执行一次
@Scheduled(fixedRate = 300000)
public void checkAndCompensateOrders() {
log.info("Starting order consistency check...");
// 查找处于中间状态的订单
List<Order> pendingOrders = orderRepository.findByStatusAndUpdateTimeBefore(
OrderStatus.PAYMENT_PENDING,
new Timestamp(System.currentTimeMillis() - TimeUnit.MINUTES.toMillis(15))
);
for (Order order : pendingOrders) {
try {
compensateOrderIfNeeded(order);
} catch (Exception e) {
log.error("Error compensating order {}: {}", order.getId(), e.getMessage(), e);
}
}
log.info("Order consistency check completed, processed {} orders", pendingOrders.size());
}
private void compensateOrderIfNeeded(Order order) {
// 检查支付状态
Payment payment = paymentRepository.findByOrderId(order.getId());
if (payment == null) {
// 支付记录不存在,可能是创建失败
log.warn("Payment record not found for order {}, marking as payment failed", order.getId());
updateOrderStatus(order, OrderStatus.PAYMENT_FAILED);
notificationService.notifyCustomer(order.getCustomerId(),
"您的订单支付未完成,请重新尝试支付或联系客服");
return;
}
switch (payment.getStatus()) {
case COMPLETED:
// 支付已完成,但订单状态未更新
log.info("Payment completed but order status not updated for order {}", order.getId());
updateOrderStatus(order, OrderStatus.PAID);
// 更新库存和其他相关业务逻辑
try {
completeOrderProcessing(order);
} catch (Exception e) {
log.error("Error processing paid order {}: {}", order.getId(), e.getMessage(), e);
// 记录错误,但不影响订单状态更新
notificationService.notifyAdmin("订单处理异常: " + order.getId());
}
break;
case FAILED:
// 支付失败,更新订单状态
log.info("Payment failed for order {}", order.getId());
updateOrderStatus(order, OrderStatus.PAYMENT_FAILED);
notificationService.notifyCustomer(order.getCustomerId(),
"您的订单支付失败,请重新尝试或选择其他支付方式");
break;
case PENDING:
// 支付仍在处理中,但已超时
if (payment.getCreateTime().before(
new Timestamp(System.currentTimeMillis() - TimeUnit.MINUTES.toMillis(30)))) {
log.warn("Payment pending too long for order {}, marking as timeout", order.getId());
payment.setStatus(PaymentStatus.TIMEOUT);
paymentRepository.save(payment);
updateOrderStatus(order, OrderStatus.PAYMENT_TIMEOUT);
notificationService.notifyCustomer(order.getCustomerId(),
"您的订单支付处理超时,请联系客服确认支付状态");
}
break;
default:
log.info("No compensation needed for order {} with payment status {}",
order.getId(), payment.getStatus());
}
}
private void updateOrderStatus(Order order, OrderStatus newStatus) {
// 更新数据库
order.setStatus(newStatus);
order.setUpdateTime(new Timestamp(System.currentTimeMillis()));
orderRepository.save(order);
// 更新缓存
String cacheKey = "order:" + order.getId();
Order cachedOrder = redisTemplate.opsForValue().get(cacheKey);
if (cachedOrder != null) {
cachedOrder.setStatus(newStatus);
cachedOrder.setUpdateTime(order.getUpdateTime());
redisTemplate.opsForValue().set(cacheKey, cachedOrder, 1, TimeUnit.HOURS);
}
log.info("Order {} status updated to {}", order.getId(), newStatus);
}
private void completeOrderProcessing(Order order) {
// 执行订单支付成功后的业务逻辑
// 如更新库存、创建物流单、发送确认邮件等
}
}