【问题标题】:Kafka producer callback ExceptionKafka生产者回调异常
【发布时间】:2020-08-22 13:25:56
【问题描述】:

当我们产生消息时,我们可以定义一个回调,这个回调可以期待一个异常:

kafkaProducer.send(producerRecord, new Callback() {
  public void onCompletion(RecordMetadata recordMetadata, Exception e) {
    if (e == null) {
      // OK
    } else {
      // NOT OK
    }
  }
});

考虑到生产者内置的重试逻辑,不知道开发者应该明确处理哪种异常?

【问题讨论】:

    标签: apache-kafka kafka-producer-api


    【解决方案1】:

    根据Callback Java Docs,回调期间可能发生以下异常:

    处理这条记录时抛出的异常。如果没有发生错误,则为 Null。可能抛出的异常包括:

    Non-Retriable 异常(致命,永远不会发送消息):

    • InvalidTopicException
    • OffsetMetadataTooLargeException
    • RecordBatchTooLargeException
    • RecordTooLargeException
    • UnknownServerException

    可重试异常(暂时的,可以通过增加#.retries 来覆盖):

    • CorruptRecordException
    • InchvalidMetadataException
    • NotEnoughReplicasAfterAppendException
    • NotEnoughReplicasException
    • OffsetOutOfRangeException
    • 超时异常
    • UnknownTopicOrPartitionException

    也许这是一个不令人满意的答案,但最终哪些异常以及如何处理它们完全取决于您的用例和业务需求。

    处理生产者重试

    但是,作为开发人员,您还需要处理 Kafka Producer 的重试机制本身。重试主要由以下因素驱动:

    retries:设置大于零的值将导致客户端重新发送任何发送失败并可能出现暂时性错误的记录。请注意,此重试与客户端在收到错误后重新发送记录没有什么不同。 允许重试而不将 max.in.flight.requests.per.connection(默认值:5)设置为 1 可能会改变记录的顺序,因为如果将两个批次发送到单个分区,并且第一个失败并重试但第二次成功,则第二批中的记录可能首先出现。另外请注意,如果 delivery.timeout.ms 配置的超时在成功确认之前首先到期,则在重试次数用完之前生产请求将失败。 用户通常应该不设置此配置,而是使用 delivery.timeout.ms 来控制重试行为。

    retry.backoff.ms:尝试重试对给定主题分区的失败请求之前等待的时间。这样可以避免在某些失败场景下,在紧密的循环中重复发送请求。

    request.timeout.ms:配置控制客户端等待请求响应的最长时间。如果在超时之前没有收到响应,客户端将在必要时重新发送请求,或者如果重试次数用尽,则请求失败。这应该大于replica.lag.time.max.ms(代理配置),以减少由于不必要的生产者重试而导致消息重复的可能性。

    建议保留上述三个配置的默认值,而是关注由定义的硬性上限时间限制

    delivery.timeout.ms:调用 send() 返回后报告成功或失败的时间上限。这限制了记录在发送之前将被延迟的总时间、等待代理确认的时间(如果预期)以及可重试发送失败所允许的时间。如果遇到不可恢复的错误,重试已用尽,或者记录被添加到到达较早交付到期期限的批次中,生产者可能会报告未能在此配置之前发送记录。此配置的值应大于或等于request.timeout.mslinger.ms 之和。

    【讨论】:

      【解决方案2】:

      您可能会收到BufferExhaustedExceptionTimeoutException

      制作人制作一张唱片后,只需将您的 Kafka 放下即可。然后继续制作唱片。一段时间后,您应该会在回调中看到异常。

      这是因为,当您发送第一条记录时,会获取元数据,之后,记录将被批处理和缓冲,并且它们最终会在您可能会看到这些异常的超时后过期。

      我想超时是delivery.timeout.ms,当它过期时会给你一个TimeoutException异常。

      【讨论】:

        【解决方案3】:

        试图在@Mike 的回答中添加更多信息,我认为回调接口中只有少数例外是枚举。

        在这里你可以看到整个列表:kafka.common.errors

        在这里,您可以看到哪些是可重试的,哪些是可重试的 不是:kafka protocol guide

        代码可能是这样的:

        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") //it's non-retriable
                    writeDiscardRecordsToSomewhereElse
                 } else {
                    log.warn("It's no that cold") //it's retriable
                 }
              } else {
                log.debug("It's summer. Everything is fine")
              }
            }
        }
        

        【讨论】:

          猜你喜欢
          • 2023-03-03
          • 2017-05-27
          • 2014-10-29
          • 1970-01-01
          • 2020-10-21
          • 2018-04-01
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          相关资源
          最近更新 更多