Put message to error topic when exceed max retries in kafka mqt (#1885)

Co-authored-by: Vishal <vishal-biyani@users.noreply.github.com>
Co-authored-by: Rahul Bhati <rjbhati009@gmail.com>
This commit is contained in:
Jacob
2021-02-02 23:20:40 +05:30
committed by GitHub
co-authored by Vishal Rahul Bhati
parent 3b4ca09931
commit ef83d4314e
+15 -14
View File
@@ -286,20 +286,6 @@ func kafkaMsgHandler(kafka *Kafka, producer sarama.SyncProducer, trigger *fv1.Me
}
}
if resp == nil {
kafka.logger.Warn("every function invocation retry failed; final retry gave empty response",
zap.String("function_url", url),
zap.String("trigger", trigger.ObjectMeta.Name))
return
}
defer resp.Body.Close()
body, err := ioutil.ReadAll(resp.Body)
kafka.logger.Debug("got response from function invocation",
zap.String("function_url", url),
zap.String("trigger", trigger.ObjectMeta.Name),
zap.String("body", string(body)))
generateErrorHeaders := func(errString string) []sarama.RecordHeader {
var errorHeaders []sarama.RecordHeader
if kafka.version.IsAtLeast(sarama.V0_11_0_0) {
@@ -314,6 +300,21 @@ func kafkaMsgHandler(kafka *Kafka, producer sarama.SyncProducer, trigger *fv1.Me
return errorHeaders
}
if resp == nil {
errorString := fmt.Sprintf("request exceed retries: %v", trigger.Spec.MaxRetries)
errorHeaders := generateErrorHeaders(errorString)
errorHandler(kafka.logger, trigger, producer, url,
fmt.Errorf(errorString), errorHeaders)
return
}
defer resp.Body.Close()
body, err := ioutil.ReadAll(resp.Body)
kafka.logger.Debug("got response from function invocation",
zap.String("function_url", url),
zap.String("trigger", trigger.ObjectMeta.Name),
zap.String("body", string(body)))
if err != nil {
errorString := string("request body error: " + string(body))
errorHeaders := generateErrorHeaders(errorString)