Jeff Bolle
02/08/2023, 2:22 PMsetProcessingGuarantees to MANUAL and then inside my function all of my processing is async, including publishing to some output topics. The final callback I have on the `CompletableFuture`from the async processing calls this method:
public <T> T handle(T obj, Throwable ex, Context context){
if (ex == null) {
context.getCurrentRecord().ack();
return obj;
} else {
context.getLogger().error("Error Processing record.", ex);
context.getCurrentRecord().fail();
return null;
}
}
Does all of that look like the right way to be acking / failing the record from inside a function using the Context? I ask because in my logs I'm seeing a lot of the following:
2023-02-08T14:13:09,961+0000 [pulsar-timer-6-1] INFO org.apache.pulsar.client.impl.UnAckedMessageTracker - [ConsumerBase{subscription='bnr-url-function-sub', consumerName='45c52', topic='<persistent://q6/default/bnr_legacy'}>] 2634 messages will be re-delivered
However I'm not seeing error logs indicating that there were an equivalent number of message failures. I'm not seeing any failures right now, and I'm seeing many copies of the UnAckedMessageTracker logging that.