【发布时间】:2021-01-07 15:15:58
【问题描述】:
为什么我需要这个:
我正在实施一个系统来安排集群上的虚拟机。 虚拟机从集群请求资源,我负责将给定的 RAM 和 CPU 调度到一个且只有一个虚拟机。我想保证这一点的唯一方法是逐个处理请求。
架构
创建 VM 的请求会发布到 requests 主题(时间线上方)。集群状态(已用/总资源)作为一系列更新存储在 cluster 主题(下)中。
@some-time 就像事件时间戳。
requests 和cluster 主题是基于cluster_id 进行分区的,因此对同一个集群的请求将按顺序排列,并且可以按顺序处理。我正在使用 Kafka Streams。
问题
如果请求之间的间隔至少为 50-100 毫秒,我很好。
但是。 假设有一些连续(在几毫秒内)创建 VM 的请求
如果我将来自 requests 的事件作为 KStream 使用并与 cluster KTable 加入它们,并在调度 VM 后将新的集群状态发布到 cluster,那么第二个请求将不会看到此更新,因为它比集群更新事件来得更快(并且读取第二个请求比推送集群更新然后使用它更快)。
我想要什么
每个请求都会看到前一个请求的集群更新。无论是通过延迟请求处理还是任何其他方式,这都是我想要的。
如何做到这一点?
希望卡夫卡已经有机制做类似的事情,你可以指点我!
以下是我的猜测:
-
将元数据添加到
requests和cluster主题。即cluster中的事件将包含last_request_id——最后处理的请求。last_request_id也将存储在线程局部变量中并传递给下一个请求。使用last_request_id丰富请求并转发到新的delayed-requests主题。 然后可能可以在last_request_id上加入cluster和delayed-requests并进行处理。 -
使用有关给定分区中集群的数据创建实例前瞬态状态存储(in-mem?)。请求读取和写入此存储,以及发布到
cluster主题 - 持久存储。在启动状态存储从cluster主题发起。
更新
看看this question,会尝试,希望这对我有用
【问题讨论】:
标签: apache-kafka apache-kafka-streams