在 Kafka 之上并发写入事件源

Jes*_*spc 5 cqrs event-sourcing apache-kafka

我一直在考虑在事件源配置中使用 Apache Kafka 作为事件存储。发布的事件将与特定资源相关联,传递到与资源类型相关联的主题,并按资源 id 分片到分区中。因此,例如,创建类型为 Folder 且 ID 为 1 的资源将产生一个 FolderCreate 事件,该事件将传递到分区中的“文件夹”主题,该分区通过将 id 1 分片到主题中的分区总数来给出。即使我不知道如何处理使日志不一致的并发事件。

最简单的场景是有两个并发操作,它们可以使彼此无效,例如一个更新文件夹,一个销毁同一个文件夹。在这种情况下,该主题的分区最终可能包含无效序列 [FolderDestroy, FolderUpdate]。这种情况通常通过对事件进行版本控制来解决,如此处所述,但 Kafka 不支持此类功能。

在这些情况下,如何确保 Kafka 日志本身的一致性?

Tom*_*omW 4

我认为可能可以使用 Kafka 进行聚合(在 DDD 意义上)或“资源”的事件溯源。一些注意事项:

  1. 序列化每个分区的写入,使用每个分区(或多个分区)的单个进程来管理它。确保在同一个 Kafka 连接上串行发送消息,如果无法承受回滚,则在向命令发送者报告成功之前使用 ack=all 。确保生产者进程跟踪每个资源的当前成功事件偏移量/版本,以便它可以在发送消息之前自行进行乐观检查。
  2. 由于即使写入实际上成功,也可能返回写入失败,因此您需要重试写入并通过在每个事件中包含 ID 来处理重复数据删除,或者通过重新读取流(其中的最新消息)来重新初始化生产者以查看写入是否真正有效。
  3. 以原子方式编写多个事件 - 只需发布包含事件列表的复合事件。
  4. 按资源 ID 查找。这可以通过在启动时读取分区中的所有事件(或特定跨资源快照中的所有事件)并将当前状态存储在 RAM 中或缓存在数据库中来实现。

https://issues.apache.org/jira/browse/KAFKA-2260会以更简单的方式解决 1,但似乎陷入停滞。

Kafka Streams 似乎为您提供了很多这样的功能。例如,4 是一个 KTable,您可以让事件生成器在发送事件之前使用 1 来确定事件对于当前资源状态是否有效。