Flink

Senior

Как развернуть кластер Flink на Kubernetes через нативную интеграцию?

Flink не просто работает внутри Kubernetes — с нативной интеграцией он напрямую обращается к Kubernetes API, чтобы управлять собственными подами TaskManager'ов. Это значимо другой деплой, чем «Flink в контейнере».

Читать на:
Нативная интеграция Standalone Deployment Kubernetes Operator
Примерно как эти три подхода сравниваются по тому, насколько сам Flink осведомлён о Kubernetes — не измеренная статистика.

Есть два реальных способа запускать Flink на Kubernetes: standalone-деплой (Flink понятия не имеет, что он на K8s — вы управляете репликами как любым другим stateless-деплойментом) и нативная интеграция (собственный ResourceManager Flink обращается к Kubernetes API, чтобы запрашивать и освобождать поды TaskManager'ов по мере надобности задач).

Что на самом деле значит «нативная»

При standalone-деплое TaskManager'ы — просто поды в Kubernetes Deployment: Kubernetes перезапускает их при гибели, но сам Flink не может попросить больше или меньше подов исходя из потребности задачи. При нативной интеграции ResourceManager Flink напрямую вызывает Kubernetes API (создание/удаление подов), чтобы масштабировать TaskManager'ы под реальную потребность задачи — ближе по духу к тому, как всегда работала интеграция Flink с YARN, просто с другим планировщиком.

# запускает кластер, который сам управляет подами TaskManager ./bin/flink run-application \ --target kubernetes-application \ -Dkubernetes.cluster-id=my-flink-app \ -Dkubernetes.container.image=my-flink:1.20 \ local:///opt/flink/usrlib/job.jar

Какой RBAC это реально требует

Поскольку Flink сам создаёт и удаляет поды, его service account нуждается в реальных правах в namespace, где он работает — не просто доступе только на чтение, который был бы нужен обычному приложению:

# минимально жизнеспособный RBAC для нативной интеграции, один namespace rules: - apiGroups: [""] resources: ["pods", "services", "configmaps"] verbs: ["get", "list", "watch", "create", "update", "patch", "delete"]

Это заметно бо́льшая поверхность прав, чем нужна большинству прикладных нагрузок, и стоит заранее обозначить тому, кто владеет политикой безопасности кластера, прежде чем это всплывёт сюрпризом на ревью — это неотъемлемая часть того, как работает нативная интеграция, а не ошибка конфигурации.

Где здесь всё ещё остаётся выбор

Flink Kubernetes Operator (отдельный, более новый проект) работает поверх любого из режимов деплоя и добавляет управление жизненным циклом в стиле Kubernetes — kubectl apply на custom resource FlinkDeployment вместо ручного запуска flink run-application, с оператором, берущим на себя апгрейды, редеплои на основе savepoint и reconciliation. Это не столько третий режим деплоя, сколько слой управления — стоит внедрять, когда у вас больше пары задач Flink в эксплуатации, и не так очевидно того стоит как лишняя движущаяся часть для одной задачи.

Вне области opti-pipe: движок правил читает метрики, которые экспортирует работающая или завершённая задача Flink — ему всё равно, были ли TaskManager'ы этой задачи выделены нативной интеграцией, standalone Deployment'ом или Kubernetes Operator'ом, поскольку ничто из этого не меняет смысл метрик.

Посмотрите, как это выглядит на вашем собственном пайплайне.

Загрузите реальный event log Spark, run_results.json от dbt или экспорт метрик Flink — и получите конкретные рекомендации, которые нужно одобрить, а не ещё одно эмпирическое правило.