Модуль 7. Стек данных на k8s (open-source)
Типовая «data-платформа» в Kubernetes для BI/DWH выглядит так:
- Интеграция/оркестрация: Kafka (шина событий), Debezium (CDC из OLTP), Airflow/Dagster (оркестратор, SLA, зависимости).
- Хранение и запросный слой: Postgres Pro (оперативные витрины/метаданные), ClickHouse (быстрые витрины), Trino (единый SQL-движок по Lakehouse), S3-совместимое хранилище (MinIO/Ceph RGW), табличный формат Apache Iceberg + Hive Metastore (или REST-каталог Iceberg).
- Трансформации: dbt Core (в контейнере; адаптер Trino/ClickHouse/Postgres), расписания через Airflow.
- Публикация: Superset/Metabase через Ingress (см. Модуль 3).
K8s даёт масштабирование и управляемость, но данные/скорость формируются выбором форматов (Parquet/Iceberg), топологией Kafka, параметрами БД и грамотным CI/CD (см. Модуль 6).
Интеграция/оркестрация
Kafka на k8s
- Оператор Strimzi — стандарт де-факто: управляет кластерами Kafka/ZooKeeper или KRaft, топиками, Connect, MirrorMaker.
- Ресурсы и пулы: брокеры — на отдельном pool (таинт pool=kafka), быстрый диск под лог сегменты.
- Надёжность: минимум 3 брокера, фактор репликации топиков replication.factor=3, min.insync.replicas=2.
- Сетевые нюансы: externalTrafficPolicy: Local для сохранения client IP на Ingress; для межкластера — advertized listeners + TLS.
Debezium CDC (Postgres → Kafka)
- Требует логической репликации в Postgres Pro: wal_level=logical, max_replication_slots, max_wal_senders, wal_keep_size.
- Разворачивается как плагин Kafka Connect (через Strimzi KafkaConnect + KafkaConnector).
- Важно: формат ключа (pk) и стратегия дедупликации/упорядочивания событий, а также политика ретенции «мертвых» записей (tombstones).
Мини-пример (ядро конфигурации Debezium для PG):
{
"name": "pg-cdc",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "postgres.data-platform.svc",
"database.port": "5432",
"database.user": "debezium",
"database.password": "*****",
"database.dbname": "oltp",
"plugin.name": "pgoutput",
"slot.name": "debezium_slot",
"publication.name": "debezium_pub",
"tombstones.on.delete": "true",
"include.schema.changes": "true",
"table.include.list": "public.orders,public.customers",
"slot.drop.on.stop": "false"
}
}
Airflow vs Dagster
- Airflow — зрелая экосистема, много провайдеров, знаком большинству команд.
-
Dagster — типобезопасные «активы», тестируемость, прозрачная зависимость «данные→данные».
Рекомендация: начать с Airflow (операторный Helm чарт, Celery/K8sExecutor), затем внедрять Dagster там, где нужны «data-assets» и строгие декларации зависимостей.
Оркестрация CDC→Lakehouse:
- Airflow следит за здоровьем коннекторов, чистит «мелкие файлы» (compaction), запускает dbt, проверяет свежесть витрин (freshness SLI), публикует дашборды.
Хранилище и запросный слой
Postgres Pro на k8s
- Шаблон: StatefulSet + Patroni + pgBackRest/WAL-G (PITR; см. Модуль 2).
- Для CDC включить: wal_level=logical, увеличить max_replication_slots, max_wal_senders.
- Отдельный пул нод + RWO-блок (RBD/NVMe) для дисков; WAL на быстрых дисках; контролировать рост WAL (если коннектор отстаёт).
ClickHouse / Trino
-
ClickHouse — быстрые OLAP-витрины, колоночный формат, материальные представления, сверхнизкая латентность на агрегациях.
- В k8s: оператор Altinity/ClickHouse Operator; шардинг/репликация через shard/replica в манифестах.
- Укрупняйте партиции (обычно по дню/неделе) и используйте TTL/мердж-настройки под вашу нагрузку.
- Trino — федерационный SQL по S3/Iceberg, Postgres, ClickHouse, Kafka.
- Паттерн: 1 coordinator + N workers (HPA), отдельный пул нод под workers, JVM memory tuning, Spill to disk (быстрый RW).
- Драйверы: iceberg, hive (если метастор Hive), postgres, clickhouse.
Lakehouse: S3 + Iceberg + Metastore
- S3/MinIO — «истина» для файлов: сырые CDC-события, Parquet-слои, бэкапы. Включить версионирование bucket, lifecycle policy и репликацию (DR).
- Apache Iceberg — табличный формат с ACID, снапшотами, эволюцией схем; решает «малые файлы» через rewrite/compaction.
-
Каталог:
- Hive Metastore (проще, привычно);
-
или REST-каталог Iceberg (современнее, без Hive).
Рекомендация: начинать с Hive Metastore (оператор Hive/Metastore на k8s), Trino → Iceberg (каталог «hive»), Spark/Flink — по мере нужды.
Трансформации: dbt Core в контейнере
- dbt + Trino adapter: materialized='table'|'view'|'incremental', стратегия merge по unique_key.
- Храните profiles.yml и секреты подключения в External Secrets (см. Модуль 4).
- Оркестрируйте через Airflow: оператор BashOperator/KubernetesPodOperator для запуска контейнера dbt c volume мэппингом артефактов.
- Добавьте тесты dbt (not null, unique, relationships) и вычисление Freshness для ключевых источников.
Мини-пример модели dbt (смысл, не синтаксис):
-- models/fct_orders.sql
{{ config(materialized='incremental', unique_key='order_id', incremental_strategy='merge') }}
select * from {{ ref('stg_orders') }}
Дизайн конвейера CDC → S3/Iceberg → Trino → dbt → витрина → Superset
Два варианта «посадки» CDC
A. Kafka → Kafka Connect S3 Sink → Parquet в S3 → Iceberg ingestion/compaction (Spark/Flink job)
- просто, минимум спец. компонентов; – придётся компактифицировать и формировать «правильные» Iceberg-таблицы.
B. Kafka → Iceberg Sink (Kafka Connect Iceberg или Flink SQL CDC → Iceberg)
- сразу ACID-таблицы, мин. «мелких файлов»; – потребуется соответствующий коннектор/движок.
На практике многие начинают с варианта A, затем переходят на B, когда нужно ACID и эволюция схем «из коробки».
Слои данных
- raw — «как прилетело» (CDC событийный поток, Avro/JSON/Parquet).
- staging — выровненные типы/UTC, dedup (по op_ts,lsn,pk), сняты tombstones.
- curated (Iceberg) — нормализованные таблицы фактов/измерений, партиционирование по дате/бизнес-ключам.
- marts/vitrines — агрегаты, готовые к BI.
Паттерны качества/актуальности
- Dedup + upsert при сборке staging → curated (dbt incremental merge по pk + op_ts).
- Freshness-чек: задача Airflow проверяет «время последнего снапшота» Iceberg.
- Контроль схемы: включить include.schema.changes=true в Debezium, а в ingestion-job учитывать эволюцию колонок.
Безопасность и управление доступом в стеке
- Kafka: TLS/mTLS, ACL на топики; sasl.scram если требуется.
- S3: ограниченные политики на бакеты (raw, staging, curated, backups), pre-signed URLs для экспорта.
- Trino: file-based или LDAP/OIDC аутентификация, SQL-гранты по схемам; авторизацию на Iceberg-каталоге (read vs write).
- Секреты — через External Secrets + Vault; токены — краткоживущие (см. Модуль 4).
- Сетевые политики (NetworkPolicy): BI/Trino/Metastore/S3 — только необходимые направления (см. Модуль 4).
Эксплуатация и производительность
- Kafka: следите за under-replicated, задержками продьюсеров/консьюмеров, размером сегментов, ретенциями.
- S3/Iceberg: «болезнь малых файлов» — плановые компакшены (Spark/Flink/Iceberg rewrite). Партиции: чаще по дню/часу; не «перепартиционируйте».
- Trino: JVM tuning (heap, GC), localized spill на быстрый диск, отдельные пулы workers, настройка query.max-memory/max-total-memory-per-node.
- dbt: параллелизм через --threads; ограничить одновременные тяжёлые модели; тесты — nightly.
- Postgres Pro: вакуум/анализ, контроль replication slot lag (Debezium), RPO/RTO по бэкапам (см. Модуль 2).
Риски и как их снижать
|
Риск |
Проявление |
Митигировать |
|---|---|---|
|
Логическая реплика «отстаёт» |
рост WAL, диски БД заполняются |
мониторинг slot lag, алерты; масштабирование Connect; временно поднять ретенцию WAL; PITR-план |
|
Малые файлы в S3 |
медленные запросы Trino, больше метаданных |
компакшены; Iceberg rewrite data files; увеличивать размер файла/микробатчи в sink |
|
Несогласованность типов/часы |
«дубли» и неверные агрегаты |
normalize в staging: UTC, точные типы; единый op_ts/event_time |
|
Сломанная схема (evolution) |
Trino ошибки «колонка X отсутствует» |
включить schema changes; поколенческий rolling-дебют изменений через dbt; тесты схем |
|
Дубликаты CDC |
MERGE не отрабатывает |
правильный ключ (pk + op_ts/lsn), dedup правила в ingestion |
|
Бутылочное горлышко Metastore |
медленные планирования Trino |
ресурсы метастора, кеш Trino, переход на REST-каталог Iceberg |
|
«Шум» DAG/оркестрации |
сложно поддерживать |
декомпозиция jobs; SLAs и дашборды (Модуль 5); шаблоны DAG; OpenLineage |
Мини-фрагменты конфигурации (ровно столько, сколько нужно)
Trino → Iceberg (каталог через Hive Metastore)
etc/catalog/iceberg.properties (смысл):
connector.name=iceberg hive.metastore.uri=thrift://hms.metastore.svc:9083 iceberg.catalog.type=hive iceberg.file-format=PARQUET fs.native-s3.enabled=true s3.endpoint=http://minio.minio:9000 s3.path-style-access=true s3.aws-access-key=<via env/secret> s3.aws-secret-key=<via env/secret>
Kafka Connect → S3 Sink (вариант A)
Идея: писать Parquet-файлы по топикам в s3://raw/orders/, дальше компакшн/ingestion в Iceberg.
{
"name": "s3-sink-orders",
"config": {
"connector.class": "io.confluent.connect.s3.S3SinkConnector",
"topics": "public.orders",
"s3.bucket.name": "raw",
"s3.part.size": "134217728",
"format.class": "io.confluent.connect.s3.format.parquet.ParquetFormat",
"flush.size": "10000",
"rotate.interval.ms": "600000",
"schema.compatibility": "FORWARD"
}
}
Дополните аутентификацию к MinIO и префиксы путей. Позже — ingestion job в Iceberg.
Практика: построить конвейер CDC→S3/Iceberg→Trino→dbt→витрина→Superset
Цель
С нуля собрать «сквозной» путь для одной таблицы orders и зависимой customers:
- Debezium читает изменения из Postgres Pro и пишет в Kafka.
- S3 Sink сохраняет Parquet в s3://raw/....
- Ingestion job формирует Iceberg-таблицы stg_orders/stg_customers (и делает компакшн).
- Trino видит Iceberg-таблицы; dbt строит fct_orders (incremental merge) и dim_customers.
- Superset подключается к Trino и показывает витрину.
Предпосылки
- Кластер k8s с MinIO/S3, Strimzi Kafka, Postgres Pro (настроен logical), Trino + Hive Metastore, Airflow, Superset (см. предыдущие модули).
- Секреты через External Secrets, сеть ограничена NetworkPolicy.
Шаги
1) Postgres Pro (OLTP) подготовка
- Включить wal_level=logical, задать max_replication_slots.
- Создать пользователя debezium с правами REPLICATION/SELECT на нужные таблицы.
- Добавить индексы по PK/времени, чтобы не душить прод.
2) Kafka + Debezium
- Развернуть KafkaConnect со сборкой Debezium плагина.
- Создать KafkaConnector для orders/customers (включить include.schema.changes=true).
- Проверить, что события идут в топики public.orders, public.customers.
3) S3 Sink (вариант A)
- Завести S3 Sink для public.orders/public.customers → s3://raw/… в Parquet.
- Настроить flush.size, rotate.interval.ms так, чтобы файлы были ~128–256 МБ.
4) Ingestion → Iceberg (минимально жизнеспособно)
-
Запустить периодический Spark/Flink job (через Airflow) для:
- чтения raw Parquet;
- дедуп/нормализация типов/UTC;
- запись в Iceberg таблицы stg_orders, stg_customers c партиционированием по дате (op_date).
- В конце job — Iceberg rewrite data files (компакшн).
5) Trino каталог
- Убедиться, что iceberg каталог в Trino настроен, видит метастор и бакет curated.
- Проверить SHOW TABLES → stg_orders, stg_customers.
6) dbt Core (адаптер Trino)
- Подготовить проект: stg_* → fct_orders, dim_customers.
- Модели incremental с unique_key и merge on (по order_id / customer_id).
- Добавить базовые dbt tests (not null/unique) и freshness.
- Запускать из Airflow KubernetesPodOperator (контейнер с dbt).
7) BI (Superset)
- Подключить к Trino (catalog iceberg, schema marts).
- Создать dataset fct_orders; собрать дашборд (выручка, заказы, конверсия).
8) Наблюдаемость и SLO (Модуль 5)
- Метрики Debezium/Connect lag, Iceberg ingestion длительность, Trino success rate/p95, freshness витрины, Superset availability.
- Алерты: отставание CDC (lag), свежесть витрин, рост 5xx Ingress.
Критерии готовности (Definition of Done)
- CDC-лаг стабилен < 1–2 минут, WAL не растёт бесконтрольно.
- В Iceberg таблицах нет «пылевых» файлов (регулярный compaction).
- dbt run проходит < N минут, тесты зелёные.
- Trino отвечает по fct_orders с p95 < 10 с (на заданном объёме).
- Дашборд Superset показывает корректные метрики и обновляется в SLA.
Чек-лист перед продом
- Strimzi Kafka с 3+ брокерами, TLS/ACL, ретенции по политике.
- Debezium с мониторингом slot lag; параметры Postgres Pro для logical replication.
- S3/MinIO: версионирование, lifecycle, репликация DR.
- Iceberg + каталог (Hive/REST), план компакшнов; стратегии партиционирования.
- Trino: отдельные пулы workers, JVM tuning, spill-диски, каталоги настроены.
- dbt: инкрементальные модели с merge, тесты и freshness; запуск из Airflow.
- Superset: SSO/Ingress, таймауты/лимиты; источники — через Trino.
- Секреты/доступ: ESO+Vault, NetworkPolicy, роли в Trino/S3, токены краткоживущие.
- Наблюдаемость: SLO для CDC/ingestion/Trino/BI, алерты burn-rate.
- Runbooks: «Debezium отстаёт», «малые файлы», «Trino тормозит», «dbt упал», «PITR Postgres».
Надёжный стек данных в k8s — это связка: Kafka+Debezium для «живых» изменений, S3+Iceberg как устойчивый слой таблиц, Trino как универсальный SQL-движок, dbt как декларативные трансформации, Airflow как оркестратор, Superset как витрина. Секрет — в правильной «посадке» CDC, борьбе с малыми файлами, инкрементальном merge и чётких SLO. Когда всё это завернуто в GitOps и наблюдаемость, BI/DWH становится быстрым и управляемым.




