【发布时间】: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