Axon事件溯源服务重启时事件重复触发致数据重复问题求助
问题描述
我正在实践Axon与事件溯源技术,开发一个追踪音乐家排练/演出出勤情况的Web应用。功能包括添加/修改音乐家、添加/修改音乐活动,以及生成出勤槽(PresenceSlot)来标记音乐家在活动中的出勤状态。
问题是每次重启应用,Axon事件存储中的事件会被重发,导致数据库里的数据不断重复。我知道这是因为在PresenceProjector中发送命令导致的,但不知道该如何重新建模。
请问应该从何处发送新的事件/命令,才能避免重启时重复发送?
相关代码
MusicianAggregate
@Aggregate public class MusicianAggregate { @AggregateIdentifier private UUID musicianId; private String firstName; private String lastName; private String phone; private LocalDate dateOfBirth; private LocalDate joinDate; private Boolean active; private String address; protected MusicianAggregate() { } @CommandHandler public MusicianAggregate(CreateMusicianCommand command) { apply(new MusicianSignedUpEvent(command.getMusicianId(), command.getFirstName(), command.getLastName(), command.getPhone(), command.getDateOfBirth(), command.getJoinDate(), command.getActive(), command.getAddress() )).andThenApply(() -> new MusicianJoinDateCreatedEvent(command.getMusicianId(), command.getJoinDate(), command.getFirstName() + " " + command.getLastName(), command.getActive())); } @EventSourcingHandler public void on(MusicianSignedUpEvent evt) { if (evt.getJoinDate() == null) { throw new MusicianJoinDateIsEmptyException(); } this.musicianId = evt.getMusicianId(); this.firstName = evt.getFirstName(); this.lastName = evt.getLastName(); this.phone = evt.getPhone(); this.dateOfBirth = evt.getDateOfBirth(); this.joinDate = evt.getJoinDate(); this.active = evt.getActive(); this.address = evt.getAddress(); } @CommandHandler public void on(UpdateMusicianCommand command) { var currentJoinDate = this.joinDate; apply(new MusicianUpdatedEvent(command.getId(), command.getFirstName(), command.getLastName(), command.getPhone(), command.getDateOfBirth(), command.getJoinDate(), command.getActive(), command.getAddress() )).andThenApplyIf( () -> !currentJoinDate.equals(command.getJoinDate()), () -> new MusicianJoinDateUpdatedEvent(command.getId(), currentJoinDate, command.getJoinDate())); } @EventSourcingHandler public void on(MusicianUpdatedEvent evt) { this.musicianId = evt.getMusicianId(); this.firstName = evt.getFirstName(); this.lastName = evt.getLastName(); this.phone = evt.getPhone(); this.dateOfBirth = evt.getDateOfBirth(); this.joinDate = evt.getJoinDate(); this.active = evt.getActive(); this.address = evt.getAddress(); } }
MusicEventAggregate
@Aggregate public class MusicEventAggregate { @AggregateIdentifier private UUID id; private String name; private String address; private LocalDateTime dateTime; public MusicEventAggregate() { } @CommandHandler public MusicEventAggregate(CreateMusicEventCommand command) { AggregateLifecycle.apply(new MusicEventCreatedEvent( command.getMusicEventId(), command.getName(), command.getAddress(), command.getDateTime())) .andThenApply(() -> new MusicEventDateCreated(command.getMusicEventId(), command.getDateTime())); } @EventSourcingHandler public void on(MusicEventCreatedEvent event) { this.id = event.getMusicEventId(); this.name = event.getName(); this.address = event.getAddress(); this.dateTime = event.getDateTime(); } @CommandHandler public void handle(UpdateMusicEventCommand command) { var currentDateTime = this.dateTime; AggregateLifecycle.apply(new MusicEventCreatedEvent( command.getId(), command.getName(), command.getAddress(), command.getDateTime())) .andThenApplyIf( () -> !currentDateTime.equals(command.getDateTime()), () -> new MusicEventDateUpdated(command.getId(), currentDateTime, command.getDateTime())); } @EventSourcingHandler public void on(MusicEventUpdatedEvent event) { this.name = event.getName(); this.address = event.getAddress(); this.dateTime = event.getDateTime(); } }
PresenceProjector
@Component @Slf4j public class PresenceProjector { private final PresenceRepository presenceRepository; private final MusicianRepository musicianRepository; private final MusicEventRepository musicEventRepository; private final CommandGateway commandGateway; public PresenceProjector(PresenceRepository presenceRepository, MusicianRepository musicianRepository, MusicEventRepository musicEventRepository, CommandGateway commandGateway) { this.presenceRepository = presenceRepository; this.musicianRepository = musicianRepository; this.musicEventRepository = musicEventRepository; this.commandGateway = commandGateway; } @EventHandler public void on(MusicianJoinDateCreatedEvent event) { if (!event.isMusicianActive()) { return; } var musicianId = event.getMusicianId(); var joinDate = event.getCurrentJoinDate(); var musicEvents = musicEventRepository.findAllAfter(joinDate.atStartOfDay()); for (MusicEvent musicEvent : musicEvents) { var presenceId = UUID.randomUUID(); commandGateway.send( new CreatePresenceSlotCommand( presenceId, musicEvent.getId(), musicianId, Boolean.FALSE)); } } @EventHandler public void on(MusicianJoinDateUpdatedEvent event) { var musicianId = event.getMusicianId(); var earlier = event.getCurrentJoinDate().isBefore(event.getNewJoinDate()) ? event.getCurrentJoinDate() : event.getNewJoinDate(); var further = event.getCurrentJoinDate().isAfter(event.getNewJoinDate()) ? event.getCurrentJoinDate() : event.getNewJoinDate(); if (event.getCurrentJoinDate().isBefore(event.getNewJoinDate())) { var musicEvents = musicEventRepository.findAllBetween(earlier.atStartOfDay(), further.plusDays(1).atStartOfDay()); for (MusicEvent musicEvent : musicEvents) { var presenceId = UUID.randomUUID(); commandGateway.send( new CreatePresenceSlotCommand( presenceId, musicEvent.getId(), musicianId, Boolean.FALSE)); } } else { var presenceList = presenceRepository.findAllBetween(earlier, further); for(Presence p : presenceList) { commandGateway.send (new RemovePresenceCommand(p.getId())); } } } @EventHandler public void on(MusicEventDateCreated event) { var musicEventId = event.getMusicEventId(); var eventDateTime = event.getDateTime(); var musicians = musicianRepository.findAllActiveOnMusicEventTime(eventDateTime.toLocalDate()); for (Musician musician : musicians) { var presenceSlotId = UUID.randomUUID(); log.info("PresenceProjector@EventHandler.onMusicEventCreatedEvent - sending CreatePresenceEntryCommand"); commandGateway.send( new CreatePresenceSlotCommand( presenceSlotId, musicEventId, musician.getId(), Boolean.FALSE)); } } @EventHandler public void on(MusicEventDateUpdated event) { var musicEventId = event.getMusicEventId(); var earlier = event.getCurrentDateTime().isBefore(event.getNewDateTime()) ? event.getCurrentDateTime() : event.getNewDateTime(); var further = event.getCurrentDateTime().isAfter(event.getNewDateTime()) ? event.getCurrentDateTime() : event.getNewDateTime(); if (event.getCurrentDateTime().isBefore(event.getNewDateTime())) { var musicians = musicianRepository.findAllActiveOnMusicEventTime(event.getNewDateTime().toLocalDate()); for (Musician musician : musicians) { var presenceSlotId = UUID.randomUUID(); log.info("PresenceProjector@EventHandler.onMusicEventCreatedEvent - sending CreatePresenceEntryCommand"); commandGateway.send( new CreatePresenceSlotCommand( presenceSlotId, musicEventId, musician.getId(), Boolean.FALSE)); } } else { var presenceList = presenceRepository.findAllBetween(earlier.toLocalDate(), further.plusDays(1).toLocalDate()); for(Presence p : presenceList) { commandGateway.send (new RemovePresenceCommand(p.getId())); } } } @EventHandler public void on(PresenceSlotCreatedEvent event) { log.info("PresenceProjector@EventHandler.onPresenceCreatedEvent: " + event); var musician = musicianRepository.findById(event.getMusicianId()); var musicianFullName = musician.get().getFullName(); var p = new Presence(event.getPresenceId(), event.getMusicianId(), musicianFullName, event.getEventId()); presenceRepository.save(p); } @EventHandler public void on(PresenceRemovedEvent event) { presenceRepository.deleteById(event.getId()); } @EventHandler public void on(MusicianMusicEventPresenceChangedEvent event) { Presence presence = presenceRepository.findForMusicianIdAndEventId(event.getMusicianId(), event.getEventId()); presence.setPresent(event.isPresent()); presenceRepository.save(presence); } @QueryHandler public List<PresenceDto> handle(GetMusicianPresenceInMusicEventQuery query) { return presenceRepository.findAllForEventId(query.getMusicEventId()); } }
核心问题
你当前的问题根源在于在查询端的Projector中发送命令。Axon的Projector属于查询模型组件,职责是监听事件并构建用于查询的视图数据。每次应用重启时,Axon会重放所有已存储的事件,Projector会重新处理每一个事件,导致里面的commandGateway.send()被重复执行,进而重复创建PresenceSlot。
正确的建模思路
命令的发送应该来自用户交互触发或者业务流程的协调组件,而不是事件重放过程中的查询端处理器。下面提供几种可行的修正方案:
方案1:使用Saga协调关联逻辑
Saga是Axon中用于协调跨聚合业务流程的组件,它拥有状态,可以跟踪已经处理过的关联关系,避免重复触发命令。
实现步骤:
- 创建
PresenceSaga类,监听MusicianJoinDateCreatedEvent、MusicianJoinDateUpdatedEvent、MusicEventDateCreatedEvent、MusicEventDateUpdatedEvent这些事件。 - Saga内部维护一个状态表(比如用数据库存储),记录已经关联过的
musicianId-eventId对,确保同一个组合不会被重复处理。 - 当事件触发时,Saga先检查是否已处理过该关联,未处理则发送
CreatePresenceSlotCommand或RemovePresenceCommand,并标记该关联为已处理。
示例Saga结构:
@Saga @Slf4j public class PresenceSaga { @Autowired private transient CommandGateway commandGateway; @Autowired private transient PresenceAssociationRepository associationRepository; @Autowired private transient MusicEventRepository musicEventRepository; @StartSaga @SagaEventHandler(associationProperty = "musicianId") public void handle(MusicianJoinDateCreatedEvent event) { if (!event.isMusicianActive()) { return; } var musicianId = event.getMusicianId(); var joinDate = event.getCurrentJoinDate(); var musicEvents = musicEventRepository.findAllAfter(joinDate.atStartOfDay()); for (MusicEvent musicEvent : musicEvents) { var associationId = buildAssociationId(musicianId, musicEvent.getId()); if (!associationRepository.existsById(associationId)) { var presenceId = UUID.randomUUID(); commandGateway.send(new CreatePresenceSlotCommand(presenceId, musicEvent.getId(), musicianId, Boolean.FALSE)); associationRepository.save(new PresenceAssociation(associationId, musicianId, musicEvent.getId())); } } } // 其他事件处理方法类似,先检查关联是否存在再执行命令 private String buildAssociationId(UUID musicianId, UUID eventId) { return musicianId + "-" + eventId; } }
方案2:在PresenceSlot聚合中避免重复创建
修改PresenceSlotAggregate的CreatePresenceSlotCommand处理器,先检查是否已经存在相同musicianId和eventId的PresenceSlot,如果存在则直接返回,不重复创建。
实现步骤:
- 在
PresenceSlotAggregate中注入Repository(或者用查询模型检查),处理CreatePresenceSlotCommand时先查询是否有相同的musicianId和eventId记录。 - 如果不存在,再应用
PresenceSlotCreatedEvent;如果存在,忽略该命令或者抛出已存在的异常。
示例:
@Aggregate public class PresenceSlotAggregate { @AggregateIdentifier private UUID presenceId; private UUID musicianId; private UUID eventId; private Boolean present; // 空构造函数 @CommandHandler public void handle(CreatePresenceSlotCommand command, PresenceSlotRepository repository) { // 检查是否已存在相同的musician-event组合 if (repository.existsByMusicianIdAndEventId(command.getMusicianId(), command.getEventId())) { return; // 或抛出异常 } apply(new PresenceSlotCreatedEvent(command.getPresenceId(), command.getMusicianId(), command.getEventId(), command.isPresent())); } // 事件溯源处理器 }
方案3:将PresenceSlot作为查询模型动态生成
如果出勤槽的核心只是用于UI展示,且出勤状态可以单独存储,可以考虑不提前创建PresenceSlot聚合,而是在查询时动态生成:
- 查询某个音乐活动的出勤列表时,先获取所有符合条件的音乐家(加入日期早于活动日期且活跃)。
- 再查询每个音乐家在该活动中的出勤状态,组合成出勤槽数据返回给UI。
- 当用户标记出勤时,直接发送
UpdatePresenceStatusCommand,更新对应的状态记录(用musicianId-eventId作为唯一键)。
这种方式不需要提前创建PresenceSlot,从根源上避免了重复创建的问题。
关键修正步骤
无论选择哪种方案,首先要做的是:
- 移除
PresenceProjector中所有发送命令的代码,让它只专注于更新查询模型(比如处理PresenceSlotCreatedEvent、PresenceRemovedEvent来维护Presence实体)。 - 确保命令的发送逻辑只在首次触发业务操作时执行,而不是事件重放时。
内容的提问来源于stack exchange,提问作者Michał Bzowski

