Administrator
发布于 2019-05-01 / 2828 阅读
47

Spring 事件机制与异步解耦实践

加一个"送新人券"的需求,我改了注册方法

4 月底产品提了个小需求:新用户注册成功后自动发一张 20 元的无门槛券。我打开 UserService.register(),120 行,里面已经塞了校验、插库、发欢迎短信、埋点上报、初始化用户配置五件事。

我在第 87 行插了三行调用发券的代码。提 PR 的时候,带我的师兄评论了一句:"加一个运营需求,为什么要动注册主流程?万一发券接口挂了,是不是用户就注册不了了?"

这句话点醒我了。我当时没想过这个问题,虽然发券那段我 try-catch 了,但下一个加需求的人未必会。于是决定用 Spring 的事件机制把这块拆开。

原来的代码

简化之后大概是这样,每一步都直接依赖下一步的实现类:

@Service
public class UserService {

    @Autowired private UserMapper userMapper;
    @Autowired private SmsClient smsClient;          // HTTP
    @Autowired private CouponClient couponClient;    // Dubbo
    @Autowired private PointClient pointClient;      // Dubbo
    @Autowired private TrackerClient trackerClient;  // HTTP

    @Transactional
    public Long register(RegisterDTO dto) {
        checkMobile(dto.getMobile());

        User user = new User();
        user.setMobile(dto.getMobile());
        user.setStatus(1);
        userMapper.insert(user);                      // 15 ms

        initUserConfig(user.getId());                 // 30 ms
        smsClient.sendWelcome(dto.getMobile());       // P99 320 ms
        couponClient.grantNewUserCoupon(user.getId());// P99 180 ms
        pointClient.addRegisterPoint(user.getId());   // P99 60 ms
        trackerClient.report("register", user.getId());// P99 95 ms

        return user.getId();
    }
}

问题有几个,按严重程度排:

  • 可用性耦合:发券、发短信、埋点这些非核心链路的服务抖动,会直接拖慢甚至拖垮注册。
  • 响应时间:注册接口 P99 是 780ms,其中 655ms 花在这些"顺便做的事"上。
  • 事务边界混乱:这些操作都在 @Transactional 里面,一次 HTTP 超时会把事务拖长,数据库连接被多占用 700ms。
  • 改一个需求动一个方法:每次加功能都要回来改这段已经很长的方法,代码评审时冲突不断。

用事件把它拆开

Spring 的事件机制是很标准的观察者模式,三个角色:事件、发布者、监听者。

先定义事件。Spring 4.2 之后不强制继承 ApplicationEvent 了,普通 POJO 也行,但我们还是继承了,方便用 getSource() 追溯发布方:

public class UserRegisteredEvent extends ApplicationEvent {

    private final Long userId;
    private final String mobile;

    public UserRegisteredEvent(Object source, Long userId, String mobile) {
        super(source);
        this.userId = userId;
        this.mobile = mobile;
    }

    public Long getUserId() { return userId; }
    public String getMobile() { return mobile; }
}

发布方注入 ApplicationEventPublisher。改造之后主流程只剩下"必须同步完成"的部分:

@Service
public class UserService {

    @Autowired private UserMapper userMapper;
    @Autowired private ApplicationEventPublisher publisher;

    @Transactional
    public Long register(RegisterDTO dto) {
        checkMobile(dto.getMobile());

        User user = new User();
        user.setMobile(dto.getMobile());
        user.setStatus(1);
        userMapper.insert(user);

        // 发个事件就完事,谁关心、怎么处理,主流程一概不知
        publisher.publishEvent(new UserRegisteredEvent(this, user.getId(), dto.getMobile()));
        return user.getId();
    }
}

监听方各自独立,一个监听器只干一件事:

@Component
@Slf4j
public class RegisterEventListener {

    @Autowired private SmsClient smsClient;
    @Autowired private CouponClient couponClient;
    @Autowired private PointClient pointClient;

    @Async("eventExecutor")
    @EventListener
    public void sendWelcomeSms(UserRegisteredEvent event) {
        smsClient.sendWelcome(event.getMobile());
    }

    @Async("eventExecutor")
    @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
    public void grantNewUserCoupon(UserRegisteredEvent event) {
        couponClient.grantNewUserCoupon(event.getUserId());
    }

    @Async("eventExecutor")
    @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
    public void addRegisterPoint(UserRegisteredEvent event) {
        pointClient.addRegisterPoint(event.getUserId());
    }
}

现在加一个新功能(比如推一条站内信),只需要新写一个监听器方法,UserService 一个字都不用改。

三个必须注意的点

一、事件默认是同步执行的

这是我踩的第一个坑。写完一跑,注册接口还是 780ms,一点没变。

原因:publishEvent 默认走的是 SimpleApplicationEventMulticaster,它拿到监听器列表后在当前线程里串行调用。异步必须显式开:

@Configuration
@EnableAsync
public class AsyncConfig {

    @Bean("eventExecutor")
    public Executor eventExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        executor.setCorePoolSize(4);
        executor.setMaxPoolSize(8);
        executor.setQueueCapacity(200);
        executor.setThreadNamePrefix("event-");
        executor.setRejectedExecutionHandler(
                new ThreadPoolExecutor.CallerRunsPolicy());
        executor.initialize();
        return executor;
    }
}

也可以给 SimpleApplicationEventMulticaster 直接设置一个全局的 Executor,这样所有 @EventListener 都变成异步。我们没这么做,因为有些监听器(比如初始化用户配置)必须同步,混在一起不好控制。

二、事务还没提交,下游查不到数据

第二个坑更隐蔽。发券接口内部会用 userId 反查用户信息,结果偶发报"用户不存在"。

原因是 @EventListener 的调用时机在 publishEvent 那一行,也就是事务提交之前。异步线程里拿着 userId 去查从库,主库的事务还没提交,从库自然没有这条记录。这个问题只在主从有延迟、且异步线程执行得比主库提交还快的时候出现,所以是偶发的。

解决方式是 @TransactionalEventListener,它把监听器的执行挂到事务的某个阶段上:

phase执行时机我们的用法
BEFORE_COMMIT提交前没用
AFTER_COMMIT(默认)事务成功提交后发券、加积分
AFTER_ROLLBACK事务回滚后清理预占的资源
AFTER_COMPLETION事务结束后(不管成功失败)记录日志

换成 AFTER_COMMIT 之后,"用户不存在"的报错彻底消失了。

它的实现原理是注册了一个 TransactionSynchronization 回调。这里有个限制:它必须在事务上下文里才生效。如果发布事件的方法没有 @Transactional,这个监听器根本不会被触发,而且不会有任何日志提示。我们测试环境就出现过"事件没执行"的假象,最后发现是那个方法忘了加事务注解。

三、异常处理和事件丢失

同步监听器和异步监听器的异常行为完全不同:

  • 同步监听器抛异常:会向上冒泡到 publishEvent,进而中断主流程、回滚事务。这是我们不希望的。
  • 异步监听器抛异常:在另一个线程里,主流程感知不到,事务照常提交。但异常会被 SimpleAsyncUncaughtExceptionHandler 吞掉只打一行日志,容易漏掉。

我们配了统一的异常处理器,把失败的事件记到一张表里,方便重试:

@Configuration
@EnableAsync
public class AsyncConfig implements AsyncConfigurer {

    @Override
    public Executor getAsyncExecutor() {
        return eventExecutor();
    }

    @Override
    public AsyncUncaughtExceptionHandler getAsyncUncaughtExceptionHandler() {
        return (ex, method, params) ->
            log.error("异步事件执行失败, method={}, params={}",
                      method.getName(), Arrays.toString(params), ex);
    }
}

更重要的一个代价必须说清楚:异步事件存在丢失风险。如果事件刚发布、事务刚提交,进程这时候挂了,事件就永远不会被消费,券也就没发出去。我们评估后接受了这个风险,因为新人券可以补发,业务上每天跑一次对账就能发现。如果是"注册送 100 元余额"这种不能丢的场景,得用本地消息表或者消息队列,事件机制扛不住。

改造前后

压测数据(200 并发,5 分钟):

指标改造前改造后
注册接口 P99780 ms145 ms
注册接口 TPS2601180
单个请求占用数据库连接时长约 700 ms约 45 ms
发券服务宕机时注册是否可用不可用可用

TPS 提升比 P99 的改善更夸张,因为事务时间短了,数据库连接池不再成为瓶颈。

小结

  • 事件机制最大的价值是让主流程不必知道下游的存在。加一个新功能时新增一个监听器就行,核心方法不用动。
  • @EventListener 默认是同步的,要异步必须加 @Async 并开 @EnableAsync,最好用自定义线程池而不是默认的 SimpleAsyncTaskExecutor(后者每次都新建线程)。
  • 监听器里要读刚写入的数据,用 @TransactionalEventListener(AFTER_COMMIT)。注意它依赖事务上下文,没事务就不触发。
  • 异步意味着事件可能丢失,能丢的和不能丢的要分开设计,不能一律异步。

最后补一句,@Async@TransactionalEventListener 都依赖 AOP 代理,所以"同类内部调用失效"这条规则对它们同样适用。我们有个监听器写在 UserService 里面,想直接调用另一个监听器方法,结果注解完全没起作用,最后是拆成独立组件解决的。

参考