Блок 7 - Kafka, Streaming и мониторинг

25 вопросов: Kafka иерархия, log compaction, Structured Streaming triggers, Python зависимости, SparkMeasure, Airflow паттерны.

core streaming

Что такое Apache Kafka и чем отличается от обычной очереди?

Apache Kafka - распределённый брокер сообщений по модели publish/subscribe, созданный LinkedIn в 2011 году. В отличие от обычных очередей (P2P, FIFO, одно сообщение = один Consumer):

Обычная очередь (RabbitMQ) Apache Kafka
Модель Point-to-Point Publish/Subscribe
Хранение Удаляет после доставки Хранит по retention (default 7 дней)
Параллелизм Один Consumer на сообщение Несколько Consumer Groups, каждая получает всё
Порядок FIFO для всей очереди FIFO в рамках одной Partition
Масштабирование Вертикальное Горизонтальное (партиции)

Ключевая идея Kafka: append-only log с retention - данные не удаляются после чтения, Consumer сам управляет своим Offset'ом.


Опишите иерархию объектов Kafka: Cluster → Topic → Partition → Segment

Kafka Cluster
├── Broker 1 (Kafka Server)        ← физический сервер Kafka
│   ├── Topic "events"             ← логическая очередь
│   │   ├── Partition 0 [LEADER]   ← физический лог
│   │   │   ├── Segment 1.log      ← файл (~1 GB)
│   │   │   ├── Segment 2.log      ← закрытый
│   │   │   └── Segment 3.log      ← активный (текущая запись)
│   │   ├── Partition 1 [FOLLOWER] ← реплика (на другом Broker)
│   │   └── Partition 2 [LEADER]
│   └── Topic "orders"
├── Broker 2
└── Broker 3
  • Topic - логическая очередь, append-only, состоит из 1+ Partition
  • Partition - физический log; упорядоченная последовательность; key routing: hash(key) % N; без ключа → round-robin
  • Segment - физический файл внутри Partition; один активный; при превышении log.segment.bytes (1 GB) → новый
  • Broker - Kafka-сервер; один Leader + N Follower реплик на каждую Partition

Что такое Offset и Consumer Group?

Offset - монотонно возрастающий порядковый номер сообщения в Partition. Уникален в рамках одной партиции. Consumer Group хранит committed offset для каждой партиции - это «закладка» для возобновления после рестарта.

Consumer Group - группа Consumer'ов, совместно читающих Topic:

Topic "events": [Partition 0] [Partition 1] [Partition 2]

Group "analytics":
  Consumer A → Partition 0 (offset 100)
  Consumer B → Partition 1, 2 (offset 50, 75)

Group "reporting":
  Consumer X → Partition 0, 1, 2 (читает независимо от Group "analytics")

Правило: Одна Partition → максимум один Consumer в группе. Параллелизм ограничен числом партиций. Разные Consumer Group получают весь Topic независимо.


Как работает роутинг сообщений по партициям?

# С ключом: partition = hash(key) % num_partitions
producer.send("events", key=b"user_123", value=event_json)
# hash("user_123") % 3 = 1 → всегда Partition 1
# Гарантия: все события user_123 → одна партиция → сохранён порядок

# Без ключа: round-robin по партициям
producer.send("events", value=event_json)
# Message 1 → Partition 0
# Message 2 → Partition 1
# Message 3 → Partition 2
# Message 4 → Partition 0 (снова)

Когда использовать ключ: Когда нужен порядок событий для одной сущности (например, все события пользователя должны обрабатываться последовательно в одном Consumer).


Что такое Log Retention и Cleanup Policy в Kafka?

Log Retention - политика хранения данных. Два типа:

# По времени (приоритет: ms > minutes > hours)
log.retention.hours=168          # 7 дней (default)
log.retention.ms=604800000       # эквивалентно, но приоритетнее

# По размеру
log.retention.bytes=1073741824   # 1 GB на партицию (-1 = без лимита)

Cleanup Policy - что делать по истечении retention:

Policy Поведение
delete (default) Физически удалить старые сегменты
compact Оставить только последнее сообщение для каждого ключа
delete,compact Сначала compact, потом delete

Как работает Log Compaction?

Log Compaction применяется при cleanup.policy=compact. Цель: для каждого ключа гарантировать наличие хотя бы последнего значения.

До compaction (один из сегментов):
offset 1: key=user1, value={"age": 25}
offset 2: key=user2, value={"name": "Bob"}
offset 5: key=user1, value={"age": 26}   ← обновление
offset 8: key=user1, value={"age": 27}   ← ещё обновление

После compaction:
offset 2: key=user2, value={"name": "Bob"}  (единственное значение)
offset 8: key=user1, value={"age": 27}      (только последнее)

Структура лога:

  • Log Head - свежая часть с монотонными offset'ами
  • Log Tail (compacted) - старая часть с «дырами» в offset'ах
  • Cleaner Point - граница между Head и Tail

Применение: CDC (Change Data Capture), key-value state, changelog топики.


Чем Structured Streaming лучше DStreams?

DStreams (deprecated) Structured Streaming
Введён Spark 0.7 Spark 2.0
Статус Deprecated Stable (с 2.2)
API RDD-based DataFrame API
Оптимизация Нет Catalyst Catalyst + Tungsten
Fault tolerance WAL + checkpointing Checkpoint + idempotent sink
Exactly-once Сложно реализовать Встроено
Stateful ops Ограниченные Полноценные (window, watermark)

Вывод: DStreams не использовать в новом коде. Structured Streaming = тот же DataFrame API, что и для batch.


Какие Sources и Sinks есть в Structured Streaming?

Sources (источники):

Source Форматы Применение
file Parquet, JSON, CSV, ORC Чтение новых файлов из директории
kafka Kafka messages Production стриминг
socket text Отладка (нет fault tolerance)
rate synthetic rows Нагрузочное тестирование

Sinks (приёмники):

Sink Применение Output Mode
file Parquet/JSON в S3/HDFS append
kafka Запись в Kafka topic append
foreachBatch Произвольная логика + несколько outputs -
foreach Построчная обработка append
console Отладка append, update, complete
memory Тестирование (temp view) append, complete
# foreachBatch - самый гибкий вариант
def write_batch(batch_df, batch_id):
    batch_df.write.parquet(f"s3a://bucket/out/batch={batch_id}")
    batch_df.write.jdbc(url, table, "append")  # два sink'а!

stream.writeStream.foreachBatch(write_batch).start()

Какие Triggers есть в Structured Streaming и когда использовать?

Trigger API Latency Применение
Unspecified (не задавать) Минимальная Max throughput: следующий batch сразу после предыдущего
Fixed interval .trigger(processingTime="5 minutes") 5 минут Периодические отчёты
Once .trigger(once=True) - Backfill: прочитать всё накопленное и завершить
Available Now .trigger(availableNow=True) - Batches до исчерпания накопленного
Continuous .trigger(continuous="1 second") ~1ms Экспериментальный, без stateful ops
# Backfill pattern: обработать всё и завершить
query = stream.writeStream \
    .trigger(once=True) \
    .format("parquet") \
    .option("checkpointLocation", "s3a://bucket/ckpt/") \
    .start()
query.awaitTermination()

Continuous trigger экспериментален и не поддерживает aggregation и stateful операции. Для production используйте Unspecified или Fixed interval.


Как обеспечить exactly-once в Structured Streaming?

Exactly-once = каждое событие обработано ровно один раз, без потерь и дублей.

Механизм: Checkpoint + idempotent sink.

query = stream.writeStream \
    .format("iceberg") \
    .option("checkpointLocation", "s3a://bucket/checkpoints/my_stream/") \
    .start()
# Checkpoint хранит: последний обработанный offset для каждой partition/topic
# При рестарте: читаем с offset'а после last committed

При рестарте:

  1. Spark читает checkpoint → знает последний offset
  2. Kafka отдаёт данные начиная с offset+1
  3. Если sink idempotent (Iceberg MERGE, Delta MERGE) → нет дублей при повторной записи

Без idempotent sink (append-only файловый sink) → at-least-once при сбое в момент записи.


Как доставить Python-зависимости на все Executor'ы Spark?

Пять методов:

Метод Нативные libs (.so) Сложность Применение
Установка на каждую ноду Низкая Stable окружение, admin доступ
--py-files eggs/zip Низкая Pure Python, без C-расширений
Conda Pack Средняя Production с NumPy/pandas
venv + venv-pack Средняя Production, стандартный Python
PEX Средняя Один файл

Conda Pack (рекомендуется):

conda create -n myenv python=3.10 && conda activate myenv
pip install pandas numpy scikit-learn
conda pack -o myenv.tar.gz

spark-submit \
    --archives myenv.tar.gz#environment \
    --conf spark.executorEnv.PYSPARK_PYTHON=./environment/bin/python \
    job.py

venv + venv-pack:

python -m venv myenv && source myenv/bin/activate
pip install -r requirements.txt
venv-pack -o myenv.tar.gz

spark-submit \
    --archives myenv.tar.gz#venv \
    --conf spark.executorEnv.PYSPARK_PYTHON=./venv/bin/python \
    job.py

ENV propagation: Переменные окружения Driver автоматически передаются на все Executor'ы (только ENV vars, не Python-пакеты).


Что такое SparkMeasure и зачем он нужен?

SparkMeasure - open-source библиотека для программного сбора Stage и Task метрик изнутри Spark-приложения. Реализует кастомный Spark Listener.

from sparkmeasure import StageMetrics

stage_metrics = StageMetrics(spark)
stage_metrics.begin()

# Ваш код
df.groupBy("key").agg(F.sum("amount")).collect()

stage_metrics.end()
stage_metrics.print_report()
# → executorRunTime: 4512 ms
# → shuffleReadBytes: 1.2 GB
# → jvmGCTime: 312 ms
# → peakExecutionMemory: 2.8 GB
# → memoryBytesSpilled: 0

Ключевые метрики:

Метрика Что значит
peakExecutionMemory Пиковое потребление execution memory
jvmGCTime Время GC (> 10% → проблема)
shuffleBytesWritten Объём shuffle - коррелирует с числом wide трансформаций
memoryBytesSpilled Spill на диск (> 0 → нехватка памяти)

FlightRecorder-режим - сбор метрик в фоне без явного begin/end, запись в файл/InfluxDB/Kafka.


Зачем настраивать spark.metrics.namespace?

По умолчанию:

spark.metrics.namespace = spark.app.id
# Пример: spark.application_1680000000000_0001

Проблема: app.id уникален для каждого запуска → в InfluxDB/Grafana каждый запуск создаёт новый namespace. Невозможно видеть тренды и историю одного приложения.

Решение:

spark.metrics.namespace=${spark.app.name}
# daily_etl_sales → всегда spark.daily_etl_sales.executor.*

Это позволяет Grafana показывать историю daily_etl_sales через десятки запусков: тренды jvmGCTime, shuffleBytesWritten, peakExecutionMemory по времени.


Что такое GraphiteSink и куда он пишет метрики?

GraphiteSink - встроенный Spark Sink для отправки метрик через Graphite Carbon protocol. Может писать в системы, поддерживающие этот протокол:

Система Поддержка
Graphite Нативно
InfluxDB Через Graphite input plugin
ClickHouse Через graphite-clickhouse
# metrics.properties
*.sink.graphite.class=org.apache.spark.metrics.sink.GraphiteSink
*.sink.graphite.host=influxdb.internal
*.sink.graphite.port=2003
*.sink.graphite.period=10
*.sink.graphite.unit=seconds

Типичный стек: GraphiteSink → InfluxDB → Grafana Dashboard.


Какие типы Executor в Airflow? Чем отличаются?

Airflow Executor определяет ГДЕ выполняются Airflow Tasks (не путать со Spark Executor).

Executor Где Task Масштабирование Когда
LocalExecutor На той же машине (процессы) Нет Dev, малые нагрузки
CeleryExecutor На отдельных Worker-нодах (через Redis/RabbitMQ) Горизонтальное Production, несколько нод
KubernetesExecutor В отдельных K8s Pod (один Pod = один Task) Авто (K8s) Cloud-native, K8s кластеры
LocalExecutor: Scheduler → [fork] → Task Process (та же машина)
CeleryExecutor: Scheduler → [Redis Queue] → Worker Node 1, 2, N
KubernetesExecutor: Scheduler → K8s API → Pod per Task (создаётся/удаляется)

Какие антипаттерны Airflow + Spark нужно избегать?

Антипаттерн 1: PythonOperator обрабатывает данные

# ❌ Неверно: чтение больших данных на Airflow-машине
def process_data():
    import pandas as pd
    df = pd.read_csv("s3://bucket/100gb_file.csv")  # 100 GB на Airflow-сервере!
    result = df.groupby("key").sum()

task = PythonOperator(task_id="process", python_callable=process_data)

Антипаттерн 2: BashOperator запускает Spark локально

# ❌ Неверно: Spark работает прямо на Airflow-машине
task = BashOperator(
    task_id="run_spark",
    bash_command="python /opt/jobs/etl.py"  # Нет кластера, нет параллелизма!
)

Правильно:

# ✓ Вариант 1: KubernetesPodOperator - изолированный контейнер
from airflow.providers.cncf.kubernetes.operators.pod import KubernetesPodOperator

task = KubernetesPodOperator(
    task_id="etl",
    image="myregistry/spark-job:v1.2.3",
    cmds=["python", "/app/etl.py"],
    arguments=["--date", "{{ ds }}"],
)

# ✓ Вариант 2: SparkSubmitOperator - submit на кластер
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator

task = SparkSubmitOperator(
    task_id="etl",
    application="s3a://bucket/jobs/etl.py",
    conf={"spark.executor.memory": "4g"},
)

Принцип: Airflow - оркестратор, а не вычислительный узел. Если Airflow-нода упала, Task перезапустится - при этом Spark на кластере продолжает работать независимо.


Чем Cluster Mode лучше Client Mode для production?

Client Mode Cluster Mode
Driver На машине запуска spark-submit Внутри YARN Container / K8s Pod
Потеря машины запуска Job падает Job продолжается
Логи Driver В терминале Внутри кластера (yarn logs / kubectl logs)
Сетевой трафик Driver ↔ Executor через сеть клиента Всё внутри кластера
Применение Разработка, Jupyter Production

В YARN Cluster Mode:

Client → spark-submit → YARN ResourceManager
  → NodeManager → YARN Container
      ├── Spark Application Master
      └── Spark Driver (здесь, в кластере)
         → Запрашивает Executor контейнеры у ResourceManager

Driver изолирован в кластере → клиентская машина (или Airflow worker) может быть отключена, Job продолжится.


Как читать логи Executor'а в YARN?

# Список контейнеров приложения
yarn logs -applicationId application_1680000000000_0001

# Логи конкретного контейнера
yarn logs -applicationId application_1680000000000_0001 \
          -containerId container_e01_1680000000000_0001_01_000002

# Признаки OOM в stderr Executor:
# java.lang.OutOfMemoryError: Java heap space
# GC overhead limit exceeded
# Container killed by YARN for exceeding memory limits

# Просмотр в Spark UI:
# Executors tab → stdout/stderr links

Опишите полный стек мониторинга Spark в production

Spark Driver / Executor
    ↓ (GraphiteSink - metrics.properties)
InfluxDB (Graphite input plugin, порт 2003)
    ↓
Grafana Dashboard
    ↓ (алерты)
PagerDuty / Slack

Ключевые Grafana панели:

  1. jvmGCTime % - время в GC / task duration (алерт > 15%)
  2. peakExecutionMemory - пиковое потребление execution memory
  3. shuffleBytesWritten - объём shuffle по времени
  4. memoryBytesSpilled - spill на диск (алерт > 0)
  5. Task Duration Distribution - гистограмма (длинный хвост → skew)
  6. Active Jobs/Stages - нагрузка кластера
# Минимальная конфигурация для production мониторинга
spark.conf.set("spark.metrics.namespace", "${spark.app.name}")
spark.conf.set("spark.eventLog.enabled", "true")
spark.conf.set("spark.eventLog.dir", "s3a://bucket/spark-logs/")
# metrics.properties настраивается отдельно с GraphiteSink

Что такое DPP и когда он помогает?

DPP (Dynamic Partition Pruning) - оптимизация для star schema паттерна. Spark вычисляет фильтр из dimension-таблицы и инжектирует его как subquery прямо в scan таблицы фактов.

Запрос: SELECT ... FROM FACT JOIN DIM ON FACT.date_id = DIM.date_id WHERE DIM.year = 2023

Без DPP:                          С DPP:
1. Read ALL FACT (100 GB)         1. Read DIM WHERE year = 2023
2. Read DIM WHERE year = 2023        → date_id ∈ {1, 2, 5, 7, ...}
3. Join + filter                  2. Read FACT WHERE date_id IN (...) [20 GB!]
                                  3. Join (уже отфильтрованные данные)

Условия для DPP:

  • FACT-таблица партиционирована по ключу join (date_id)
  • DIM-таблица достаточно мала для broadcast
  • Есть фильтр на DIM-таблице

Результат на TPC-DS: многие запросы ускоряются в 10–20× (q25: 390 с → 20 с).

spark.conf.set("spark.sql.optimizer.dynamicPartitionPruning.enabled", "true")  # default ON

Как Kafka интегрируется со Structured Streaming?

# Чтение из Kafka
stream_df = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "kafka:9092") \
    .option("subscribe", "events,orders") \
    .option("startingOffsets", "earliest")  # или "latest"
    .option("maxOffsetsPerTrigger", 100000)  # ограничить размер batch
    .load()

# stream_df.schema: key, value, topic, partition, offset, timestamp, ...
from pyspark.sql.functions import from_json, col
from pyspark.sql.types import StructType, StringType, LongType

schema = StructType().add("user_id", LongType()).add("event", StringType())

parsed = stream_df.select(
    from_json(col("value").cast("string"), schema).alias("data"),
    col("partition"),
    col("offset")
).select("data.*", "partition", "offset")

# Запись обратно в Kafka
result = parsed.groupBy("event").count()

result.writeStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "kafka:9092") \
    .option("topic", "event_counts") \
    .option("checkpointLocation", "s3a://bucket/ckpt/kafka_out/") \
    .outputMode("update") \
    .start()

Checkpoint: Spark хранит committed offset для каждой partition каждого topic. При рестарте продолжает с последнего committed offset.


Назовите основные Catalyst логические оптимизации

Catalyst применяет rule-based оптимизации к Logical Plan:

Оптимизация Пример
Predicate Pushdown filter("year = 2023") → читает только нужные row groups Parquet
Constant Folding WHERE amount * 1.0 > 0 + 100WHERE amount > 100.0
Column Pruning SELECT id FROM t → читает только колонку id из Parquet
Empty Relation Propagation JOIN (WHERE 1=0) → skip всего join
Boolean Simplification a AND TRUEa
# Проверить что pushdown работает
df.filter("year = 2024").select("user_id").explain()
# В плане:
# PushedFilters: [EqualTo(year, 2024)]
# ReadSchema: struct<user_id:bigint>   ← только нужная колонка

Почему не рекомендуется более 5 cores на Executor?

Правило «5 cores на Executor» - рекомендация Cloudera для HDFS-кластеров.

Причины:

  1. HDFS I/O конкуренция: При большом числе одновременных Task на одном JVM-процессе HDFS-клиент испытывает конкуренцию за соединения → throughput падает нелинейно
  2. GC-паузы: Большой heap с большим числом Task → длинные GC-паузы затрагивают все одновременные Task Executor'а
  3. Memory per task: Больше ядер = меньше памяти на каждый Task → риск OOM

Формула для узла 16 CPU / 64 GB RAM:

Cores на Executor: 5
Executors на узел: (16 - 1) / 5 = 3  ← 1 core оставить на ОС/YARN
RAM на Executor: (64 - 4) / 3 ≈ 20 GB  ← 4 GB на ОС
memoryOverhead: 2 GB
executor.memory: 18 GB

На S3-кластерах (без HDFS) это правило менее критично - можно пробовать 6–8 cores при правильном тюнинге GC.