【问题标题】:How to control the number of messages that being emitted by Apache Kafka per a specific time?如何控制 Apache Kafka 在特定时间发出的消息数量?
【发布时间】:2020-12-23 10:51:36
【问题描述】:

我是 Apache Kafka 的新手,我正在尝试配置 Apache Kafka,使其尽可能多地接收来自生产者的消息,但它仅在每个特定时间向消费者发送配置的消息数量。 换句话说,如何将 Apache Kafka 配置为每“30 秒”仅发送“例如 50 条消息” 无论消息数量如何,都可以发送给消费者,并且在接下来的 30 秒内,它会从已兑现的消息中再获取 50 条消息,依此类推。

【问题讨论】:

  • 你为什么需要它?一般来说,这是个坏主意。在大多数情况下,最好尽快将消息发送到输出(以避免由于可能的应用崩溃而导致数据丢失)。
  • 我将使用消耗的消息与每次请求数量有限的 3rd 方提供商进行通信。因此,我使用 Kafka 作为兑现,定期从中提取数据,如果有两个消费者实例,Kafka 解决了两次提取相同数据的问题

标签: java spring-boot apache-kafka spring-kafka


【解决方案1】:

如果您可以控制消费者

您可以使用max.poll.records 属性来限制每个poll() 方法调用的最大记录数。然后你只需要保证poll()在30秒内被调用一次。

一般来说,您可以查看所有可用的配置属性here

如果您无法控制消费者

那么您唯一的选择就是根据您的需求编写消息 - 在 30 秒内最多编写 50 条消息。没有可用的配置选项。只有您的应用程序逻辑才能做到这一点。

更新 - 如何控制确保调用 poll

最简单的方法是:

while (true) {
   consumer.poll()
   // .. do your stuff
   Thread.sleep(30000);
}

您可以通过测量处理时间使事情变得更复杂(即在poll 之后开始调用Thread.sleep() 以完全不等待超过 30 秒。

【讨论】:

  • 我可以控制消费者,我的问题变成了示例,如何确保在 30 秒内调用一次 poll()当我使用 spring-kafka 时?
  • 您在最初的问题中没有问过这个问题。请参阅对我的回答的评论:idleBetweenPolls 适合您。它在KafkaMessageListenerContainer 中所做的正是@VladislavVarslavans 在他上次编辑中所展示的。
  • 我不能在 c# 中使用 "max.poll.records" 是否有替代 c# github.com/confluentinc/confluent-kafka-dotnet/issues/1451
【解决方案2】:

生产者确实不向消费者发送消息的问题。在生产者放置消息的地方之间有一个持久的 Kafka 主题。而且它真的不在乎另一边是否有任何消费者。从消费者的角度来看也是如此:它只是主题数据的订阅者,并不关心另一端是否有一些生产者。因此,考虑到存在消息中间件的从消费者到生产者的背压是错误的方向。

另一方面,尚不清楚这些消费消息如何影响您的第三方服务。关键是 Kafka 消费者每个分区都是单线程的。因此,来自一个分区的所有消息将(必须)在同一个线程中一一处理。这样您就不能向您的服务发送多条消息:只有在前一条已回复时才能发送下一条。所以,想一想:你的消费者应用程序怎么可能超过速率限制?

但是,如果您在消费者端有足够的分区和高并发性,那么您最终可能会收到来自不同线程的多个并行请求。为此,我建议查看速率限制器模式。这个库提供了一个很好的实现:https://resilience4j.readme.io/docs/ratelimiter。最好将消息保留在主题中,然后尝试以某种方式限制生产者。

总结:即使消费者方面不是您的项目,最好与该团队讨论如何改进他们的消费者。你做得很好:生产者向 Kafka 主题发送消息。你还能在这里做什么?

【讨论】:

  • 我完全知道有一个持久的 Kafka 主题,我的问题是如何控制消费者在 30 秒内仅调用 'poll()' 50 次或仅调用 'poll()' 一次在 30 秒内,但最多收到 50 条消息
  • 好的。请参阅消费者配置属性中的max.poll.records 和Spring Kafka 的idleBetweenPolls 中的idleBetweenPollsdocs.spring.io/spring-kafka/docs/current/reference/html/…。您在那部分的问题不清楚:卡夫卡不会向消费者发送消息。消费者执行拉动逻辑。我们也称它为delivery。 “发送”正是生产者部分。
【解决方案3】:

有趣的用例,不知道为什么需要它,但有两种可能的解决方案: 1. 为了保护集群,您可以使用配额,而不是消息量而是带宽吞吐量:https://kafka.apache.org/documentation/#design_quotas。 2. 如果您需要每个时间范围内的确切消息量,您可以在消费和暂停之间放置一个缓冲服务(速率限制器),将消息发布到消费主题。速率限制器可能会消耗下一个 50,然后暂停直到一分钟过去。由于重复消息,这将增加集群上使用的空间。您还需要注意如何暂停消费者,需要发送心跳,否则您将不断重新平衡您的消费者,即您不能只睡到下一分钟。这显然是在您无法控制最终消费者的情况下。

【讨论】:

    猜你喜欢
    • 2018-06-25
    • 1970-01-01
    • 2019-05-18
    • 2020-07-08
    • 1970-01-01
    • 2019-11-23
    • 2014-02-13
    • 1970-01-01
    相关资源
    最近更新 更多