Требование - потреблять только последние сообщения от topi c при ручном перезапуске или непредвиденном сбое
Если сбой и повторный запуск задания Flink, задание начинается с восстановленной контрольной точки, и при этом выполняется попытка обработать записи из Кафка хранится в гос. Чтобы избежать старых записей, я попытался изменить идентификатор группы. Тем не менее записи с контрольной точки обрабатываются.
Я использую следующий код для обработки только самых последних записей. Оно работает. Но единственная проблема заключается в том, что я не могу игнорировать состояние из контрольной точки для Flink Kafak Consumer в случае непредвиденного сбоя.
Код: myConsumer.setStartFromLatest ();
Документация: https://ci.apache.org/projects/flink/flink-docs-stable/dev/connectors/kafka.html#kafka -consumers-start-position-configuration
Мое единственное требование - обрабатывать последние события от Kafka.
Спасибо