Spring Modulith中检测死信事件的惯用方法是什么?
如何优雅检测Spring Modulith中反复重试失败的事件?
在探索Spring Modulith时,我们需要以惯用方式检测那些反复重试却始终无法完成的事件,并对此类情况进行通知。
Spring Modulith提供了IncompleteEventPublications抽象接口用于重新提交待处理事件,但该接口无法直接满足检测需求。我们曾尝试两种不太理想的实现方式:
不推荐的实现方式
滥用Predicate产生副作用
通过在resubmitIncompletePublications的过滤Predicate中嵌入通知逻辑,违反了函数式编程的单一职责原则,代码可读性和维护性差:@Scheduled(fixedDelayString = "5s") void retryJobs() { incompleteEventPublications.resubmitIncompletePublications(event -> { var isOld = event.getPublicationDate().isBefore(Instant.now().minus(Duration.ofHours(3))); if (isOld) { incrementDeadLetterCounter(); } return !isOld; }); }使用内部组件EventPublicationRegistry
直接依赖spring-modulith-events-core包中的EventPublicationRegistry,该组件属于内部实现,未来可能存在兼容性风险:private final org.springframework.modulith.events.core.EventPublicationRegistry registry; @Scheduled(fixedDelayString = "1m") void detectDeadLetters() { var deadLetterAmount = registry.findIncompletePublications().stream() .map(EventPublication::getPublicationDate) .filter(e -> e.isBefore(Instant.now().minus(Duration.ofHours(3)))) .count(); setGauge(deadLetterAmount); }
最优实现方式:使用公开API EventPublicationRepository
Spring Modulith提供了**EventPublicationRepository**这一公开API,专门用于查询和管理事件发布状态,是检测停滞事件的惯用方式:
示例代码
import org.springframework.modulith.events.EventPublication; import org.springframework.modulith.events.EventPublicationRepository; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import java.time.Duration; import java.time.Instant; import java.util.List; @Component public class StalledEventMonitor { private final EventPublicationRepository publicationRepository; private final EventNotificationService notificationService; // 自定义通知服务 public StalledEventMonitor(EventPublicationRepository publicationRepository, EventNotificationService notificationService) { this.publicationRepository = publicationRepository; this.notificationService = notificationService; } @Scheduled(fixedDelayString = "1m") void checkAndAlertStalledEvents() { // 设置超时阈值:3小时未完成的事件判定为停滞事件 Instant timeoutThreshold = Instant.now().minus(Duration.ofHours(3)); // 查询所有超时未完成的事件 List<EventPublication> stalledEvents = publicationRepository.findIncompletePublications() .filter(event -> event.getPublicationDate().isBefore(timeoutThreshold)) .toList(); if (!stalledEvents.isEmpty()) { // 发送通知(比如告警、日志记录、存入死信队列等) notificationService.sendStalledEventAlert(stalledEvents); // 可选:将事件标记为已完成,避免重复检测(根据业务需求决定) stalledEvents.forEach(event -> publicationRepository.markAsCompleted(event.getIdentifier())); } } }
优势说明
- 符合惯用设计:
EventPublicationRepository是Spring Modulith对外暴露的标准API,不属于内部实现,兼容性有保障 - 职责清晰:将事件检测与通知逻辑分离,代码结构更清晰,易于维护
- 灵活扩展:可根据业务需求调整过滤条件(比如结合重试次数、事件类型等),实现更细粒度的检测
扩展方案
如果需要基于重试次数而非仅时间来判定停滞事件,可以在事件对象中添加重试次数属性,发布事件时更新该属性,然后在查询时通过event.getEvent()获取事件对象并过滤:
.filter(event -> { var businessEvent = (YourCustomEvent) event.getEvent(); return businessEvent.getRetryCount() >= 5; // 重试超过5次判定为停滞 })
内容的提问来源于stack exchange,提问作者Sebastian Sprenger
相关产品推荐
相关产品推荐

