【问题标题】:Kafka enable.auto.commit false in combination with commitSync()Kafka enable.auto.commit false 结合 commitSync()
【发布时间】:2019-01-30 04:13:29
【问题描述】:

我遇到enable.auto.commit 设置为false 的情况。对于每个poll(),获得的记录被卸载到threadPoolExecutorcommitSync() 是在上下文之外发生的。但是,我怀疑这是否是正确的处理方式,因为在我提交消息时,我的线程池可能仍在处理少量消息。

while (true) {
 ConsumerRecords < String, NormalizedSyslogMessage > records = consumer.poll(100);
 Date startTime = Calendar.getInstance().getTime();
 for (ConsumerRecord < String, NormalizedSyslogMessage > record: records) {
  NormalizedSyslogMessage normalizedMessage = record.value();
  normalizedSyslogMessageList.add(normalizedMessage);
 }
 Date endTime = Calendar.getInstance().getTime();
 long durationInMilliSec = endTime.getTime() - startTime.getTime();
 // execute process thread on message size equal to 5000 or timeout > 4000
 if (normalizedSyslogMessageList.size() == 5000) {
  CorrelationProcessThread correlationProcessThread = applicationContext
   .getBean(CorrelationProcessThread.class);
  List < NormalizedSyslogMessage > clonedNormalizedSyslogMessages = deepCopy(normalizedSyslogMessageList);
  correlationProcessThread.setNormalizedMessage(clonedNormalizedSyslogMessages);
  taskExecutor.execute(correlationProcessThread);
  normalizedSyslogMessageList.clear();
 }
 consumer.commitSync();
}

【问题讨论】:

    标签: java apache-kafka


    【解决方案1】:

    完全同意拉利特所说的。目前我正在经历同样的情况,我的处理发生在不同的线程中,消费者和生产者发生在不同的线程中。我使用 ConcurrentHashMap 在生产者和消费者线程之间共享,更新特定偏移是否已被处理。

    ConcurrentHashMap<OffsetAndMetadata, Boolean>
    

    在消费者方面,可以使用本地 LinkedHashMap 来维护从 Topic/Partition 消费记录的顺序,并在消费者线程本身中进行手动提交。

    LinkedHashMap<OffsetAndMetadata, TopicPartition>
    

    如果您的处理线程正在维护任何已消费的记录顺序,您可以参考以下链接。 Transactions in Kafka

    在我的方法中要提到一点,如果发生任何故障,数据将有可能被复制。

    【讨论】:

      【解决方案2】:

      我想这里有几件事要解决。

      首先是偏移量不同步 - 这可能是由以下任一原因引起的:

      1. 如果poll() 获取的消息数量没有达到normalizedSyslogMessageList 的大小为5000,则commitSync() 仍将运行,无论当前批次的消息是否已被处理。

      2. 1234563无论如何都要提交偏移量。

      第二部分(我相信这是您真正关心/问题) - 这是否是处理此问题的最佳方式。我会说不,因为上面的第 2 点,即correlationProcessThread 在这里以一种即发即弃的方式被调用,所以你不知道这些消息的处理何时完成以便你能够安全地提交偏移量。

      这是来自《卡夫卡权威指南》的声明 -

      重要的是要记住 commitSync() 将提交最新的 poll() 返回的偏移量,因此请确保在之后调用 commitSync() 您已完成对集合中所有记录的处理,否则您将面临风险 缺少消息。

      第 2 点尤其难以修复,因为:

      • 向池中的线程提供消费者引用基本上意味着多个线程试图访问一个消费者实例(This post 提到了这种方法和问题 - 主要是 Kafka 消费者不是线程安全的)。
      • 即使您尝试通过使用submit() 方法而不是ExecutorService 中的execute() 在提交偏移量之前获取处理线程的状态,那么您也需要进行阻塞get() 方法调用以correlationProcessThread。因此,在多个线程中处理可能不会获得很多好处。

      解决此问题的选项

      由于我不了解您的上下文和确切要求,我只能提出概念性想法,但可能值得考虑

      • 根据消费者实例需要执行的处理中断消费者实例并在同一线程中执行处理或
      • 您可以探索在地图中维护消息偏移量的可能性(在处理消息时),然后提交这些特定的偏移量 (this method)

      我希望这会有所帮助。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2017-08-30
        • 1970-01-01
        • 1970-01-01
        • 2017-12-09
        相关资源
        最近更新 更多