【发布时间】:2019-01-24 04:08:31
【问题描述】:
这是我们物联网平台中的入站消息流:
Device ---(MQTT)---> RabbitMQ Broker ---(AMQP)---> Apache Storm ---> Kafka
我正在寻求实施一种解决方案,该解决方案可以有效地限制/限制每个客户端每秒发布到 Kafka 的数据量。
当前的策略是利用 Guava 的 RateLimiter,每个设备都有自己的本地缓存实例。当接收到设备消息时,从缓存中获取映射到该 deviceId 的 RateLimiter 并调用 tryAquire() 方法。如果成功获得许可,则元组照常转发到 Kafka,否则,超出配额并静默丢弃消息。这种方法相当繁琐,在某些时候注定会失败或成为瓶颈。
我一直在阅读 Kafka 的字节速率配额,并相信这在我们的案例中会非常有效,尤其是因为 Kafka 客户端可以动态配置。在我们的平台中创建虚拟设备后,应在client.id == deviceId 的位置添加一个新的client.id。
假设以下用例为例:
- 管理员创建 2 个虚拟设备:湿度和温度传感器
- 触发规则以在 Kafka 中为上述设备创建新的 user/clientId 条目
- 通过 Kafka CLI 设置生产者配额值
- 两个设备都发出入站事件消息
- ...?
这是我的问题。如果使用单个 Producer 实例,是否可以在调用 send() 之前在 ProducerRecord 或 Producer 中的某个位置指定 client.id?如果一个生产者只允许一个client.id,这是否意味着每个设备都必须有自己的生产者?如果只允许一对一的映射,那么缓存潜在的数百个(如果不是数千个)Producer 实例是否明智,每个设备一个实例?有没有更好的方法我还不知道?
注意:我们的平台是一个“开门系统”,这意味着客户永远不会收到错误响应,例如“Rate Exceeded”或任何与此相关的错误。这对最终用户来说都是透明的。出于这个原因,我不能干扰 RabbitMQ 中的数据或将消息重新路由到不同的队列。我集成这些东西的唯一选择在于 Storm 或 Kafka 之间。
【问题讨论】:
标签: apache-kafka apache-storm messaging throttling