【问题标题】:How to get all aggregates with Axon framework?如何使用 Axon 框架获取所有聚合?
【发布时间】:2017-11-13 07:06:19
【问题描述】:

我开始使用 Axon 框架并遇到了一些障碍。

虽然我可以使用它们的 ID 加载单个聚合,但我不知道如何获取所有聚合的列表或所有聚合 ID 的列表。

EventSourcingRepository 类只有返回一个聚合的 load() 方法。

有没有办法对所有聚合 (ID) 进行聚合,或者我是否应该在轴突之外保留所有聚合 ID 的列表?

为了简单起见,我现在只使用InMemoryEventStorageEngine。 我正在使用 Axon 3.0.7。

【问题讨论】:

    标签: axon


    【解决方案1】:

    首先我想知道您为什么要从Repository 中检索所有聚合的完整列表。 设置了Repository 接口,以便您可以加载Aggregate 来处理命令或创建新的Aggregate。

    问你的问题,我几乎猜你是用它来查询而不是命令处理。 然而,这不是EventSourcingRepository 的预期用途。

    我能想到的你想要这个的一个原因是,你想实现一个 API 调用来向你的应用程序中特定类型的所有 Aggregates 发布一个命令。 考虑到这种情况,是的,您需要自己存储 aggregateId 引用。

    但以我之前的问题结束:为什么要通过Repository 接口检索聚合列表?

    答案更新

    关于您的评论,我在回答中添加了以下内容:

    Axon 可帮助您在设置应用程序时考虑到事件溯源以及 CQRS(命令查询职责分离)。 因此,这意味着您的应用程序的命令端和查询端是​​分开的。

    Aggregate Repository 是应用程序的命令端,您可以在其中请求执行操作。 因此,它不提供聚合列表,因为命令是对 a 聚合的意图表达。因此,它只需要Repository 用户检索一个聚合或创建一个聚合。

    您需要的聚合列表示例是应用程序的查询端。 查询端(您的视图/实体)通常根据事件(通过事件获取)进行更新。 对于您的应用程序中的任何查询要求,您通常会引入一个针对您的需求量身定制的单独视图。

    在您的示例中,这意味着您将引入一个事件处理组件,监听您的聚合事件,该组件使用聚合的查询模型更新存储库。

    【讨论】:

    • 感谢您的回答。由于我对 Axon 和事件溯源一般都是新手,因此我对它们有“错误”的想法并非不可能。我不一定想通过Repository 接口检索聚合列表。这正是我期望这成为可能的地方。我需要一个聚合列表,以便能够在 UI 中选择它们,然后编辑它们中的每一个。
    • 我想了很多,因此我给出了答案的形式:) Axon 可以帮助您在设置应用程序时考虑到事件溯源,但也考虑到 CQRS。因此,这意味着您的应用程序的命令和查询端被分开了。聚合repository 是应用程序的命令端,您可以在其中请求执行操作。然而,聚合列表的需要是应用程序的查询端,它会根据事件进行更新。因此,您通常会为 any UI 视图提供一个单独的存储库,在您的示例中是聚合列表。
    • 我已经调整了答案以反映我的评论。
    【解决方案2】:

    传递给EventSourcingRepository 的EventStore 实现了StreamableMessageSource<M extends Message<?>>,这是一种访问聚合的方法。

    虽然使用事件处理组件的框架方式可能会更好地扩展(取决于它的使用方式/上下文),但我很确定事件处理组件无论如何都是由StreamableMessageSource<M extends Message<?>> 驱动的。因此,如果我们想跳过框架并直接进入,我们可以这样做:

        List<String> aggregates(StreamableMessageSource<Message<?>> eventStore) {
            return immediatelyAvailableStream(eventStore.openStream(
                    eventStore.createTailToken() /* All events in the event store */
            ))
                    .filter(e -> e instanceof DomainEventMessage)
                    .map(e -> (DomainEventMessage) e)
                    .map(DomainEventMessage::getAggregateIdentifier)
                    .distinct()
                    .collect(Collectors.toList());
        }
    
        /*
            Note that the stream returned by BlockingStream.asStream() will block / won't terminate
            as it waits for future elements.
         */
        static <M> Stream<M> immediatelyAvailableStream(final BlockingStream<M> messageStream) {
            Iterator<M> iterator = new Iterator<M>() {
                @Override
                public boolean hasNext() {
                    return messageStream.hasNextAvailable();
                }
    
                @Override
                public M next() {
                    try {
                        return messageStream.nextAvailable();
                    } catch (InterruptedException e) {
                        Thread.currentThread().interrupt();
                        throw new IllegalStateException("Didn't expect to be interrupted");
                    }
                }
            };
    
            Spliterator<M> spliterator = Spliterators.spliteratorUnknownSize(iterator, Spliterator.ORDERED);
            Stream stream = StreamSupport.stream(spliterator, false);
            return (Stream)stream.onClose(messageStream::close);
        }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-04-08
      • 1970-01-01
      • 1970-01-01
      • 2021-12-11
      • 1970-01-01
      相关资源
      最近更新 更多