【发布时间】:2017-10-03 08:38:50
【问题描述】:
我想使用 Lagom 构建数据处理管道。此管道的第一步是使用 Twitter 客户端来订阅 Twitter 消息流的服务。对于每条新消息,我都希望将消息保留在 Cassandra 中。
我不明白的是,例如,我将 Aggregare 根建模为 TwitterMessages 列表,运行一段时间后,这个 agggregare 根的大小将达到数 GB。无需将所有 TwitterMessages 存储在内存中,因为该服务的目标只是持久化每个传入消息,然后将消息发布到 Kafka 以供下一个服务处理。
如何将聚合根建模为消息流的持久实体而不消耗无限资源?如果 Lagom,是否有任何示例代码显示这种用法?
【问题讨论】:
-
关于该聚合的业务规则是什么?它保护的不变量是什么?
-
没什么,它应该只是将每个传入的消息附加到数据库中。
-
那么您不需要聚合,至少不需要事件来源的聚合。也许某种流处理器?
-
我的想法是,使用 Lagom,您可以使用事件溯源,在我的情况下,每次传入消息到达时都会触发一些命令。该命令将触发一个包含传入消息的事件,然后将其保留。如果我使用流处理器,我将只使用常规的 Cassandra 客户端,并在传入流上以 forEach 类型的方式写入每条传入消息。我不应该使用事件溯源来满足我所有的持久性需求吗?
-
我不明白这一点,因为您没有任何不变量需要保护。
标签: scala event-sourcing lagom