【问题标题】:Kafka async Commit Offset ReplicationKafka异步提交偏移复制
【发布时间】:2017-06-07 04:29:39
【问题描述】:

我们偶尔会遇到副本领导者和其他 ISR 节点之间的高延迟,这会导致消费者收到以下错误:

org.apache.kafka.clients.consumer.RetriableCommitFailedException: Commit offsets failed with retriable exception. You should retry committing offsets.
Caused by: org.apache.kafka.common.errors.TimeoutException: The request timed out.

我可以增加offsets.commit.timeout.ms,但我不想这样做,因为它可能会导致额外的副作用。 但从更广泛的角度来看,我不希望代理等待同步所有其他副本上的提交偏移量,而是在本地提交并异步更新其余部分。 仔细检查代理配置,我发现 offsets.commit.required.acks 看起来就是这样配置的,但文档也隐晦地指出:the default (-1) should not be overridden

为什么?我什至尝试过查看代理源代码,但几乎没有找到其他信息。

知道为什么不推荐这样做吗?有没有不同的方法可以达到相同的结果?

【问题讨论】:

  • 如果这个问题是暂时的,你可能会增加'request.timeout.ms'。
  • 我们在 Kafka 2.1.1 上遇到了完全相同的问题...解决此问题是否成功?

标签: apache-kafka kafka-consumer-api


【解决方案1】:

我建议实际重试提交偏移量。

让您的消费者异步提交偏移量并实现重试机制。但是,重试异步提交可能会导致您在提交较大的偏移量后提交较小的偏移量,这应该通过各种方式避免。

在《Kafka - The Definitive Guide》一书中,有一个关于如何缓解这个问题的提示:

重试异步提交:为异步重试获取正确的提交顺序的一种简单模式是使用单调递增的序列号。每次提交时增加序列号,并将提交时的序列号添加到 commitAsync 回调中。当您准备发送重试时,检查回调获得的提交序列号是否等于实例变量;如果是,则没有更新的提交,可以安全地重试。如果实例序列号更高,请不要重试,因为已经发送了更新的提交。

作为一个例子,你可以在下面的 Scala 中看到这个想法的实现:

import java.util._
import java.time.Duration
import org.apache.kafka.clients.consumer.{ConsumerConfig, ConsumerRecord, KafkaConsumer, OffsetAndMetadata, OffsetCommitCallback}
import org.apache.kafka.common.{KafkaException, TopicPartition}
import collection.JavaConverters._

object AsyncCommitWithCallback extends App {

  // define topic
  val topic = "myOutputTopic"

  // set properties
  val props = new Properties()
  props.put(ConsumerConfig.GROUP_ID_CONFIG, "AsyncCommitter5")
  props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092")
  // [set more properties...]
  

  // create KafkaConsumer and subscribe
  val consumer = new KafkaConsumer[String, String](props)
  consumer.subscribe(List(topic).asJavaCollection)

  // initialize global counter
  val atomicLong = new AtomicLong(0)

  // consume message
  try {
    while(true) {
      val records = consumer.poll(Duration.ofMillis(1)).asScala

      if(records.nonEmpty) {
        for (data <- records) {
          // do something with the records
        }
        consumer.commitAsync(new KeepOrderAsyncCommit)
      }

    }
  } catch {
    case ex: KafkaException => ex.printStackTrace()
  } finally {
    consumer.commitSync()
    consumer.close()
  }


  class KeepOrderAsyncCommit extends OffsetCommitCallback {
    // keeping position of this callback instance
    val position = atomicLong.incrementAndGet()

    override def onComplete(offsets: util.Map[TopicPartition, OffsetAndMetadata], exception: Exception): Unit = {
      // retrying only if no other commit incremented the global counter
      if(exception != null){
        if(position == atomicLong.get) {
          consumer.commitAsync(this)
        }
      }
    }
  }

}

【讨论】:

    猜你喜欢
    • 2022-10-04
    • 2017-12-10
    • 2017-08-22
    • 2018-03-27
    • 2020-10-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-11-14
    相关资源
    最近更新 更多