Apache Flink 1.15 KafkaSink error handling

Viewed 64

We are using the KafkaSink at the end of our dispatcher pipeline. I was wondering if there is a way to handle errors getting Metadata for topics, which do not exist on our broker. Normally the topics should be present, but we may have some race conditions, where data is processed and should be written to topics, which were deleted in the meantime.

So in this case, the KafkaProducer tries to get Metadata for this topic with a warning UNKNOWN_TOPIC_OR_PARTITION and fails after a configured timeout with

Topic <topicname> not present in metadata after 60000 ms.

and the complete job fails. Autocreation of topics is no option, because we do not want orphaned topics to be recreated.

I could not find a way to configure or extend the KafkaSink accordingly to ignore such cases and continue processing other events.

Is there something I overlooked or maybe another way to handle such cases?

0 Answers
Related