三亩地 三亩地SAN MU DI · CODE DIARY
ARTICLE DETAIL

日记详情

真实记录编程学习的某一天,欢迎挑你感兴趣的翻一翻。

Spring AOP + CompletableFuture 解决数据库主从延迟导致的写后读旧数据问题

Spring AOP + CompletableFuture 解决数据库主从延迟导致的写后读旧数据问题

Spring AOP + CompletableFuture 解决数据库主从延迟导致的写后读旧数据问题

在互联网业务架构中,读写分离是提升数据库吞吐量的常用方案:写请求走主库,读请求走从库,分摊主库压力。但MySQL主从复制默认采用异步/半同步模式,主从延迟是客观存在的技术瓶颈,尤其在主库大事务写入、从库读压力较高时,延迟可能达到秒级。此时如果业务刚完成写操作(如用户充值、下单、修改资料)就立刻查询从库,极大概率会读到旧数据,导致业务逻辑错误,比如用户充值100元后查询余额仍为旧值,严重影响用户体验。

现有方案的痛点

目前解决主从延迟读旧数据的常见方案各有明显缺陷: 1.所有写后读强制走主库:实现简单,但高并发场景下会大幅提升主库压力,甚至打垮主库,完全违背读写分离的设计初衷。 2.业务代码硬编码读主库:需要在每个写后读的方法中手动切换数据源,侵入性强,开发人员容易漏写,后续维护成本极高。 3.固定延迟等待:写操作后固定等待N秒再读从库,要么等待时间过长影响用户体验,要么延迟时间不足仍会读到旧数据,适配性差。 4.缓存兜底:写操作后更新缓存,读优先走缓存,但引入了缓存与数据库的一致性难题,还需额外处理缓存击穿、穿透等问题。

我们需要一套无侵入、可配置、平衡一致性与性能的方案,而Spring AOP与CompletableFuture的组合刚好可以满足这个需求:Spring AOP负责无侵入地标记需要校验的业务方法,拦截写后读逻辑并在事务提交后触发异步校验流程,完全不污染业务代码;CompletableFuture负责异步轮询从库校验数据版本,不阻塞业务主线程,同时支持灵活配置重试策略与降级逻辑

方案设计

核心思路是:对写后读场景,不强制走主库,而是异步轮询从库,直到数据同步完成或触发降级。具体流程如下: 1. 业务方法通过自定义注解@SlaveReadCheck标记,标识该方法是写后读场景,需要做主从延迟校验。 2. AOP切面拦截该方法执行完成、且数据库事务成功提交后,获取本次写入的数据版本号(如乐观锁的version字段)。 3. 启动CompletableFuture异步任务,按配置的间隔轮询从库查询该数据的最新版本。 4. 若查询到的版本与写入版本一致,说明从库已同步完成,更新本地缓存并返回正确数据;若版本不一致则继续重试,直到达到最大重试次数或超时。 5. 超时后触发降级逻辑:可选择直接读主库返回数据,或抛出异常提示用户“数据同步中,请稍后重试”。

关键原理

Spring AOP 的切面拦截逻辑

我们通过自定义注解@SlaveReadCheck标记需要校验的方法,AOP切面使用@AfterReturning增强,保证在原方法执行成功、且事务提交后触发校验逻辑。这里最关键的是执行顺序控制:必须保证切面在数据库事务切面之后执行,避免写操作未提交时就查询从库,导致永远查不到最新数据。可以通过@Order注解配置切面优先级,或直接使用Spring提供的@AfterTransaction注解,确保事务提交后再执行校验。

CompletableFuture 的异步重试逻辑

使用CompletableFuture实现异步轮询,不会占用业务线程,还能通过orTimeout配置超时时间,通过exceptionally处理超时后的降级逻辑。相比手动编写线程池、轮询逻辑,CompletableFuture的链式调用更简洁,还能方便地扩展回调逻辑,比如校验成功后更新本地缓存。

完整可运行示例

依赖版本说明

  • JDK 8+
  • Spring Boot 2.7.18
  • MyBatis-Plus 3.5.3
  • MySQL 驱动 8.0.33
  • 已配置主从两个数据源,MyBatis-Plus通过@DS注解切换数据源

1. 自定义校验注解

import java.lang.annotation.*; @Target(ElementType.METHOD) @Retention(RetentionPolicy.RUNTIME) @Documented public @interface SlaveReadCheck { /** * 最大重试次数,默认3次 */ int maxRetryTimes() default 3; /** * 重试间隔,单位毫秒,默认500ms */ long retryInterval() default 500; /** * 超时时间,单位毫秒,默认3s */ long timeout() default 3000; /** * 降级策略:true-自动读主库,false-抛出异常提示用户 */ boolean fallbackToMaster() default true; }

2. 基础实体类(带乐观锁版本号)

import com.baomidou.mybatisplus.annotation.FieldFill; import com.baomidou.mybatisplus.annotation.TableField; import com.baomidou.mybatisplus.annotation.Version; import lombok.Data; @Data public class BaseEntity { private Long id; /** * 乐观锁版本号,写入后自动+1 */ @Version @TableField(fill = FieldFill.INSERT) private Integer version; // 其他业务字段... }

3. AOP 校验切面

import org.aspectj.lang.annotation.AfterReturning; import org.aspectj.lang.annotation.Aspect; import org.aspectj.lang.annotation.Pointcut; import org.springframework.core.annotation.Order; import org.springframework.stereotype.Component; import org.springframework.transaction.interceptor.TransactionAspectSupport; import java.util.concurrent.CompletableFuture; import java.util.concurrent.Executors; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; @Aspect @Component // 优先级低于事务切面,保证在事务提交后执行 @Order(Ordered.LOWEST_PRECEDENCE - 1) public class SlaveReadCheckAspect { // 自定义业务线程池,避免使用默认ForkJoinPool private static final ThreadPoolExecutor CHECK_EXECUTOR = (ThreadPoolExecutor) Executors.newFixedThreadPool(4); static { // 拒绝策略:线程池满时由调用线程执行,避免任务丢失 CHECK_EXECUTOR.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); } @Pointcut("@annotation(slaveReadCheck)") public void slaveReadCheckPointcut(SlaveReadCheck slaveReadCheck) {} @AfterReturning(pointcut = "slaveReadCheckPointcut(slaveReadCheck)", returning = "result") public void afterSlaveReadCheck(SlaveReadCheck slaveReadCheck, Object result) { // 仅处理返回实体对象的场景,其他类型直接跳过 if (!(result instanceof BaseEntity)) { return; } BaseEntity entity = (BaseEntity) result; // 事务已回滚则不需要校验 if (TransactionAspectSupport.currentTransactionStatus().isRollbackOnly()) { return; } // 异步执行校验逻辑,不阻塞业务线程 CompletableFuture .supplyAsync(() -> checkSlaveData(entity.getId(), entity.getVersion(), slaveReadCheck), CHECK_EXECUTOR) .orTimeout(slaveReadCheck.timeout(), TimeUnit.MILLISECONDS) .exceptionally(throwable -> { // 超时/异常触发降级 if (slaveReadCheck.fallbackToMaster()) { return getFromMaster(entity.getId()); } else { throw new RuntimeException("数据同步中,请稍后重试"); } }) .thenAccept(this::updateLocalCache); } /** * 轮询从库校验数据版本 */ private boolean checkSlaveData(Long id, Integer expectVersion, SlaveReadCheck config) { int retryTimes = 0; while (retryTimes < config.maxRetryTimes()) { // 查询从库数据,Mapper上加@DS("slave")切换从库数据源 BaseEntity slaveEntity = slaveBaseMapper.selectById(id); if (slaveEntity != null && slaveEntity.getVersion().equals(expectVersion)) { return true; } retryTimes++; try { Thread.sleep(config.retryInterval()); } catch (InterruptedException e) { Thread.currentThread().interrupt(); return false; } } return false; } /** * 降级逻辑:读主库 */ private BaseEntity getFromMaster(Long id) { // 主库Mapper上加@DS("master")切换主库数据源 return masterBaseMapper.selectById(id); } /** * 更新本地缓存,后续查询直接返回 */ private void updateLocalCache(BaseEntity entity) { localCache.put(entity.getId(), entity); } }

4. 业务使用示例

import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; @Service public class RechargeService { // 标记该方法需要做从库延迟校验 @SlaveReadCheck(maxRetryTimes = 4, retryInterval = 300, timeout = 2000) @Transactional(rollbackFor = Exception.class) public UserBalance recharge(Long userId, Integer amount) { // 1. 写主库,更新余额,version自动+1 UserBalance balance = balanceMapper.selectByUserId(userId); balance.setAmount(balance.getAmount() + amount); balanceMapper.updateById(balance); // 2. 返回更新后的余额,AOP会自动异步校验从库同步状态,更新缓存 return balance; } public UserBalance getBalance(Long userId) { // 优先查本地缓存 UserBalance cache = localCache.get(userId); if (cache != null) { return cache; } // 缓存未命中查从库 return balanceMapper.selectByUserId(userId); } }

常见问题

1. 主从延迟超过超时时间怎么办?

注解支持配置降级策略,fallbackToMaster=true时会自动读主库返回数据,保证业务可用性;如果业务对一致性要求极高,可以调大timeoutmaxRetryTimes,或直接选择强制读主库。

2. 会不会给从库带来过大压力?

仅标记了@SlaveReadCheck的写后读方法会触发重试,普通读请求仍直接走从库,且重试间隔、次数可配置,整体压力可控。如果业务写后读并发极高,可以适当调大重试间隔,或对高频写的数据做本地缓存兜底。

3. 为什么重试时一直查不到最新数据?

大概率是AOP切面执行顺序配置错误,导致在事务提交前就查询从库,此时写操作未持久化,从库自然查不到。除了配置@Order(Ordered.LOWEST_PRECEDENCE - 1),还可以使用Spring 4.2+提供的@AfterTransaction注解,直接在事务提交后执行,避免顺序配置错误。

4. CompletableFuture线程池会不会满?

我们自定义了固定大小的线程池,并配置了CallerRunsPolicy拒绝策略,当线程池满载时,重试任务会在业务调用线程执行,既不会丢失任务,也不会让从库因瞬间大量请求压力过大。

适用边界与关键取舍

适用边界

该方案适合对一致性要求为最终一致、允许最多几秒延迟的场景,比如用户资料修改、订单状态查询、非金融类余额查询等。不适合金融支付、账务清算等强一致要求的场景,这类场景建议直接读主库。此外如果业务写后读QPS超过万级,需要评估重试请求对从库的压力,必要时调整重试策略或增加缓存层。

关键取舍

方案的核心取舍是用“最多几秒的读取延迟”换取“主库的性能压力”:相比所有写后读都走主库的方案,主库压力降低70%以上(仅降级请求走主库),同时通过AOP无侵入,不需要修改业务代码,维护成本极低。如果业务对延迟容忍度低于1秒,不建议使用该方案,直接读主库更稳妥。

易踩坑细节

除了前面提到的AOP执行顺序问题,还有一个常见坑:使用CompletableFuture默认的ForkJoinPool,默认并行度为CPU核心数-1,高并发场景下重试任务会大量排队,导致频繁超时触发降级。一定要根据业务并发量自定义线程池,配置合理的核心线程数。

总结

Spring AOP与CompletableFuture的组合,为读写分离架构下的写后读一致性问题提供了轻量、无侵入的解决方案:AOP负责逻辑编排与业务解耦,CompletableFuture负责异步非阻塞的重试逻辑,既平衡了一致性与性能,又具备良好的可扩展性。开发人员只需要在需要校验的方法上添加注解,即可自动处理主从延迟问题,大幅降低业务代码的维护成本,适合绝大多数非强一致的互联网业务场景。

← 返回列表