【问题标题】:Usage of Java Kafka Consumer in multiple threadsJava Kafka Consumer在多线程中的使用
【发布时间】:2020-09-01 09:52:25
【问题描述】:

我正在考虑在线程池中使用 Kafka Consumer。我提出了这种方法。现在它似乎工作正常,但我正在考虑缺点以及这种方法会带来什么问题。基本上我需要的是将记录处理与消费分离。此外,我需要有一个强有力的保证,即只有在处理完所有记录后才会提交。有人可以就如何更好地做到这一点提出建议或建议吗?

  final var consumer = new KafkaConsumer<String, String>(props);
    consumer.subscribe(topics);
    final var threadPool = Executors.newFixedThreadPool(32);

    while(true) {

        ConsumerRecords<String, String> records;

        synchronized (consumer) {
            records = consumer.poll(Duration.ofMillis(100));
        }

        CompletableFuture.runAsync(this::processTask, threadPool).thenRun(() -> {
            synchronized (consumer) {
                consumer.commitSync();
            }
        });
    }

【问题讨论】:

  • 这里可以写完整的代码吗?我假设会有一个 while 循环轮询 kafka 消费者以连续获取您希望以异步方式处理的记录。
  • 是的,我更新了问题。

标签: java apache-kafka kafka-consumer-api


【解决方案1】:

问题

此解决方案不适合所述要求:

另外,我需要有一个强有力的保证,即只有在处理完所有记录后才会提交

场景:

  1. 轮询读取 100 条记录,开始异步处理
  2. 轮询读取 5 条记录,开始异步处理
  3. 立即处理 5 条记录,并在处理 100 条记录时完成消费者提交
  4. 消费者崩溃

当消费者再次启动时,最后一次提交将对应第 105 条记录。因此它将开始处理第 106 条记录,我们错过了成功处理 1-100 条记录。

您只需通过以下方式提交您在该轮询中处理的偏移量:

void commitSync(Map<TopicPartition, OffsetAndMetadata> offsets);

此外,还需要保证顺序,以便首先提交第一个轮询,然后是第二个,依此类推。这会相当复杂。

主张

我相信您正在尝试在消息处理中实现并发。这可以通过更简单的解决方案来实现。增加您的 ma​​x.poll.records 以读取一个体面的批次,将其分成更小的批次并异步运行以实现并发。完成所有批次后,提交给 kafka 消费者。

【讨论】:

  • 是的,我就是害怕这个。文档中很少提到 KafkaConsumer 对于多线程是不安全的,但我不确定它在同步条件下的行为。非常感谢。
【解决方案2】:

我遇到了以下文章,它解耦了 kafka 中记录的消费和处理。您可以通过显式调用poll() 方法并在pause()resume() 方法的帮助下处理记录来实现此目的。

Processing kafka records in Multi-threaded env

【讨论】:

    猜你喜欢
    • 2017-04-20
    • 2014-04-12
    • 2020-05-19
    • 2016-10-12
    • 2019-07-23
    • 1970-01-01
    • 2023-03-02
    • 1970-01-01
    • 2018-10-28
    相关资源
    最近更新 更多