Модуль 30. Jobs, CronJobs и обработка данных
Batch-задачи в k8s — это «одноразовые» или периодические поды, которые что-то посчитали — завершились. Они отлично ложатся на ETL/ELT (экстракт/трансформ), генерацию отчётов, загрузки в витрины, компакты в lakehouse, проверки качества (DQ).
Инструменты слоя:
- Job/CronJob — базовые примитивы batch.
- Airflow/Dagster — оркестрация пайплайнов (KubernetesPodOperator, K8sExecutor).
- Spark на k8s (или Spark Operator) — тяжёлые распределённые трансформации.
- KEDA/ScaledJob — событийное масштабирование по очередям/метрикам.
- Volcano/YuniKorn (по ситуации) — специализированный планировщик под батч.
Когда Job/CronJob, а когда «что-то ещё»
Job/CronJob — хороши когда:
- Задача самодостаточна (один контейнер/пара контейнеров), вход/выход ясен (файлы/S3/БД), длительность от секунд до часов.
- Нужен жёсткий контроль повторов/таймаутов и простая конкарренси (параллельные воркеры).
- Сохранение результатов — внешнее (БД/S3); сам под — «одноразовый».
Airflow — когда:
- Много шагов/зависимостей, нужны SLA, ретраи с backoff, окна расписаний, уведомления и каталоги метаданных.
- Хотите единый «контрольный центр» за пайплайнами BI/DWH.
Spark (или оператор) — когда:
- Данные «большие/распределённые», важно распараллеливание и shuffle, нужны коннекторы к lakehouse, стриминг.
Анатомия Job/CronJob (что действительно важно)
Job (batch/v1) — ключевые поля:
- parallelism — сколько подов одновременно.
- completions — сколько «успешных подов» нужно (mode Indexed подходит для «массивов» задач).
- backoffLimit — число ретраев.
- activeDeadlineSeconds — общий таймаут job.
- ttlSecondsAfterFinished — авто-уборка завершённых job/подов.
- podFailurePolicy — условный retry/не retry по exit-кодам/причинам (чтобы не «молотить» бессмысленно).
- template.spec.restartPolicy = OnFailure / Never.
- suspend — «поставить на паузу».
CronJob (batch/v1) — ключевые поля:
- schedule — cron-маска (ориентируйтесь на время контроллера; сложные TZ — через оркестратор).
- concurrencyPolicy — Forbid (новый не стартует, пока идет предыдущий), Replace (убить старый, запустить новый), Allow.
- startingDeadlineSeconds — сколько ждать «пропущенный» запуск.
- successfulJobsHistoryLimit / failedJobsHistoryLimit — сколько историй хранить, остальное удалять.
Мини-эскиз (идея):
apiVersion: batch/v1
kind: CronJob
spec:
schedule: "5 * * * *" # каждый час в :05
concurrencyPolicy: Forbid # не копим параллельные запуски
startingDeadlineSeconds: 600 # если опоздали <10 мин — запустим
successfulJobsHistoryLimit: 1
failedJobsHistoryLimit: 2
jobTemplate:
spec:
backoffLimit: 3
activeDeadlineSeconds: 3600
ttlSecondsAfterFinished: 600
template:
spec:
restartPolicy: OnFailure
serviceAccountName: etl-sa
Практика: сразу включайте TTL Controller и хистори-лимиты — это главный антидот от «кладбища job».
Паттерны batch-обработки
A. Parameter sweep / «массив задач»
Используйте Indexed Jobs: completionMode: Indexed, поды получают индекс (JOB_COMPLETION_INDEX) и берут свой сегмент данных.
B. Work-queue (очередь задач)
Deployment «пишет» в очередь (Kafka, RabbitMQ, SQS), а ScaledJob (KEDA) поднимает Job-воркеры по глубине очереди. Плюсы: эластичность, автоскейл «по событию».
C. «Тяжёлые» ETL
Разбивайте на шаги: extract (Job/Pod), transform (SparkApplication), load (Job). Стейт — только вовне (S3/БД). Все шаги идемпотентны.
D. DQ/GE проверки
Запуск Great Expectations/dbt tests как Job/CronJob; результаты — в BI/каталог.
Интеграция с Airflow
Подходы:
- KubernetesPodOperator — каждый таск = изолированный pod (образ со всем внутри). Простая безопасность/изолирование, чёткая трассировка.
- KubernetesExecutor — сам Airflow шедулит таски в k8s. Хорошо, когда почти всё в кластере.
- CeleryExecutor + k8s — если история/опыт с Celery, но окружение — k8s.
Практические нюансы:
- Образы не «жирные», зависимости версионируем, артефакты — в S3.
- serviceAccount с минимальным RBAC (читать секреты/ConfigMap, писать в нужные бакеты/БД — и точка).
- Секреты — через External Secrets/Vault, а не в env-файлах.
- SLA: Airflow «видит» фактические статусы pod, логирует stdout/stderr в централизованный лог.
Spark на Kubernetes (и оператор)
Варианты:
- Native spark-submit: драйвер и экзекьюторы — поды; настройки через --conf spark.kubernetes.*.
- Spark Operator (CRD SparkApplication): декларативно, с retry-логикой и «контроллером жизни» джобы (удобно для GitOps).
Практические акценты:
- Dynamic Allocation: в k8s используйте shuffle tracking (без внешнего shuffle-сервиса).
- NodePools: отдельный пул под executors (taints), чтобы не «пинать» прод-микросервисы.
- Хранилище: чтение/запись в S3/объектку (Parquet/Iceberg), локальные PV — только для спиллов и очень аккуратно.
- Разграничение: PriorityClass, ResourceQuota/LimitRange на ns, чтобы джобы не «выели» кластер.
Оптимизация под ресурсы (чтобы «летало», а не душило)
-
Requests/Limits:
- для JVM-тасков избегайте агрессивных CPU limits (throttling убивает throughput и GC-паузы); часто лучше задать requests≈limits или вовсе без limit на CPU в worker-нодах для batch-пула; память — всегда с запасом.
- -XX:MaxRAMPercentage / Spark memoryOverhead согласовать с k8s-limit, иначе OOMKilled.
- QoS-класс: критичные ночные джобы — Guaranteed (requests=limits), фоновые — Burstable/BestEffort.
- Affinity/Taints: выделенный node pool для batch (taints batch=true:NoSchedule + tolerations в job).
- Horizontal параллелизм: в Job — parallelism; в Spark — --executor-cores/--num-executors/dynamic allocation.
- Warm-образы: используйте «base-образ + слой с кодом» и локальные кэши (registry-mirror), чтобы старт подов был быстрым.
Надёжность: ретраи, дедлайны, идемпотентность
- Retry-логика: backoffLimit + podFailurePolicy (например, не ретраить ExitCode=2 для «ошибка входных данных»).
- Таймауты: activeDeadlineSeconds (цепочка «залипла — убили — перезапустили в следующий слот»).
- ConcurrencyPolicy Forbid для CronJob — чтобы не множились параллельные запуски одного окна.
- Idempotency: шаг можно повторить без дублей (меркеры в S3/БД, UPSERT/мерджи, транзакции).
- Run-to-completion: под должен завершиться сам (не держим «демоны» в Job).
Мониторинг и алертинг batch
Собираем:
- kube-метрики: kube_job_status_failed/succeeded/active, kube_job_complete, cronjob_next_schedule_time, длительность kube_pod_container_*.
- Airflow/Spark: «таски фейлятся», длина очереди, длительность этапов.
- SLI/SLO: «% успешных job за окно», «p95 длительности задач», «задержка начала относительно расписания».
Алерты (идеи):
- CronJobMissed (нет запусков > 2 интервала),
- JobBackoffStorm (высокий rate ретраев),
- TooManyActiveJobs в ns/team,
- SparkAppFailed / AirflowDAGFailed > N%,
- QueueDepthHigh (для KEDA).
Риски и анти-паттерны
|
Риск |
Проявление |
Контрмера |
|---|---|---|
|
Накопление невыполненных job |
«Кладбище» в ns, тормоза контроллера |
ttlSecondsAfterFinished, historyLimit, ConcurrencyPolicy=Forbid/Replace, дедлайны, алерты |
|
«Вечные ретраи» при плохих данных |
Бесполезная утилизация кластера |
podFailurePolicy: не ретраить по известным exit-кодам; валидация входа |
|
CPU throttling (JVM/Spark) |
Долгая работа, GC-пилы |
Ослабить/убрать CPU-limits в batch-пуле; тюнинг JVM/Spark memoryOverhead |
|
Монолитные DAG’и |
Починка = «перекатывать всё» |
Декомпозиция на независимые шаги/артефакты; идемпотентность |
|
Секреты в env/образе |
Утечки |
ESO/Vault, монтирование через CSI; RBAC минимум |
|
Смешение batch и прод-сервисов |
Прод «подвис» из-за ночных ETL |
Отдельные нод-пулы с taints, PriorityClass, Quota |
|
Отсутствие SLO |
«Нам кажется, что тормозит» |
Завести SLO (успешность/длительность/отложенные старты) и алерты burn-rate |
Практика (лабораторка за 1–2 дня)
- CronJob отчёта: каждые 30 мин запускает dbt-модель → пишет витрину в Postgres/ClickHouse. Включить: Forbid, startingDeadlineSeconds, ttlSecondsAfterFinished, historyLimit.
- KEDA ScaledJob: масштабирование воркеров по глубине Kafka/RabbitMQ-очереди; ограничить maxReplicaCount.
- Airflow: DAG из трёх задач (extract→transform→load) на KubernetesPodOperator, секреты — через ESO.
- Spark Operator: SparkApplication на чтение Parquet из S3, dynamic allocation, ограниченный PriorityClass.
- Мониторинг: дашборд SLO batch + алерты «missed schedule», «job backoff storm», «queue depth».
- Анти-тест: сломать входные данные → показать, что job не ретраит бесконечно (podFailurePolicy), и уведомить on-call.
Чек-лист «готово к продакшену (batch)»
- Везде включён TTL Controller + ttlSecondsAfterFinished; у CronJob — historyLimit.
- concurrencyPolicy задана; дедлайны/ретраи продуманы; podFailurePolicy настроен.
- Отдельный нод-пул/taints для batch; PriorityClass/Quota выставлены.
- Секреты — через Vault/ESO/CSI; SA/RBAC — минимальные.
- Airflow/Spark интегрированы декларативно (GitOps); образы версионируем.
- Дашборды/алерты по SLO batch; очереди/ретраи/длительности под контролем.
- Док-пакет: runbook’и «CronJob пропустил окно», «Job в backoff», «Spark не получил executors», «Очередь растёт».
Вопрос-ответ
В: Что выбрать: CronJob или Airflow по расписанию?
О: Если задача атомарная и не требует управлять зависимостями — CronJob. Если есть цепочки, SLA, зависимости, уведомления — Airflow.
В: Как не допустить лавины параллельных запусков?
О: У CronJob поставить concurrencyPolicy: Forbid (или Replace), startingDeadlineSeconds, и таймаут activeDeadlineSeconds в job.
В: Как чистить «хвосты» из завершённых job?
О: Включить TTL (ttlSecondsAfterFinished) и successful/failedJobsHistoryLimit у CronJob. Плюс периодический kubectl/оператор-уборщик — как страховка.
В: Когда нужен KEDA/ScaledJob?
О: Когда объём задач не по времени, а по событию (длина очереди/лаг). KEDA автоматически создаёт нужное число job-воркеров.
В: Как обезопасить прод-сервисы от «ночных» ETL?
О: Выделенный node pool с taints для batch, PriorityClass ниже, ResourceQuota на ns; Ingress/DB лимиты соединений.
В: Почему мои Java/Scala job «тормозят» в k8s сильнее, чем на VM?
О: Чаще всего — CPU throttling из-за жёстких limits и несогласованных настроек JVM. Ослабьте CPU-limit в batch-пуле, настройте MaxRAMPercentage/memoryOverhead.
В: Как делать «ровно один» запуск, даже если было несколько триггеров?
О: Дедупликация: внешний лейк-маркер (S3/БД), «upsert-семантика» загрузок, блокировка (advisory lock/Redis) — и идемпотентный код.
Batch-уровень в Kubernetes прост, если помнить три вещи: дисциплина расписаний/ретраев/таймаутов, изоляция и экономия ресурсов (node-пулы, QoS, KEDA) и идемпотентность шагов. С Airflow и Spark вы получите «скелет» для сложных конвейеров, а с правильными настройками Job/CronJob — ни один ночной прогон не превратит кластер в «кладбище задач».



