Облачная платформаAdvanced

Как предотвратить потерю данных после перезапуска Flink джобы?

Язык статьи: Русский
Показать оригинал
Страница переведена автоматически и может содержать неточности. Рекомендуем сверяться с английской версией.

Механизм чекпоинтов/сейвпоинтов DLI Flink является полным и надёжным. Вы можете использовать этот механизм для предотвращения потери данных при ручном перезапуске джобы или её перезапуске из‑за исключения.

  • Чтобы предотвратить потерю данных, вызванную перезапуском джобы из‑за системных сбоев, выполните следующие действия:
    • Для Flink SQL джоб выберите Enable Checkpointing и задайте подходящий интервал чекпоинтов, учитывающий влияние на производительность сервиса и длительность восстановления после исключения. Выберите Auto Restart upon Exception и Restore Job from Checkpoint. После настройки, если джоба будет перезапущена аномально, внутреннее состояние и позиция потребления будут восстановлены из последнего файла чекпоинта, что гарантирует отсутствие потери данных и точную и согласованную семантику внутреннего состояния, например операторов агрегации. Кроме того, чтобы обеспечить отсутствие дублирования данных, используйте базу данных или файловую систему с первичным ключом в качестве источника данных. В противном случае добавьте логику дедупликации (данные, сгенерированные от последнего успешного чекпоинта до момента возникновения исключения, будут потребляться повторно) для downstream‑процессов.
    • Для Flink Jar джоб необходимо включить чекпоинтинг в коде. Кроме того, если у вас есть пользовательское состояние для сохранения, нужно реализовать интерфейс ListCheckpointed и задать уникальный ID для каждого оператора. В конфигурации джобы выберите Restore Job from Checkpoint и укажите путь к чекпоинту.
      Note

      Чекпоинтинг Flink гарантирует точность и согласованность данных внутреннего состояния. Однако для пользовательских Source/Sink или операторов с состоянием необходимо реализовать API ListCheckpointed, чтобы обеспечить надёжность данных сервиса.

  • Чтобы предотвратить потерю данных после ручного перезапуска джобы из‑за изменения сервиса, выполните следующие действия:
    • Для джоб без внутреннего состояния можно задать время начала или позицию потребления источника данных Kafka на момент, предшествующий остановке джобы.
    • Для джоб с внутренним состоянием можно выбрать Trigger Savepoint при остановке джобы. Включите Restore Savepoint при повторном запуске джобы. Джоба восстановит позицию потребления и состояние из выбранного файла сейвпоинта. Механизм генерации и формат чекпоинтов Flink совпадают с таковыми у сейвпоинтов. Перейдите к списку Flink джоб и выберите More > Import Savepoint в колонке Operation Flink джобы, чтобы импортировать последний чекпоинт в OBS и восстановить джобу из него.