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
相关产品推荐
相关产品推荐

