【问题标题】:Distinguish how to handle exceptions in async Kafka producer区分如何处理异步 Kafka 生产者中的异常
【发布时间】:2020-10-21 02:46:20
【问题描述】:

向 Kafka 生成消息时,您可能会遇到两种错误:可重试和不可重试。处理时应该如何区分?

我想异步生成记录,将callback object 收到不可重试异常的记录保存在另一个主题(或 HBase)中,并让生产者为我处理所有收到可重试异常的记录(最多尝试次数上限)并且,当它最终到达它时,它会成为第一个)。

我的问题是:尽管有callback object,生产者仍会自行处理可检索的异常吗? 因为在Interface Callback 中说:

可重试异常(瞬态,可以通过增加# .重试)

可能是这样的代码吗?

producer.send(record, callback)

def callback: Callback = new Callback {
    override def onCompletion(recordMetadata: RecordMetadata, e: Exception): Unit = {
      if(null != e) {
         if (e == RecordTooLargeException || e == UnknownServerException || ..) {
            log.error("Winter is comming")
            writeDiscardRecordsToSomewhereElse
         } else {
            log.warn("It's no that cold") //it's retriable. The producer will keep trying by itself?
         }
      } else {
        log.debug("It's summer. Everything is fine")
      }
    }
}

Kafka 版本:0.10.0

任何灯光都将不胜感激! :)

【问题讨论】:

    标签: apache-kafka kafka-producer-api


    【解决方案1】:

    正如卡夫卡圣经(又名Kafka-The Definitive Guide)所说:

    缺点是虽然 commitSync() 会重试提交,直到它 成功或遇到不可重试的失败,commitAsync() 不会重试。

    原因:

    它不会重试是在 commitAsync() 收到一个 来自服务器的响应,可能有一个后来的提交 已经成功了。

    假设我们发送了一个提交偏移量的请求 2000. 有一个临时的通信问题,所以代理永远不会收到请求,因此永远不会响应。同时,我们处理 另一个批次并成功提交偏移量 3000。如果 commitA sync() 现在重试之前失败的提交,它可能会成功 在偏移量 3000 已经处理后提交偏移量 2000 并且 坚定的。在重新平衡的情况下,这将导致更多 重复。

    除此之外,您仍然可以创建一个递增的序列号,您可以在每次提交时增加该序列号并将该编号添加到回调对象。当重试时间到来时,只需检查 Acc 的当前值是否等于您给 Callback 的数字。如果是这样,它是安全的,您可以执行提交。否则,已经有一个新的提交,你不应该重试这个偏移量的提交。

    看起来很麻烦,那是因为如果你在考虑这个,你应该改变你的策略。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-08-22
      • 1970-01-01
      • 2021-04-30
      • 1970-01-01
      • 2020-05-19
      • 2012-01-31
      • 2023-03-03
      • 2017-05-27
      相关资源
      最近更新 更多