【发布时间】:2019-03-12 14:09:40
【问题描述】:
是否可以使用usingSubscribingEventProcessors,并且在投影事件时,总是从头开始重新投影所有事件。含义 - 我从不将预测保存到数据库,而是在聚合发出新事件时重新投影所有事件?
【问题讨论】:
是否可以使用usingSubscribingEventProcessors,并且在投影事件时,总是从头开始重新投影所有事件。含义 - 我从不将预测保存到数据库,而是在聚合发出新事件时重新投影所有事件?
【问题讨论】:
当然可以!
但是,您无法通过使用订阅事件处理器来实现这一点。
您应该利用跟踪事件处理器,但在其后面有一个InMemoryTokenStore。这样做,应用程序永远无法从它停止的地方开始,因为它停止的地方的知识,TrackingToken,不存在。
因此,您最终会在每次启动时重新创建您的预测。
您可以采取的另一种方法略有不同。
您仍将使用跟踪事件处理器,但使用实际持久的TokenStore 实现。其次,在您的应用程序启动时,您可以使用TrackingEventProcessor#resetTokens() 函数重播给定的跟踪事件处理器。
采用这种方法,您可以在事件处理组件中添加@ResetHandler 注释函数,以在再次处理所有事件之前清除投影表。
希望这能给博扬一些见解!
【讨论】:
InMemoryTokenStore 和configurer.registerTokenStore("inMemoryProcessor", cfg -> inMemoryTokenStore()); 进行了设置,并使用@ProcessingGroup("inMemoryProcessor") 进行了注释投影仪。现在我可以看到跟踪事件处理器不断地轮询来自 DB 的事件。是否可以仅在调用 @QueryHandler 注释方法而不是之前进行投影?基本上,我发送例如 id=1 的请求,并且此方法知道在那个时间点它需要重新投影此聚合中的所有事件并返回当前投影(不保存到数据库)?
queryResult.updates().blockFirst() 更改为queryResult.initialResult().block(),它似乎返回了正确的投影结果。
FROM DomainEventEntry e WHERE e.globalIndex > :token。
@QueryHandler 带注释的函数上出现查询时才创建投影。不过,这将需要一些个人实施。在其内部深处,Axon 使用 AnnotationEventListenerAdapter 来包装 @EventHandler 带注释的类,以便能够使用事件流调用它。您可以将您的投影包装在其中,并在您的@QueryHandler 上从EventStore 中提取事件流并自己处理此过程。我以前做过,但只有你可以决定这个家庭作业是否可取。
@Steven 你觉得这个解决方案怎么样?
public class ReplayingSubscribingEventProcessor extends SubscribingEventProcessor {
private final SubscribableMessageSource<? extends EventMessage<?>> messageSource;
protected ReplayingSubscribingEventProcessor(
Builder builder) {
super(builder);
this.messageSource = builder.messageSource;
}
public static Builder builder() {
return new Builder();
}
/**
* Whenever there is a need to process event messages, ignore all of them and since already inside messageSource,
* just take all messages from event source and re-project all from beginning for this aggregate root
* @param eventMessages
*/
@Override
protected void process(List<? extends EventMessage<?>> eventMessages) {
try {
//reprocess all previous events for this aggregate (get id from current event)
GenericDomainEventMessage gdem = (GenericDomainEventMessage) eventMessages.get(0);
List<? extends EventMessage<?>> prevEvs = ((EventStore)messageSource).readEvents(gdem.getAggregateIdentifier()).asStream()
.collect(Collectors.toList());
processInUnitOfWork(prevEvs, new BatchingUnitOfWork<>(prevEvs), Segment.ROOT_SEGMENT);
} catch (RuntimeException e) {
throw e;
} catch (Exception e) {
throw new EventProcessingException("Exception occurred while processing events", e);
}
}
public static class Builder extends SubscribingEventProcessor.Builder{
private SubscribableMessageSource<? extends EventMessage<?>> messageSource;
@Override
public Builder messageSource(
SubscribableMessageSource<? extends EventMessage<?>> messageSource) {
super.messageSource(messageSource);
this.messageSource = messageSource;
return this;
}
@Override
public ReplayingSubscribingEventProcessor.Builder name(String name) {
super.name(name);
return this;
}
@Override
public ReplayingSubscribingEventProcessor.Builder eventHandlerInvoker(
EventHandlerInvoker eventHandlerInvoker) {
super.eventHandlerInvoker(eventHandlerInvoker);
return this;
}
@Override
public ReplayingSubscribingEventProcessor.Builder processingStrategy(
EventProcessingStrategy processingStrategy) {
super.processingStrategy(processingStrategy);
return this;
}
public ReplayingSubscribingEventProcessor build() {
return new ReplayingSubscribingEventProcessor(this);
}
}
}
和配置:
@Autowired
public void configure(EventProcessingConfigurer configurer){
configurer.registerEventProcessor("inMemoryProcessor",
(n, c, ehi) -> replayingSubscribingEventProcessor(n, c, ehi, org.axonframework.config.Configuration::eventBus));
}
public ReplayingSubscribingEventProcessor replayingSubscribingEventProcessor(
String name,
org.axonframework.config.Configuration conf,
EventHandlerInvoker eventHandlerInvoker,
Function<org.axonframework.config.Configuration, SubscribableMessageSource<? extends EventMessage<?>>> messageSource) {
return ReplayingSubscribingEventProcessor.builder()
.name(name)
.eventHandlerInvoker(eventHandlerInvoker)
.messageSource(messageSource.apply(conf))
.processingStrategy(DirectEventProcessingStrategy.INSTANCE)
.build();
}
【讨论】: