【发布时间】:2019-01-22 22:29:45
【问题描述】:
上下文
我们使用 Kafka 来处理大型消息,偶尔会达到 10MB,但大多在 500 KB 范围内。处理一条消息最多可能需要 30 秒左右,但有时需要一分钟。
问题
使用较少数量的消费者(最多约 50 个)处理数据不会导致经纪人重复进行重新平衡,并且处理工作正常。这种规模的任何重新平衡也相当快,根据代理日志,通常不到一分钟。
一旦消费者规模扩大到 100 或 200,消费者就会不断地重新平衡,间隔最长约为 5 分钟。这导致 5 分钟的工作/消耗,然后是 5 分钟的重新平衡,然后再次相同。消费者服务不会失败,只是在没有明显原因的情况下重新平衡。当扩大消费者规模时,这会导致吞吐量降低。
当扩展到 2oo 个消费者时,处理以每个消费者每分钟 2 条消息的平均速率执行。单个消费者在不重新平衡时的处理速度约为每分钟 6 条消息。
我不怀疑数据中心的网络是个问题,因为我们有一些消费者对消息执行不同类型的处理,他们每分钟传递 100 到 1000 条消息没有问题。
其他人是否经历过这种模式并找到了一个简单的解决方案,例如更改特定的配置参数?
其他信息
Kafka 代理是 2.0 版,其中有 10 个跨不同的数据中心。 Replication 设置为 3。此主题的分区为 500。特定 broker 配置的摘录,以适应更好地处理大型消息的情况:
- compression.type=lz4
- message.max.bytes=10000000 # 10 MB
- replica.fetch.max.bytes=10000000 # 10 MB
- group.max.session.timeout.ms=1320000 # 22 分钟
- offset.retention.minutes=10080 # 7 天
在消费者方面,我们使用带有重新平衡侦听器的 java 客户端,该侦听器清除来自已撤销分区的所有缓冲消息。这个缓冲区有 10 条消息。消费者客户端运行客户端 API 版本 2.1,Java 客户端从 2.0 到 2.1 的更新似乎显着减少了这些较大消费者数量上的以下类型的代理日志(我们为几乎每个客户端和之前的每次重新平衡都得到了这些):
INFO [GroupCoordinator 2]: Member svc002 in group my_group has failed, removing it from the group (kafka.coordinator.group.GroupCoordinator)
消费者与代理位于不同的数据中心。偏移量的提交是异步执行的。循环轮询在一个线程中执行,该线程以 15 秒的超时时间填充缓冲区;一旦缓冲区已满,线程就会休眠几秒钟,并且仅在缓冲区有可用空间时才进行轮询。大消息用例的配置摘录:
- max.partition.fetch.bytes.config=200000000 # 200 MB
- max.poll.records.config=2
- session.timeout.ms.config=1200000 # 20 分钟
日志文件
以下是在 30 分钟时间范围内管理此特定组的代理日志文件的摘录。命名简化为 my_group 和 mytopic。还有一些来自不相关主题的条目。
19:47:36,786] INFO [GroupCoordinator 2]: Stabilized group my_group generation 3213 (__consumer_offsets-18) (kafka.coordinator.group.GroupCoordinator)
19:47:36,810] INFO [GroupCoordinator 2]: Assignment received from leader for group my_group for generation 3213 (kafka.coordinator.group.GroupCoordinator)
19:47:51,788] INFO [GroupCoordinator 2]: Preparing to rebalance group my_group with old generation 3213 (__consumer_offsets-18) (kafka.coordinator.group.GroupCoordinator)
19:48:46,851] INFO [GroupCoordinator 2]: Stabilized group my_group generation 3214 (__consumer_offsets-18) (kafka.coordinator.group.GroupCoordinator)
19:48:46,902] INFO [GroupCoordinator 2]: Assignment received from leader for group my_group for generation 3214 (kafka.coordinator.group.GroupCoordinator)
19:50:10,846] INFO [GroupMetadataManager brokerId=2] Removed 0 expired offsets in 0 milliseconds. (kafka.coordinator.group.GroupMetadataManager)
19:54:29,365] INFO [GroupCoordinator 2]: Preparing to rebalance group my_group with old generation 3214 (__consumer_offsets-18) (kafka.coordinator.group.GroupCoordinator)
19:57:22,161] INFO [ProducerStateManager partition=unrelated_topic-329] Writing producer snapshot at offset 88002 (kafka.log.ProducerStateManager)
19:57:22,162] INFO [Log partition=unrelated_topic-329, dir=/kafkalog] Rolled new log segment at offset 88002 in 11 ms. (kafka.log.Log)
19:59:14,022] INFO [GroupCoordinator 2]: Stabilized group my_group generation 3215 (__consumer_offsets-18) (kafka.coordinator.group.GroupCoordinator)
19:59:14,061] INFO [GroupCoordinator 2]: Assignment received from leader for group my_group for generation 3215 (kafka.coordinator.group.GroupCoordinator)
20:00:10,846] INFO [GroupMetadataManager brokerId=2] Removed 0 expired offsets in 0 milliseconds. (kafka.coordinator.group.GroupMetadataManager)
20:02:57,821] INFO [GroupCoordinator 2]: Preparing to rebalance group my_group with old generation 3215 (__consumer_offsets-18) (kafka.coordinator.group.GroupCoordinator)
20:06:51,360] INFO [GroupCoordinator 2]: Stabilized group my_group generation 3216 (__consumer_offsets-18) (kafka.coordinator.group.GroupCoordinator)
20:06:51,391] INFO [GroupCoordinator 2]: Assignment received from leader for group my_group for generation 3216 (kafka.coordinator.group.GroupCoordinator)
20:10:10,846] INFO [GroupMetadataManager brokerId=2] Removed 0 expired offsets in 0 milliseconds. (kafka.coordinator.group.GroupMetadataManager)
20:10:46,863] INFO [ReplicaFetcher replicaId=2, leaderId=8, fetcherId=0] Node 8 was unable to process the fetch request with (sessionId=928976035, epoch=216971): FETCH_SESSION_ID_NOT_FOUND. (org.apache.kafka.clients.FetchSessionHandler)
20:11:19,236] INFO [GroupCoordinator 2]: Preparing to rebalance group my_group with old generation 3216 (__consumer_offsets-18) (kafka.coordinator.group.GroupCoordinator)
20:13:54,851] INFO [ProducerStateManager partition=mytopic-321] Writing producer snapshot at offset 123640 (kafka.log.ProducerStateManager)
20:13:54,851] INFO [Log partition=mytopic-321, dir=/kafkalog] Rolled new log segment at offset 123640 in 14 ms. (kafka.log.Log)
20:14:30,686] INFO [ProducerStateManager partition=mytopic-251] Writing producer snapshot at offset 133509 (kafka.log.ProducerStateManager)
20:14:30,686] INFO [Log partition=mytopic-251, dir=/kafkalog] Rolled new log segment at offset 133509 in 1 ms. (kafka.log.Log)
20:16:01,892] INFO [GroupCoordinator 2]: Stabilized group my_group generation 3217 (__consumer_offsets-18) (kafka.coordinator.group.GroupCoordinator)
20:16:01,938] INFO [GroupCoordinator 2]: Assignment received from leader for group my_group for generation 3217 (kafka.coordinator.group.GroupCoordinator)
非常感谢您对此问题的任何帮助。
【问题讨论】:
标签: java apache-kafka kafka-consumer-api