【发布时间】:2018-10-02 22:07:41
【问题描述】:
我正在使用spring-kafka 并使用来自kafka 主题的批处理记录并提交AbstractMessageListenerContainer.AckMode.BATCH 的偏移量。
在我的情况下,处理批处理记录需要时间(大约 20 秒),消费者线程等待批处理完成,然后再次进行轮询(在此轮询中提交偏移量)。在这种情况下,我会将List 的记录分配给一个线程(名称:ProcessThread),该线程将处理所有记录并将结果返回给消费者线程,然后消费者线程将记录结果。 (在所有这个过程中,消费者线程会一直等待,直到它从ProcessThread 获得结果,这会导致性能低下。
ProcessThread 有没有办法处理向 kafka 提交偏移量?这样消费者线程就不需要等待了,每次轮询都会创建一个新的processThread
在我的情况下,我有 20 个分区和 10 个 pod 的主题,每个 pod 有 2 个消费者线程(spring kafka 并发消费者),每个轮询 100 条记录(使用 spring boot 线程处理所有这些记录@Async)
通过上述配置,我可以在 2 小时内处理 100 万条记录,我需要将其拖到至少 40 分钟。
感谢任何帮助????
【问题讨论】:
标签: java performance spring-kafka