【问题标题】:How to align (delay processing) Kafka events from two topics (by inner property)? OR How to sequentially process events from two topics?如何对齐(延迟处理)来自两个主题的 Kafka 事件(通过内部属性)?或如何顺序处理来自两个主题的事件?
【发布时间】:2021-01-07 15:15:58
【问题描述】:

为什么我需要这个:

我正在实施一个系统来安排集群上的虚拟机。 虚拟机从集群请求资源,我负责将给定的 RAM 和 CPU 调度到一个且只有一个虚拟机。我想保证这一点的唯一方法是逐个处理请求。

架构

创建 VM 的请求会发布到 requests 主题(时间线上方)。集群状态(已用/总资源)作为一系列更新存储在 cluster 主题(下)中。 @some-time 就像事件时间戳。

requestscluster 主题是基于cluster_id 进行分区的,因此对同一个集群的请求将按顺序排列,并且可以按顺序处理。我正在使用 Kafka Streams。

问题

如果请求之间的间隔至少为 50-100 毫秒,我很好。

但是。 假设有一些连续(在几毫秒内)创建 VM 的请求

如果我将来自 requests 的事件作为 KStream 使用并与 cluster KTable 加入它们,并在调度 VM 后将新的集群状态发布到 cluster,那么第二个请求将不会看到此更新,因为它比集群更新事件来得更快(并且读取第二个请求比推送集群更新然后使用它更快)。

我想要什么

每个请求都会看到前一个请求的集群更新。无论是通过延迟请求处理还是任何其他方式,这都是我想要的。

如何做到这一点?

希望卡夫卡已经有机制做类似的事情,你可以指点我!

以下是我的猜测:

  • 将元数据添加到 requestscluster 主题。即cluster 中的事件将包含last_request_id——最后处理的请求。 last_request_id 也将存储在线程局部变量中并传递给下一个请求。使用last_request_id 丰富请求并转发到新的delayed-requests 主题。 然后可能可以在last_request_id 上加入clusterdelayed-requests 并进行处理。

  • 使用有关给定分区中集群的数据创建实例前瞬态状态存储(in-mem?)。请求读取和写入此存储,以及发布到 cluster 主题 - 持久存储。在启动状态存储从cluster 主题发起。

更新

看看this question,会尝试,希望这对我有用

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    使用 DSL,可能无法实现您想要的。不过我想你可以使用处理器 API:对于这种情况,你使用两个状态存储,一个用于“使用表”,另一个用于缓冲“供应请求”。

    每次处理“供应请求”时,您都会将“使用表”标记为已阻止(可能使用一些特殊键)。如果第二个“供应请求”进入并且表被“阻塞”,则将其缓冲在“请求缓冲存储”中。每次更新“使用表”时,您都会检查“请求缓冲存储”中是否有任何缓冲事件并处理一个请求。如果请求存储为空,您可以“解锁”该表。

    【讨论】:

      猜你喜欢
      • 2015-03-30
      • 1970-01-01
      • 1970-01-01
      • 2022-10-26
      • 1970-01-01
      • 1970-01-01
      • 2020-04-24
      • 2014-10-23
      • 1970-01-01
      相关资源
      最近更新 更多