Потребитель Flink Kafka игнорирует новый идентификатор группы при восстановлении с контрольной точки - PullRequest
0 голосов
/ 04 апреля 2020

Требование - потреблять только последние сообщения от 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.

Спасибо

Добро пожаловать на сайт PullRequest, где вы можете задавать вопросы и получать ответы от других членов сообщества.
...