【问题标题】:How can I instantly get the result of my produced event in kafka and kafka-streams?如何立即获得我在 kafka 和 kafka-streams 中制作的事件的结果?
【发布时间】:2020-02-23 21:02:07
【问题描述】:

我正在通过以下场景简化我的问题: 3 个朋友共享一张会员卡。该卡有两个限制

  1. 最多可以使用10次(不管哪个用卡,即 friend_a 可以使用 10 次。
  2. 卡中的最大金额为 200。因此,如果 1 个“事件”的价值 = 200,则卡已“完成”。

我正在使用一个 kafka 生产者,它像这样在 kafka 集群中发送事件

{ "name": "friend_1", "value": 10 }

{ "name": "friend_3", "value": 20 }

这些事件被发布到一个与 kafka 流连接的主题,该流按键分组并进行聚合以汇总所花费的钱。这似乎可行,但是我面临“并发问题”

让我们想象一下这张卡使用了 9 次,所以只剩下 1 次可以使用,总共花费了 190,也就是说还有 10 个单位可以花费。

因此,朋友_2 想购买价格为 11 单位的东西(不应该被允许),而朋友_3 想购买价格为 9 单位的东西,而这应该是允许的。 Friend_3 将第 10 次修改使用卡片的状态。所有其他未来的尝试都不应修改任何内容。

因此卡用户知道他发送的事件是否修改了最大使用数和总计数似乎是合理的。我怎么能在卡夫卡做到这一点?使用流聚合我总是可以增加值,但我怎么知道我的操作是否“修改了卡片的状态”?

更新:如果交易验证了规则,卡用户应该立即得到负面反馈。

【问题讨论】:

  • 如果以下答案之一解决了您的问题,请将其标记为已接受。

标签: apache-kafka kafka-consumer-api apache-kafka-streams kafka-producer-api


【解决方案1】:

根据我对您问题的理解,有几个选项。

一种选择是在聚合之后派生一个新流,您 filter() 用于修改卡片“状态”的数据,例如过滤所有已花费 > 200 单位或 > 10 使用的事件。然后,此流可用于通知卡用户该卡已被使用,例如通过发送电子邮件。这种方法可以单独使用 DSL 来实现。

当需要更大的灵活性或更严格的控制时,另一种选择是使用处理器 API(您可以将其与 DSL 集成,以便您的大部分代码可以继续使用 DSL),您可以自己实现聚合步骤(使用您附加到TransformerProcessor 的状态存储)。在聚合期间,您可以实现检查传入事件是否有效的逻辑(在您的示例中:9 个单位的朋友_3 有效,11 个单位的朋友_2 无效)。如果它是有效的,聚合会增加卡的计数器(单位和用途),就是这样。如果它无效,该事件将被丢弃并且不会修改计数器,Transformer/Processor 可以向另一个流发出一个新事件,告诉卡用户某些事情不起作用。您可以类似地实现该功能以通知用户一张卡已被完全使用、一张卡不再可用或该卡的任何其他“状态更改”。

另外,根据您想要做什么,看看 Kafka Streams 的 interactive queries 功能。有时其他应用程序可能想要对某物的最新状态(例如卡片的状态)进行快速点查找(查询),这可以通过交互式查询来完成,例如一个 REST API。

希望这会有所帮助!

【讨论】:

  • 嗨,Michael,我的案例要求在交易成功或“不成功”时立即回复卡用户。电子邮件通知,即使是几秒钟后,也不是这种情况。
  • 当您使用事件流设置(如 Kafka)时,您选择的架构主要是异步而不是同步。因此,我建议您重新考虑您的“立即响应”要求。 (是同步请求/响应吗?是否可以接受一些异步延迟,如果可以,多少?)
猜你喜欢
  • 2020-07-30
  • 2020-03-24
  • 2020-04-01
  • 1970-01-01
  • 2022-08-19
  • 1970-01-01
  • 2017-12-22
  • 1970-01-01
相关资源
最近更新 更多