【问题标题】:How to submit kafka offset after processing batch records处理批处理记录后如何提交kafka offset
【发布时间】: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


    【解决方案1】:

    您应该至少升级到 1.3.7。当前版本是 2.1.10。其中没有单独的线程。只有消费者线程可以发送偏移量。消费者不是线程安全的。

    您不应该移交给单独的线程。使用更高的并发和更多的分区。

    【讨论】:

    • 你的意思是我应该将spring kafka升级到1.3.7,它对性能有何影响?你能提供一些文档吗
    • 是的,但当前是 2.1.10。
    • 您的意思是获取更多分区主题并将每个消费者分配给每个分区,对吗?这到底是什么意思There is no separate thread in those.
    • 我还有一个疑问,如果消费者长时间没有响应,比如 25 或 30 秒,kafka 会重新平衡吗?
    • 在 1.3.0 之前的版本中,每个消费者有 2 个线程。一个消费者线程和一个侦听器线程。这是为了避免由于侦听器速度慢而导致的重新平衡。 KiP-62 解决了这个问题。现在 max.poll.interval 用于 rebalance 。默认为 5 分钟。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-03-17
    • 2015-09-12
    • 1970-01-01
    • 2017-01-11
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多