【问题标题】:Design Pattern - Spring KafkaListener processing 1 million records in 1 hour设计模式 - Spring KafkaListener 在 1 小时内处理 100 万条记录
【发布时间】:2022-06-10 22:25:18
【问题描述】:

我的 Spring Boot 应用程序将每小时从 kafka 代理收听 100 万条记录。每条消息的整个处理逻辑需要 1-1.5 秒,包括数据库插入。 Broker有64个分区,这也是我@KafkaListener的并发。

在我每小时收听大约 50k 条记录的较低环境中,我当前的代码只能在一分钟内处理 90 条记录。以下是代码,所有其他配置参数,如 max.poll.records 等都是默认值:

@KafkaListener(id="xyz-listener", concurrency="64", topics="my-topic")
public void listener(String record) {

// processing logic 

}

我确实得到每小时 7-8 次“消费者很可能被踢出群组”。我认为这两个问题都可以通过隔离侦听器方法和对每条消息进行多线程处理来解决,但我不知道该怎么做。

【问题讨论】:

    标签: spring multithreading spring-kafka


    【解决方案1】:

    这里有几点需要考虑。首先,对于单个应用程序而言,64 个消费者似乎有点过多。

    考虑到默认情况下每个轮询每次都为每个消费者获取500 records,如果单个批次花费的时间超过max.poll.timeout.ms 的默认值 5 分钟,您的应用可能会过载并导致消费者被踢出组待处理。

    首先,我会考虑scaling the application horizontally,以便每个应用程序处理较少数量的分区/线程。

    增加吞吐量的第二种方法是使用batch listener,并按您在this answer 中看到的那样批量处理处理和数据库插入。

    同时使用这两种方法,您应该在每个应用程序中并行处理大量工作,并且应该能够实现您想要的吞吐量。

    当然,您应该使用不同的数据对每种方法进行负载测试,以获得适当的指标。

    编辑:针对您的评论,如果您想实现此吞吐量,我还不会放弃批处理。如果您逐行执行数据库操作,您将需要更多资源才能获得相同的性能。

    如果您的规则引擎不执行任何 I/O,您可以通过它迭代批次中的每条记录而不会损失性能。

    关于数据一致性,可以尝试一些策略。例如,您可以使用lock 来确保即使通过重新平衡,也只有一个实例将在给定时间处理给定批次的记录 - 或者在 Kafka 中使用重新平衡挂钩可能有一种更惯用的处理方式。

    有了这个,您可以在收到记录时批量加载您需要过滤掉重复/过期记录的所有信息,通过内存中的规则引擎迭代每条记录,然后批量持久化所有结果,然后释放锁。

    当然,如果不了解流程的更多细节,就很难制定出理想的策略。关键是通过这样做,您应该能够在每个实例中处理大约 10 倍以上的记录,所以我肯定会试一试。

    【讨论】:

    • 另一个注意事项:如果您使用concurrency="64",您需要确保在您要运行它的机器上拥有 64 个 CPU 内核。从技术上讲,Java 中的并发只是委托给本机操作系统线程,它不超过可用 CPU 的数量。所以,是的,如果您想要如此出色的性能,请考虑横向扩展您的应用程序。因此,不同的分区将在不同的机器上处理。
    • 当然,如果处理主要是CPU-bound,那么线程多于内核是没有意义的。但是,如果处理涉及I/O,则线程数多于内核数通常是有益的,因为在某些时候线程将被卡住等待数据,而另一个线程可以在此期间从内核中受益。 This answer 对此进行了更进一步的解释,但我未能找到一个关于该主题并不太冗长的合适资源。
    • 是的!很高兴知道。无论如何,我记得几天前,我在具有 8 个 CPU 的机器上的程序在 10 个并发线程后停止显示良好的性能。所以,是的,可能 IO 绑定假设是有道理的......
    • 是的,例如,在NIO 这样的服务器Tomcat 能够在具有更少内核的机器上处理one thread per request model 上的100 concurrent requests 之前就是这样。 Loom 在这方面会发生很大变化;实际上有一个有趣的twitter thread,Loom 的项目负责人解释了一下
    • @Tomaz 即使我使用批处理侦听器,我仍然必须通过规则引擎分别处理每条消息,并且需要单独插入以跟踪我从 Kafka 收到的重复/过时记录。我的流程不涉及 I/O,所以也许我可以尝试降低并发并尝试水平扩展。
    猜你喜欢
    • 1970-01-01
    • 2020-01-24
    • 2017-08-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-11-11
    • 2017-01-31
    • 1970-01-01
    相关资源
    最近更新 更多