Capstone overview — production-grade Flink platform
В первых 20 модулях курса мы разобрали internals Flink: архитектуру JobManager/TaskManager, network stack, state backends, checkpoint механику, exactly-once 2PC, disaggregated state, SQL planner, CEP, AI-фичи, performance profiling. Capstone объединяет это в production-grade платформу, на которой несколько команд могут запускать свои Flink-jobs независимо.
Capstone: MySQL CDC -> Flink enrichment -> Paimon + Kafka CI/CD для Flink: savepoint-driven deploymentВ этом модуле мы построим end-to-end Flink-платформу с пятью ключевыми возможностями: multi-tenant K8s setup, observability stack, autoscaling, savepoint automation, и runbook’и для типовых инцидентов. Финальный урок — миграция legacy Flink 1.20 джоба на 2.2 с переходом на disaggregated state.
Архитектура целевой платформы
Платформа состоит из четырёх слоёв:
Главный принцип: self-service для tenants, automation для platform team. Tenant команда не должна писать PagerDuty-pages в Slack platform team каждый раз, когда нужен новый Flink job или savepoint. Платформа должна предоставлять Kubernetes-native интерфейс (FlinkDeployment CRD), на котором tenants работают самостоятельно.
Требования и trade-offs
Multi-tenancy. Несколько команд, разные SLA, разные требования к ресурсам. Tenant A — fraud detection с строгим RTO < 5 min. Tenant B — ads attribution с bursty workload (10x пиков на Black Friday). Tenant C — analytics с relaxed requirements. Платформа должна обслужить всех без cross-tenant interference.
Observability. Каждый Flink job должен экспонировать standard metrics (records-in-per-second, checkpoint duration, backpressure), которые автоматически попадают в Prometheus и появляются в Grafana dashboards. On-call видит проблемы до того, как пользователи начинают звонить.
Autoscaling. Static parallelism не работает для bursty workload. Платформа должна автоматически scale up при Kafka lag и scale down при quiet period. Идеально — predict bursts перед ними (ML-based prediction).
Savepoint automation. Каждые N минут — periodic savepoint в S3 с retention policy. Cross-region replication для DR. Automated rollback playbook: «вернуть job на savepoint от 2 часов назад» должно занимать секунды, не часы.
Runbook automation. Типовые инциденты (job crashed, checkpoint failed, RocksDB compaction stuck) должны иметь автоматизированные runbook’и с четкими шагами для on-call: «запусти этот скрипт, и job восстановится».
Trade-offs:
- Operator vs Session cluster. Operator с per-job clusters даёт лучшую изоляцию (один job — один K8s pod set), но больший overhead. Session cluster с несколькими jobs дешевле, но один сбой влияет на все jobs. Мы выбираем Operator.
- Open-source vs managed. AWS Managed Flink (Kinesis Data Analytics) проще, но vendor lock-in и хуже для multi-tenancy. Self-hosted on K8s сложнее, но полная гибкость. Мы делаем self-hosted.
- Reactive mode vs explicit autoscaling. Reactive mode (Flink 1.13+) автоматически использует все доступные resources. Explicit autoscaling через Operator API более controllable. Мы используем combination — reactive для compute, custom controller для memory.
Tools stack
Kubernetes. Базовая платформа. На любом managed K8s (EKS, GKE, AKS) или on-prem. Минимум 1.25+ для Flink Operator 1.10+.
Flink Kubernetes Operator 2.x. Cluster-scoped operator, управляющий FlinkDeployment CRD. Поддерживает Application Mode, Session Mode, periodic savepoints, autoscaling.
Apache Flink 2.2. Latest LTS на момент написания. ForstDB для disaggregated state, AI features (CREATE MODEL, ML_PREDICT, VECTOR_SEARCH), Adaptive Scheduler default.
Prometheus + Grafana. Stack для metrics. Flink exposes metrics через PrometheusReporter, Prometheus scrapes из pods (через ServiceMonitor + Prometheus Operator).
Loki + Promtail. Centralized logging. Flink job логи попадают в Loki через Promtail sidecar в TaskManager pods. Запросы через Grafana logs panel.
Vault. Secrets management. Kafka credentials, S3 keys, DB passwords — всё в Vault. K8s ServiceAccount-based auth через Kubernetes Vault Auth method.
ArgoCD. GitOps deployment. Все FlinkDeployments в Git, ArgoCD синхронизирует state cluster со state в Git. Tenant команды делают PR в общий platform repo, ArgoCD applies.
PagerDuty + Slack. Alerting. Critical алерты — PagerDuty (wake on-call), warning — Slack channel. Runbook URLs в alert annotations.
Не существует «правильного» tools stack — выбирайте на основе вашей organization. Если у вас уже Datadog для metrics, не вводите Prometheus только потому что «так делают другие». Принцип: используйте standard в company, добавляйте только то, что необходимо для Flink-specific нужд.
Урок-по-урок plan
Следующие уроки модуля построят все компоненты платформы:
Урок 2: Multi-tenant K8s setup. Namespace-per-team, RBAC, ResourceQuota, NetworkPolicy. Flink Operator config для cluster-wide management. Tenant onboarding процедура.
Урок 3: Autoscaling setup. Flink Autoscaler (FLIP-271) для compute. Custom controller для memory-based scaling. ML-based prediction для bursty workloads через ARIMA на исторических Kafka lag метриках.
Урок 4: Savepoint automation. Periodic savepoint через CronJob, retention policy через Lambda (удаляем savepoints старше 30 дней). Cross-region copy в DR S3. Rollback playbook script.
Урок 5: Migration 1.20 -> 2.2 + disaggregated state. Real-world migration scenario: legacy Flink 1.20 job, 2TB state, требует переход на 2.2 с disaggregated state (ForstDB) для лучшего checkpoint performance. Шаги, риски, validation.
Что не входит в capstone
Несколько важных тем, которые мы не покрываем здесь (по причинам complexity / out-of-scope):
- CI/CD pipeline для Flink jobs. Это отдельная тема — Jenkins/GitHub Actions для build, integration tests, canary deployments. Зависит от вашей сборочной инфраструктуры.
- Schema migration. Если ваш Flink job читает Avro/Protobuf — schema evolution требует отдельной стратегии (compatibility checks, dual-read transitions). Это тема Confluent Schema Registry, а не Flink internals.
- Cost optimization beyond autoscaling. Spot instances для TaskManager, scheduled scale-down ночью, multi-cluster placement. Платформа должна позволять, но конкретные стратегии зависят от cloud provider.
- Disaster recovery testing automation. В модуле 20 разобрали DR, но автоматизированный chaos engineering для testing failover (kill primary region) — отдельная тема.
Success criteria
Платформа считается успешной, если:
- Tenant может задеплоить новый Flink job без помощи platform team — через PR с FlinkDeployment yaml в Git.
- Алерты для типовых инцидентов автоматизированы — checkpoint failure, job restart, lag spike — все имеют runbook URL в alert annotation.
- Автоматический failover при disaster — primary region упал, DR-кластер берёт нагрузку без manual intervention (или с одобрением через ChatOps).
- Multi-tenancy без cross-interference — тяжёлый tenant не «съедает» ресурсы лёгкого. ResourceQuota и priorityClass обеспечивают изоляцию.
- Cost transparency — каждый tenant видит свой cost (через label-based cost allocation) и может оптимизировать сам.
В следующем уроке начнём building с фундамента: multi-tenant K8s setup.
Building production Flink platform — это не «выходные с туториалом». Минимум 6-12 месяцев team work для зрелой платформы. Capstone модуль показывает building blocks и trade-offs, но реальная имплементация требует итераций, learning from incidents, и адаптации к specifics вашей компании.
Итоги
Capstone-платформа объединяет 5 ключевых building blocks: multi-tenant K8s, observability, autoscaling, savepoint automation, runbook’и. Каждый блок — отдельный урок с практическими конфигурациями.
Главный принцип — self-service для tenants, automation для platform team. Целевой стейт: новый Flink job деплоится за минуты, инциденты разруливаются автоматически или через runbook за минуты, миграция между версиями Flink безопасна и предсказуема.