Блок 7 - Kafka, Streaming и мониторинг
25 вопросов: Kafka иерархия, log compaction, Structured Streaming triggers, Python зависимости, SparkMeasure, Airflow паттерны.
Что такое 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
При рестарте:
- Spark читает checkpoint → знает последний offset
- Kafka отдаёт данные начиная с offset+1
- Если 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 панели:
- jvmGCTime % - время в GC / task duration (алерт > 15%)
- peakExecutionMemory - пиковое потребление execution memory
- shuffleBytesWritten - объём shuffle по времени
- memoryBytesSpilled - spill на диск (алерт > 0)
- Task Duration Distribution - гистограмма (длинный хвост → skew)
- 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 + 100 → WHERE amount > 100.0 |
| Column Pruning | SELECT id FROM t → читает только колонку id из Parquet |
| Empty Relation Propagation | JOIN (WHERE 1=0) → skip всего join |
| Boolean Simplification | a AND TRUE → a |
# Проверить что 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-кластеров.
Причины:
- HDFS I/O конкуренция: При большом числе одновременных Task на одном JVM-процессе HDFS-клиент испытывает конкуренцию за соединения → throughput падает нелинейно
- GC-паузы: Большой heap с большим числом Task → длинные GC-паузы затрагивают все одновременные Task Executor'а
- 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.