How to send avro message to new Topic using functional Kstream (processor)

Viewed 55

I am getting "could not be established. Broker may not be available" when the message sent to new topic "failureTopic" or any exception. I am using Kafka version: 3.0.0

@Bean
@SuppressWarnings("unchecked")
public Function<KStream<String, AvroClass>, KStream<String, AvroClass>> process() {

final AvroClass[] finalMessage = {null};
return input -> input.branch( (k, avroMessage) -> {
try {
finalMessage[0] = subprocess();
if (finalMessage[0] != null)
     return true;
else {
   Message<NewAvroClass> mess = MessageBuilder.withPayload(avroMessage).build();
   streamBridge.send("failureTopic", mess);
   return false;
}
} catch (Exception e) {
    handleProcessingException(k, avroMessage);
    return false;
}
},
(k, v) -> true)[0].map((k, v) -> new KeyValue<>(k, finalMessage[0] ) );
}
0 Answers
Related