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] ) );
}