【问题标题】:Kafka Producer QuotasKafka 生产者配额
【发布时间】: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。

假设以下用例为例:

  1. 管理员创建 2 个虚拟设备:湿度和温度传感器
  2. 触发规则以在 Kafka 中为上述设备创建新的 user/clientId 条目
  3. 通过 Kafka CLI 设置生产者配额值
  4. 两个设备都发出入站事件消息
  5. ...?

这是我的问题。如果使用单个 Producer 实例,是否可以在调用 send() 之前在 ProducerRecord 或 Producer 中的某个位置指定 client.id?如果一个生产者只允许一个client.id,这是否意味着每个设备都必须有自己的生产者?如果只允许一对一的映射,那么缓存潜在的数百个(如果不是数千个)Producer 实例是否明智,每个设备一个实例?有没有更好的方法我还不知道?

注意:我们的平台是一个“开门系统”,这意味着客户永远不会收到错误响应,例如“Rate Exceeded”或任何与此相关的错误。这对最终用户来说都是透明的。出于这个原因,我不能干扰 RabbitMQ 中的数据或将消息重新路由到不同的队列。我集成这些东西的唯一选择在于 Storm 或 Kafka 之间。

【问题讨论】:

    标签: apache-kafka apache-storm messaging throttling


    【解决方案1】:

    您可以通过应用程序配置client.idproperties.put ("client.id", "humidity")properties.put ("client.id", "temp") 根据每个client.id可以设置值

    producer_byte_rate = 1024, consumer_byte_rate = 2048,
    request_percentage = 200
    

    怀疑我是否与此配置有关 (producer_byte_rate = 1024, consumer_byte_rate = 2048, request_percentage = 200),生产者不会假设插入的配置,因为消费者工作正常

    【讨论】:

      【解决方案2】:

      虽然您可以在 Producer 对象上指定 client.id,但请记住它们是重量级的,您可能不愿意创建它们的多个实例(尤其是在每个设备一个的基础上)。

      关于减少Producer 的数量,您是否考虑过为每个用户而不是每个设备创建一个,或者甚至有一个有限的共享池?然后可以使用 Kafka 消息头来识别实际生成数据的设备。缺点是您需要限制自己的消息生成,这样一个设备就不会从其他设备获取所有资源。

      但是,您可以限制 Kafka 代理端的用户,并将配置应用于默认用户/客户端:

      > bin/kafka-configs.sh  --zookeeper localhost:2181 --alter --add-config 'producer_byte_rate=1024,consumer_byte_rate=2048,request_percentage=200' --entity-type clients --entity-default
      Updated config for entity: default client-id.
      

      有关更多示例和深入解释,请参阅 https://kafka.apache.org/documentation/#design_quotas

      如何识别消息取决于您的架构,可能的解决方案包括:

      【讨论】:

      • “然后可以使用Kafka消息头来识别实际产生数据的设备。”请您详细说明一下并解释如何实现它?即使我是针对每个用户进行的,我仍然需要弄清楚如何告诉 Kafka 消息 X 来自客户端 1,消息 Y 来自客户端 2,等等。所有这些都通过 单个,共享生产者实例。
      猜你喜欢
      • 2023-03-03
      • 1970-01-01
      • 2019-11-30
      • 1970-01-01
      • 1970-01-01
      • 2018-02-05
      • 2013-01-23
      • 2019-10-31
      • 2016-01-15
      相关资源
      最近更新 更多