You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Quarkus中@RunOnVirtualThread场景下如何替代sleep实现可靠异步邮件发送?

解决Quarkus异步邮件发送的事务一致性问题

你当前用TimeUnit.MILLISECONDS.sleep(2000)的临时方案本质是试图等待主事务完成数据库提交,但这种方式不可靠——事务提交时长可能因负载波动变化,短了还是会出现查不到数据的情况,长了又会浪费系统资源。以下是两种可靠的替代方案:

方案一:事务提交后再发送事件

核心思路是确保事件仅在主事务成功提交后才发送,这样事件处理器查询数据库时,相关数据已经持久化完成。

修改invitePerson方法,利用Quarkus的事务API添加事务完成回调:

private void invitePerson(Kurs kurs, User dancer) {
    Kursinvitation kursinvitation = createInvitation(dancer, kurs);
    
    // 注册事务提交后的执行动作
    Transaction.current().addCompletion(() -> {
        CourseInvitationEvent event = new CourseInvitationEvent(
                kurs.getId(),
                dancer.getId(),
                kursinvitation.getId()
        );
        eventBus.send("course-invitation", event);
        log.debug("Sent CourseInvitationEvent after transaction commit for dancer {} and kurs {}",
                dancer.getId(), kurs.getId());
    });
}

方案二:给事件处理器添加重试机制(配合方案一使用更佳)

即使确保事务提交后发送事件,也可能因数据库读写延迟、网络抖动等偶发问题导致查询失败。此时可以用Quarkus的@Retry注解添加重试逻辑,提升可靠性。

修改事件处理器代码,移除sleep并添加重试:

@ApplicationScoped
@Slf4j
public class CourseInvitationEventListener {

    @Inject
    KursDao kursDao;
    @Inject
    UserDao userDao;
    @Inject
    KursMailerService kursMailerService;
    @Inject
    KursinvitationDao kursinvitationDao;

    @ConsumeEvent(value = "course-invitation", blocking = true)
    @RunOnVirtualThread
    @Retry(maxRetries = 3, delay = 500) // 最多重试3次,每次间隔500ms
    public void handleCourseInvitationEvent(CourseInvitationEvent event) {
        try {
            log.debug("Processing CourseInvitationEvent for dancer {} and kurs {}",
                    event.getDancerId(), event.getKursId());

            Kurs kurs = kursDao.findById(event.getKursId());
            User dancer = userDao.findById(event.getDancerId());
            Kursinvitation kursinvitation = kursinvitationDao.findById(event.getKursinvitationId());

            if (kurs == null || dancer == null || kursinvitation == null) {
                log.error("Could not find entities for CourseInvitationEvent: kursId={}, dancerId={}, invitationId={}",
                        event.getKursId(), event.getDancerId(), event.getKursinvitationId());
                throw new IllegalStateException("Required entities not found"); // 抛出异常触发重试
            }

            kursMailerService.sendEinladungZumKurs(kurs, dancer, kursinvitation);

            log.info("Successfully sent invitation email for dancer {} and kurs {}",
                    event.getDancerId(), event.getKursId());

        } catch (Exception e) {
            log.error("Error sending invitation email for event: {}", event, e);
            throw e; // 抛出异常让重试机制生效
        }
    }
}

额外说明

  • 方案一从根源上解决了事务未提交就触发事件的问题,是首选方案;
  • 方案二作为兜底,处理偶发的异常情况,两者结合能最大化可靠性;
  • 可移除处理器上的@Transactional注解:邮件发送通常是非事务性操作,不必要的事务会增加系统开销。

内容的提问来源于stack exchange,提问作者Paul

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.12 05:27:02