【发布时间】:2018-08-31 01:18:06
【问题描述】:
我正在尝试在 Kafka streams 之上实现一个简单的 CQRS/事件溯源概念证明(如 https://www.confluent.io/blog/event-sourcing-using-apache-kafka/ 中所述)
我有 4 个基本部分:
-
commands主题,使用聚合 ID 作为键,对每个聚合的命令进行顺序处理 -
events主题,聚合状态的每次更改都会发布到该主题(同样,key 是聚合 ID)。此主题的保留政策是“永不删除” -
KTable 减少聚合状态并将其保存到状态存储中
事件主题流-> 按聚合 ID 分组到 Ktable -> 将聚合事件减少到当前状态 -> 实体化为国营商店 命令处理器 - 命令流,左连接聚合状态 KTable。对于结果流中的每个条目,使用函数
(command, state) => events生成结果事件并将它们发布到events主题
问题是 - 有没有办法确保我在状态存储中拥有最新版本的聚合?
如果违反业务规则,我想拒绝命令(例如 - 如果实体被标记为已删除,则修改实体的命令无效)。但是,如果发布了DeleteCommand,紧随其后的是ModifyCommand,则删除命令将生成DeletedEvent,但是在处理ModifyCommand 时,来自状态存储的加载状态可能还没有反映出来,并且将发布冲突事件。
我不介意牺牲命令处理吞吐量,我宁愿获得一致性保证(因为所有内容都按相同的键分组并且应该最终在同一个分区中)
希望这很清楚 :) 有什么建议吗?
【问题讨论】:
标签: apache-kafka event-sourcing apache-kafka-streams