【问题标题】:Understanding Persistent Entities with streams of data了解具有数据流的持久实体
【发布时间】:2017-10-03 08:38:50
【问题描述】:

我想使用 Lagom 构建数据处理管道。此管道的第一步是使用 Twitter 客户端来订阅 Twitter 消息流的服务。对于每条新消息,我都希望将消息保留在 Cassandra 中。

我不明白的是,例如,我将 Aggregare 根建模为 TwitterMessages 列表,运行一段时间后,这个 agggregare 根的大小将达到数 GB。无需将所有 TwitterMessages 存储在内存中,因为该服务的目标只是持久化每个传入消息,然后将消息发布到 Kafka 以供下一个服务处理。

如何将聚合根建模为消息流的持久实体而不消耗无限资源?如果 Lagom,是否有任何示例代码显示这种用法?

【问题讨论】:

  • 关于该聚合的业务规则是什么?它保护的不变量是什么?
  • 没什么,它应该只是将每个传入的消息附加到数据库中。
  • 那么您不需要聚合,至少不需要事件来源的聚合。也许某种流处理器?
  • 我的想法是,使用 Lagom,您可以使用事件溯源,在我的情况下,每次传入消息到达时都会触发一些命令。该命令将触发一个包含传入消息的事件,然后将其保留。如果我使用流处理器,我将只使用常规的 Cassandra 客户端,并在传入流上以 forEach 类型的方式写入每条传入消息。我不应该使用事件溯源来满足我所有的持久性需求吗?
  • 我不明白这一点,因为您没有任何不变量需要保护。

标签: scala event-sourcing lagom


【解决方案1】:

事件溯源是一个很好的默认选择,但并不是所有事情的正确解决方案。在您的情况下,这可能不是正确的方法。首先,您需要将推文持久化,还是可以直接将它们发布到 Kafka?

假设您需要它们持久化,聚合应该存储在内存中验证传入命令和生成新事件所需的任何内容。根据您的描述,您的聚合不需要任何数据来执行此操作,因此您的聚合不会是 Twitter 消息列表,而是可能只是 NotUsed。每次收到命令时,它都会为该推文发出一个新事件。这里的问题是,它并不是真正的聚合,因为你没有聚合任何状态,你只是发出事件以响应没有不变量或任何东西的命令。因此,您并没有真正将 Lagom 持久实体 API 用于它的用途。尽管如此,无论如何以这种方式使用它可能是有意义的,它是一个高级 API,带有一些有用的东西,包括流功能。但是也有一些你应该注意的问题,你将所有的推文放在一个实体中,你将吞吐量限制在一个节点上的一个核心可以一次按顺序执行的操作。因此,也许您可​​以期望每秒处理 20 条推文,如果您期望它超过此数量,那么您使用了错误的方法,并且您至少需要将您的推文分发到多个实体。

另一种方法是直接将消息直接存储在 Cassandra 中,然后直接发布到 Kafka。这会简单得多,涉及的机制也少得多,而且它应该可以很好地扩展,只要确保您明智地选择 Cassandra 中的分区键列 - 我可能会按用户 ID 进行分区。

【讨论】:

  • 也许您可以添加一个示例,说明如何在没有 Persistent 模块的情况下将 Cassandra 与 Lagom 一起使用?
猜你喜欢
  • 2020-07-05
  • 2015-10-28
  • 2013-08-26
  • 1970-01-01
  • 1970-01-01
  • 2019-02-01
  • 2012-01-30
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多