【问题标题】:Event sourcing with Kafka streams使用 Kafka 流进行事件溯源
【发布时间】:2018-08-31 01:18:06
【问题描述】:

我正在尝试在 Kafka streams 之上实现一个简单的 CQRS/事件溯源概念证明(如 https://www.confluent.io/blog/event-sourcing-using-apache-kafka/ 中所述)

我有 4 个基本部分:

  1. commands 主题,使用聚合 ID 作为键,对每个聚合的命令进行顺序处理
  2. events 主题,聚合状态的每次更改都会发布到该主题(同样,key 是聚合 ID)。此主题的保留政策是“永不删除”
  3. KTable 减少聚合状态并将其保存到状态存储中

    事件主题流-> 按聚合 ID 分组到 Ktable -> 将聚合事件减少到当前状态 -> 实体化为国营商店
  4. 命令处理器 - 命令流,左连接聚合状态 KTable。对于结果流中的每个条目,使用函数 (command, state) => events 生成结果事件并将它们发布到 events 主题

问题是 - 有没有办法确保我在状态存储中拥有最新版本的聚合?

如果违反业务规则,我想拒绝命令(例如 - 如果实体被标记为已删除,则修改实体的命令无效)。但是,如果发布了DeleteCommand,紧随其后的是ModifyCommand,则删除命令将生成DeletedEvent,但是在处理ModifyCommand 时,来自状态存储的加载状态可能还没有反映出来,并且将发布冲突事件。

我不介意牺牲命令处理吞吐量,我宁愿获得一致性保证(因为所有内容都按相同的键分组并且应该最终在同一个分区中)

希望这很清楚 :) 有什么建议吗?

【问题讨论】:

    标签: apache-kafka event-sourcing apache-kafka-streams


    【解决方案1】:

    我不认为 Kafka 对 CQRS 和事件溯源有好处,正如您所描述的那样,因为它缺乏一种(简单的)方法来确保防止并发写入。这个article 详细讨论了这个问题。

    我的意思是你描述它的方式是你期望一个命令生成零个或多个事件或失败并出现异常的事实;这是带有事件溯源的经典 CQRS。大多数人都期待这种架构。

    您可以使用不同的样式进行事件溯源。您的命令处理程序可以为接收到的每个命令产生事件(即DeleteWasAccepted)。然后,事件处理程序最终可以以事件来源的方式处理该事件(通过从其事件流重建聚合的状态)并发出其他事件(即ItemDeletedItemDeletionWasRejected)。因此,命令被触发并忘记,异步发送,客户端不等待立即响应。然而,它等待描述其命令执行结果的事件。

    一个重要的方面是事件处理程序必须以串行方式(完全一次并按顺序)处理来自同一聚合的事件。这可以使用单个 Kafka 消费者组来实现。您可以在 video 中了解此架构。

    【讨论】:

      【解决方案2】:

      请阅读我的同事 Jesper 的这篇文章。 Kafka 是一款很棒的产品,但实际上根本不适合事件溯源

      https://medium.com/serialized-io/apache-kafka-is-not-for-event-sourcing-81735c3cf5c

      【讨论】:

      • 我读过,他提出了很好的观点。但是,它们不适用于我描述的使用 Kafka 流的设计: 加载当前状态 - 使用 KTable 完成;一致的写入 - 由 Kafka 的分区模型和容错保证处理
      【解决方案3】:

      我想出的一个可能的解决方案是实现一种乐观锁定机制:

      1. 在命令中添加expectedVersion 字段
      2. 使用 KTable Aggregator 为每个处理的事件增加聚合快照的版本
      3. 如果expectedVersion 与快照的聚合版本不匹配,则拒绝命令

      这似乎提供了我正在寻找的语义

      【讨论】:

      • 很高兴你找到了你想要的东西,但这不是乐观锁定的工作原理。该命令应该重试,而不是拒绝。
      • 重试一个命令跟乐观锁没关系...确实乐观锁不是在数据库的事务管理器中做的,但是观察到的效果是一样的——如果是预期的版本不匹配的更改将不会被保留。由于Kafka按分区顺序处理消息,因此不应该有并发写入en.wikipedia.org/wiki/Optimistic_concurrency_control
      • 如果命令由于并发写入而被拒绝,您的客户会看到什么?
      • 由于基础架构或架构不适合,您无法执行有效命令。
      • 定义如何定义“有效”。如果客户端正在使用旧版本,您可以选择拒绝该命令并要求客户端重试。或者,您可以选择通过将命令重新添加到当前版本的命令主题来自动重试该命令。其实很简单
      猜你喜欢
      • 2019-12-13
      • 2016-06-02
      • 2019-02-03
      • 2019-06-13
      • 2020-05-27
      • 1970-01-01
      • 1970-01-01
      • 2022-08-19
      • 1970-01-01
      相关资源
      最近更新 更多