Введение в PySpark и Spark SQL —
от архитектуры до production
Мастерство PySpark: от новичка до продвинутого уровня. Free PySpark course на русском — 301 лекций, Catalyst, AQE, Iceberg, Kafka, ClickHouse и self-host платформа без SaaS-привязки.
21 этап обучения — от новичка до продвинутого уровня
Лучший бесплатный PySpark tutorial на русском: от введения в PySpark и Spark internals до оптимизации, Lakehouse и стриминга. Каждый этап — отдельный модуль с лекциями и практикой.
- Что такое Apache Spark: место в экосистеме Big Data
- Архитектура кластера: Driver, Executor и Cluster Manager
- Spark vs Hadoop MapReduce: почему in-memory и lazy evaluation
- RDD, DataFrame и Dataset: UnsafeRow, Encoders, эволюция абстракций
- Transformations vs Actions: почему Spark ленив до последнего
- Narrow vs Wide трансформации: когда возникает shuffle
- SparkSession: точка входа, конфигурация и master URL
- spark-submit: упаковка зависимостей и запуск на кластере
- Локальная установка PySpark и первый job с нуля
- Глоссарий: ~100 терминов - от Action/DAG/Shuffle до Salting, DPP, Schema Evolution
- Введение в PySpark: что это и зачем
- Настройка SparkSession: конфигурации ресурсов, памяти и логирования
- Чтение и запись файлов: CSV, JSON, Parquet, Hive, JDBC
- Чтение CSV: опции header, inferSchema, delimiter, nullValue, режимы
- Чтение JSON: однострочный и многострочный форматы, обработка ошибок
- Обращение к колонкам: строки, dot-нотация, col(), скобочная нотация, lit()
- Выборка колонок: select(), withColumn(), alias(), drop(), вложенные структуры
- Фильтрация данных: filter(), where(), isin(), like(), rlike(), isNull()
- Группировка данных: groupBy(), agg(), count(), sum(), avg(), countDistinct()
- Объединение данных: inner, left, right, outer, cross, semi, anti join
- Пивотирование: groupBy().pivot().agg() — длинный формат в широкий
- Работа с NULL: na.fill(), dropna(), na.replace(), фильтрация по null
- Функции работы с датами: to_date, date_format, datediff, date_add, trunc
- Математические функции: арифметика, abs, round, sqrt, pow, log, тригонометрия
- Строковые функции: concat, substring, regexp_extract, upper, lower, trim, lpad
- Оконные функции: Window, partitionBy, orderBy, rank, dense_rank, row_number
- Lead и Lag: доступ к следующей и предыдущей строке в окне
- rowsBetween: скользящие суммы и средние, управление диапазоном окна
- DAG Scheduler: как трансформации превращаются в задачи
- Stage Boundary: почему shuffle создаёт новый Stage
- Task Scheduling: FIFO vs FAIR, locality levels
- Catalyst Optimizer: 4 фазы от AST до Physical Plan
- Rule-Based Optimization: constant folding, predicate pushdown, column pruning
- Cost-Based Optimizer: ANALYZE TABLE и статистика колонок
- Tungsten: whole-stage codegen и почему Java bytecode быстрее
- Unified Memory Model: execution vs storage регионы и граница
- Spill на диск: причины, диагностика в Spark UI и как избежать
- GC Pressure: G1GC настройки, off-heap память и сигналы в логах
- Py4J Gateway и Python Workers: как Python-код переходит в JVM
- Data Locality: PROCESS_LOCAL → ANY, spark.locality.wait и влияние S3
- Жизненный цикл приложения: CoarseGrainedExecutorBackend, Block Manager и ExternalShuffleService
- Speculative Execution: straggler detection, конфигурация и когда отключать
- Spark Connect: gRPC thin client - удалённое подключение без локального Spark и архитектура протокола
- DataFrame API: select, filter, withColumn, alias и Column expressions
- Агрегации: groupBy, agg, rollup, cube и grouping sets
- Window Functions: frame specification, ROWS vs RANGE, UNBOUNDED
- Window Functions Advanced: сложные примеры в DWH и Data Lake
- String functions: regexp_replace, concat_ws, split, substring, length - когда заменяют UDF
- Date/Time functions: to_date, date_add, datediff, year/month/dayofweek, from_unixtime, timezone
- Conditional columns: when/otherwise, coalesce, nullif и NULL-безопасные сравнения
- sample() и randomSplit(): стратифицированная выборка и train/test split
- Complex Types: работа с ArrayType, MapType и функциями higher-order
- StructType: вложенные схемы, schema inference и schema evolution
- UDF: Python UDF и их цена - когда стоит, когда не стоит
- Pandas UDF (Arrow): SCALAR, GROUPED_MAP, GROUPED_AGG
- Join стратегии: broadcast, sort-merge, shuffle hash и join hints
- Форматы данных: Parquet, ORC, Avro - выбор, настройка, pushdown
- Read options: inferSchema vs explicit schema, PERMISSIVE/DROPMALFORMED/FAILFAST, badRecordsPath
- Actions в деталях: foreachPartition connection-pool, checkpoint для длинных DAG, toPandas safety
- Тестирование PySpark: pytest fixtures, chispa, паттерн testable transforms
- Pandas API on Spark (pyspark.pandas): pandas-синтаксис на кластере, Arrow, ограничения
- DataFrame Wrangling: left_semi/left_anti, set operations, pivot, posexplode и sampleBy
- EXPLAIN FORMATTED: чтение логического и физического плана пошагово
- Аналитические запросы: GROUPING SETS, ROLLUP, CUBE без UDF
- LATERAL VIEW и explode: unnesting массивов в строки
- PIVOT: транспонирование строк в колонки и обратно
- JSON: from_json, to_json, get_json_object, schema_of_json
- Semi-structured данные: вложенные struct, flatten и нормализация
- Catalog API: CREATE TABLE, ALTER TABLE, DESCRIBE DETAIL
- Views и Temp Views: persistence, scoping и подводные камни
- Hive Metastore: external vs managed tables и partitioned tables DDL
- ANSI Mode: строгая типизация, implicit cast, division by zero и миграция
- Spark Config: тюнинг памяти, shuffle, AQE и PySpark-специфика
- Shuffle Partitions: как выбрать spark.sql.shuffle.partitions вручную
- AQE: adaptive coalescing - автоматический подбор числа партиций
- AQE: автоматическое переключение join стратегий в рантайме
- Dynamic Partition Pruning: механизм star-schema и условия срабатывания
- Bucketing: sort-merge join без shuffle - настройка и ограничения
- Storage Partitioning: выбор колонок для partitionBy и анти-паттерны
- Small Files Disease: причины, последствия и влияние на планировщик
- Compaction: стратегии укрупнения файлов через Spark и Iceberg
- Data Layout: Z-ordering, Hilbert curves и file skipping
- S3A Connector: архитектура, конфигурация и параметры производительности
- Rename Problem: почему S3 не POSIX и как это ломает Spark
- Committer Algorithms: magic committer vs staging committer vs partitioned
- Parquet на S3: vectorized reader, column pruning и row group filtering
- Predicate Pushdown в Object Storage: что пушится, что нет
- Small Files на S3: LIST latency и ограничения при большом числе файлов
- Compaction стратегии: Spark job, Iceberg optimize, auto-vacuum
- MinIO: установка, настройка bucket policy и интеграция со Spark
- Boto3: Python SDK для S3-admin операций - листинг, удаление объектов, Presigned URLs и metadata
- LocalStack и Testcontainers: локальная S3-эмуляция - ETL-тесты без реального кластера
- Сравнение S3-совместимых хранилищ: MinIO, Ceph, SeaweedFS, Garage, Apache Ozone
- HDFS Architecture: NameNode, DataNode, блоки, репликация и heartbeat
- NameNode HA: QJM quorum, Standby NameNode и автоматический failover
- Block Locality: как Spark scheduler использует data-locality при планировании
- Short-Circuit Reads: чтение DataNode без сетевого стека через Unix socket
- Erasure Coding: Reed-Solomon коды, overhead vs репликация, hot/cold данные
- Small Files в HDFS: нагрузка на NameNode heap и метаданные
- HDFS CLI: диагностика через hdfs dfs, fsck, dfsadmin и балансировка
- HDFS vs S3: когда что выбирать - latency, consistency, стоимость
- Hive Metastore: HMS как thrift-сервер метаданных, enableHiveSupport() и spark.catalog API
- Managed vs External Tables: lifecycle данных, LOCATION, DROP TABLE поведение и Data Lake паттерн
- Dynamic Partition Overwrite: partitionOverwriteMode=dynamic - идемпотентная перезапись разделов
- Форматы файлов в HDFS: ORC vs Parquet vs Avro, Snappy vs ZSTD vs Gzip - когда что выбирать
- Hive DDL: STORED AS, ROW FORMAT SERDE, MSCK REPAIR TABLE и StorageDescriptor в метастере
- Диагностика: чтение Spark UI - Jobs, Stages, Event Timeline, SQL plan
- Skew Detection: как увидеть skew в Task Duration и shuffle read
- AQE: включение, параметры и что оно реально делает
- AQE Skew Join: автоматическое обнаружение и разбиение hot partitions
- Data Skew: ручной salting - генерация ключей, двухпроходная агрегация
- Executor Sizing: 4-компонентная формула памяти и правило 5 cores
- Memory Fractions: spark.memory.fraction, storageFraction и tuning
- Broadcast Join: autoBroadcastJoinThreshold, hints, Task too large
- Sort-Merge Join vs Shuffle Hash Join: когда каждый применяется
- KryoSerializer: регистрация классов и измеримый выигрыш
- mapPartitions: connection pool паттерн для JDBC и HTTP клиентов
- Bloom Filter Join и Dataset vs RDD: узкоспециализированные случаи
- Кэширование: cache() vs persist(), Storage Levels и когда unpersist()
- Все 5 стратегий JOIN: BHJ/SMJ/SHJ/BNLJ/Cartesian и join hints
- Hardware tuning: spark.local.dir на NVMe и memoryOverhead для Python-воркеров
- Runtime Filtering: AQE Bloom Filter push к сканированию - dynamic filtering в star-schema joins
- Spark UI: вкладка Jobs - DAG визуализация и поиск узкого места
- Spark UI: вкладка Stages - task skew, shuffle write/read и GC time
- Spark UI: вкладка SQL - чтение physical plan и метрики операторов
- Spark UI: вкладка Executors - memory breakdown и task сводка
- History Server: настройка event log на S3/HDFS и retention политика
- Spark Metrics System: Sources, Sinks и namespace таксономия
- JMX + Prometheus: pull-based сбор метрик через jmx_exporter
- Grafana Dashboard: ключевые метрики Spark для production-алертинга
- SparkMeasure: programmatic сбор Stage и Task метрик в коде
- Алерты: SLO для job duration, failure rate и executor OOM
- Профилирование: чтение Java Stack Traces, PySpark Profiler и WholeStageCodegen в SQL tab
- Medallion Architecture: bronze, silver, gold - принципы и границы слоёв
- Idempotency: OVERWRITE, MERGE, dedup паттерны для safe retry
- Schema Evolution: стратегии добавления и удаления колонок без поломки
- Deduplication: dropDuplicates vs Window row_number - когда и что выбрать
- Anti-Join dedup (left_anti): загрузка инкремента без дублей - паттерн для batch ETL
- Hash-ключ для составных PK: sha2/xxhash64 - dedup по 10+ колонкам без shuffle penalty
- Bloom Filter в Spark SQL: probabilistic dedup и ускорение join на больших таблицах
- Airflow + Spark: SparkSubmitOperator vs KubernetesPodOperator
- DAG Design: декомпозиция job, зависимости и параллелизм в Airflow
- Testing PySpark: pyspark.testing.assertDataFrameEqual и Unit-тесты
- Testing: изоляция от внешних источников без mock-антипаттернов
- Quality Gates: встраивание row count, null check и drift-проверок
- Dead Letter Queue: паттерн Quarantine для плохих записей
- Observability: структурированное логирование и трассировка Spark job
- SCD Type 2: valid_from/valid_to паттерн через MERGE INTO и Iceberg row-level deletes
- Apache Iceberg: зачем нужен table format и что он добавляет к Parquet
- Iceberg Snapshot Model: дерево файлов метаданных и атомарные коммиты
- Manifest Files и Manifest List: как Iceberg планирует scan без листинга
- Hidden Partitioning: partition transforms без изменения SQL-запросов
- Partition Evolution: изменение схемы партиционирования без перезаписи данных
- Copy-on-Write: принцип работы, latency при write, быстрый read
- Merge-on-Read: delete files, position deletes и equality deletes
- MERGE INTO: upsert паттерн для CDC - структура запроса и настройки
- Time Travel: запросы к историческим снапшотам и откат таблицы
- Table Maintenance: rewrite_data_files, expire_snapshots, orphan files
- Delta Lake: transaction log и сравнение с Iceberg snapshot model
- Delta Lake: OPTIMIZE, VACUUM, Z-ordering и auto-optimize - production-операции для managed tables
- Выбор формата: Iceberg vs Delta Lake vs Hudi - матрица по кейсам
- JDBC Connector: базовое чтение, параметры подключения и JDBC URL
- Parallel JDBC Read: partitionColumn, lowerBound, upperBound, numPartitions
- Predicates: Custom Query и фильтрация на стороне PostgreSQL
- Connection Pool: ограничение числа соединений и fetchsize
- Write Modes: append, overwrite и staging-таблица для upsert
- PostgreSQL WAL: как работает Write-Ahead Log и logical replication
- Debezium CDC: replication slots, snapshot mode и структура Kafka-сообщений
- CDC Events: envelope формат, before/after поля и обработка операций
- ClickHouse MergeTree: физическая организация - part, chunk, granule
- ORDER BY как разреженный индекс: granules, mark files и query pruning
- ReplacingMergeTree: когда и как происходит дедупликация
- AggregatingMergeTree: материализованные агрегаты и partial state
- TTL: автоматическое удаление и перемещение данных по времени
- Projection: вторичный порядок сортировки для BI-запросов
- Spark → ClickHouse: JDBC vs clickhouse-spark-connector, batch tuning
- Партиционирование в ClickHouse vs Iceberg: разные роли и не путать
- Kafka Internals: producer routing, partition leader, ISR и acknowledgment
- Log Segment Lifecycle: retention.ms, retention.bytes и segment rolling
- Log Compaction: cleanup.policy=compact, tombstone и retention гарантии
- Structured Streaming: micro-batch модель - trigger, batch, commit
- Kafka Source: startingOffsets, failOnDataLoss, maxOffsetsPerTrigger
- Output Modes: append, update, complete - что работает с чем
- Checkpointing: WAL, offset commit и recovery при рестарте
- Watermarks: late data handling и state expiration
- Stateful Operations: mapGroupsWithState и управление state store
- Stateful Deduplication: dropDuplicates + withWatermark - ограничение размера state и TTL
- Kafka Sink: exactly-once через idempotent producer и транзакции
- RocksDB State Store: stateStore.providerClass, snapshotting и sizing для production-grade stateful streaming
- dbt Core: models, refs, sources и materialization types (table, view, incremental)
- dbt-spark адаптер: методы подключения thrift, http, session и profiles.yml
- dbt + Iceberg: file_format=iceberg, tblproperties и schema evolution
- dbt Incremental Models: append, insert_overwrite и merge с unique_key
- dbt macros для Spark: partition_by, clustered_by, location_root и submit_timeout
- dbt Tests и Source Freshness: generic tests, singular tests, data contracts
- dbt + Airflow: BashOperator dbt run vs SparkSubmitOperator, CI/CD pipeline
- dbt vs PySpark: матрица выбора по сложности логики и объёму данных
- Polars: Lazy API и query optimizer - как работает планирование
- Polars: streaming mode для датасетов больше RAM
- Polars: Arrow memory model и zero-copy обмен данными
- Polars vs pandas: многопоточность, SIMD и почему быстрее
- DuckDB: SQL поверх Parquet без сервера - установка и синтаксис
- DuckDB + Iceberg: чтение lakehouse локально без Spark
- Выбор инструмента: Polars vs DuckDB vs Spark - матрица по сценариям
- Great Expectations: Expectations, Suites и базовая проверка датасета
- Great Expectations: Checkpoint и интеграция в Airflow/Spark pipeline
- Soda Core: checks.yaml синтаксис и CLI-driven подход к тестированию
- Schema Contracts: контракты на уровне схемы и drift-детекция
- Reconciliation: подсчёт расхождений между источником и таргетом
- OpenLineage: автоматический lineage через Spark listener
- OpenLineage + Airflow + dbt: единый граф зависимостей данных
- Data Catalog: OpenMetadata vs DataHub - установка и базовая настройка
- MLflow Tracking: runs, metrics, params и artifact store
- MLflow Autologging: поддерживаемые фреймворки и что логируется автоматически
- Custom Metrics и Artifacts: логирование Spark DataFrames и кастомных графиков
- MLflow Model Registry: Register, Staging, Production жизненный цикл
- CI/CD для моделей: автоматический promote через GitHub Actions
- Spark для Feature Engineering: распределённые трансформации для ML
- Batch Inference: Spark + MLflow pyfunc для scoring больших датасетов
- Feature Store паттерны: offline store на Iceberg и online serving
- Spark on Kubernetes: SparkOperator, submit mode и cluster mode
- Driver и Executor Pods: resource requests, limits и node affinity
- Python Dependencies: conda-pack, venv-pack и PEX - упаковка окружений
- Docker Image для Spark: базовый image, кастомизация и registry
- Hive Metastore: PostgreSQL backend, схема, HA и failover
- Iceberg REST Catalog: установка, конфигурация и подключение из Spark
- Project Nessie: Git-like branching для Iceberg - установка в Docker, ветки и merge
- Service Accounts и Secrets: безопасный доступ к MinIO и PostgreSQL
- Network Policies: изоляция Spark namespace и egress правила
- History Server на K8s: event log в S3 и PersistentVolume
- Мониторинг платформы: Prometheus, Alertmanager и базовые алерты
- YARN: Resource Manager, Node Manager, client vs cluster mode и Dynamic Allocation
- CI/CD для Spark: GitHub Actions, pytest gate, venv-pack и автодеплой
- Проект 1: Batch Lakehouse - PostgreSQL → Bronze → Silver → Gold → ClickHouse
- Проект 1: Airflow DAG - идемпотентность, retry и quality gate
- Проект 2: CDC Pipeline - Debezium → Kafka → Spark Streaming → Iceberg MERGE
- Проект 3: Optimization Lab - воспроизвести skew, устранить через AQE + salting
- Проект 4: Production Runbook - мониторинг, алерты и incident response план
- Архитектурный разбор: code review решений и защита технических решений
- Проблема JVM: почему Tungsten не может использовать SIMD и при чём тут колоночное хранение
- SIMD и векторные регистры: AVX-512, как процессор обрабатывает Parquet-колонку пачкой за один такт
- Apache Arrow: in-memory columnar format, zero-copy передача и Record Batch как граница Python ↔ JVM ↔ C++
- Pandas UDF с Arrow: @pandas_udf, Series-батч вместо строки, скорость vs Row-at-a-time UDF
- Spark Plugin API: SparkPlugin и QueryStage интерфейс для замены executor
- Photon (Databricks): SIMD-движок на C++, что заменяет в Tungsten, 3–10× на Join и агрегациях
- Apache DataFusion Comet: архитектура Rust + Arrow, как подключается к Spark
- Comet: установка jar, конфигурация, поддерживаемые операторы, ограничения и запуск TPC-H бенчмарка
- Gluten + Velox: Substrait IR и offloading вычислений в C++ Velox
- RAPIDS cuDF: GPU acceleration, CUDA требования и профиль данных
- RAPIDS: когда GPU не даёт прироста - малые данные, IO-bound запросы
- Serverless Spark на K8s: EMR Serverless, Dataproc Serverless, pod templates и auto-scaling
- Ray + RayDP: Spark ETL → Ray Actor для distributed ML, сравнение с MLlib
- Выбор движка: Comet vs Velox vs RAPIDS vs Photon - матрица по кейсам
- Sail: Rust-нативный движок как полная замена JVM - Spark Connect, Arrow Flight, DataFusion и TPC-H 4× быстрее
- Блок 1 - Архитектура и основы: Driver/Executor, RDD, DAG, Catalyst, Tungsten, lifecycle, Data Locality
- Блок 2 - Оптимизация и производительность: skew, AQE, broadcast join, bucketing, CBO, spill, Z-Ordering, NotSerializableException
- Блок 3 - DataFrame API и SQL: joins, window functions, null handling, SCD Type 2, external tables
- Блок 4 - Python и расширенные возможности: UDF, Pandas UDF, Arrow, Structured Streaming, output modes
- Блок 5 - Практические кейсы: 10TB pipeline, ExecutorLostFailure, Pandas→PySpark, K8s, Data Quality
- Блок 6 - S3, HDFS и Apache Iceberg: S3A committers, CoW vs MoR, компакция, CDC, Iceberg vs Delta
- Блок 7 - Kafka, Streaming и мониторинг: иерархия Kafka, log compaction, triggers, Python зависимости, SparkMeasure, Airflow
- Блок 8 - Практические вопросы: 70 вопросов в формате Databricks-экзамена - кэширование, партиционирование, DataFrame API, join, UDF, I/O
- 101 упражнение: DataFrame из списков, фильтрация, агрегации, оконные функции, pivot/unpivot, нормализация, ML-пайплайны, UDF, работа с датами
- collect() и toPandas(): когда весь датасет едет в Driver и убивает JVM
- explode() misuse: N×M взрыв строк, вложенный explode и когда использовать flatten
- repartition() abuse: лишний shuffle и когда coalesce() не вариант
- Кэш без стратегии: persist до фильтрации, нет unpersist и cache в цикле
- UDF Pitfalls: захват closure, null без защиты, non-determinism и row-at-a-time налог
- Shuffle Explosion: COUNT DISTINCT на TB, multiple GROUP BY, sort без причины
- Schema Anti-Patterns: inferSchema в продакшне, from_json без схемы, triple-nested struct
- Cartesian Join Trap: accidental cross join, BNLJ на больших данных, non-equi joins
- Schema-first prompting: printSchema(), DDL и sample data как обязательный контекст для LLM - почему AI 'галлюцинирует' без схемы
- Dev Guides и System Prompt: cursorrules / CLAUDE.md со Spark-специфичными правилами (no collect(), prefer broadcast, no Python UDF без аппрувала)
- AI-readable репозиторий: README hierarchy, ADRs (Architecture Decision Records), каталог /examples с join_patterns / cdc_patterns / scd2_patterns
- Итеративный промптинг для Spark задач: разбивка ETL на шаги - чтение+фильтрация → трансформация → оптимизация (bucketing/partitioning) → запись
- AI-assisted testing: генерация pytest + chispa тестов, автоматическое написание Great Expectations rules по схеме таблицы
- AI для рефакторинга: legacy notebooks → production modules, SQL optimization, deduplication, SQL-first mindset почему декларативный код лучше для агентов
- Опасные зоны LLM в Spark: что AI делает плохо - skew, cluster economics, data semantics, hidden assumptions, необходимость human review
- Что такое OpenCode: terminal-native open-source coding agent, model-agnostic, MCP/skills/plugins, LSP integration, parallel sessions - почему интересен для DE
- Настройка среды: VSCode + OpenCode CLI + Dockerized Spark + Local Lakehouse (docker-compose: spark, postgres, minio, nessie/iceberg)
- Agent context - что дать агенту: architecture docs, spark_guidelines.md, coding standards, /examples, schemas, data contracts, ADRs
- Custom Skills для DE (часть 1): SCD2 pipeline generator, CDC pipeline generator, Delta Lake MERGE patterns, bronze ingestion шаблон
- Custom Skills для DE (часть 2): Spark optimization reviewer, SQL reviewer, Airflow DAG generator, data contract validator
- MCP ecosystem для DE: catalog (DataHub / OpenMetadata), lineage (OpenLineage), schema registry - агент читает lineage и схемы напрямую
- Agentic feature workflow: ticket → repo analysis → implementation plan → code generation → tests → validation → PR opening → human review
- Autonomous code review workflow: AI проверяет shuffle explosion, cartesian joins, hardcoded paths, отсутствие broadcast hints, UDF без необходимости
- Production agentic стек: solo/small team (OpenCode + Docker + GitHub Actions + Great Expectations) vs enterprise (MCP + Unity Catalog + Databricks + DataHub)
- de-agent-skills: готовый open-source репозиторий с pyspark_etl и spark_sql skills, enterprise specs и guides - подключаем как базу знаний агента
- Medallion Architecture: зачем три слоя, design-принципы каждого, anti-patterns - почему 'папки в S3 без метаслоя' уже legacy
- Bronze Layer: append-only immutable ingestion, обязательные audit-колонки (_ingest_ts, _source_system, _source_file, _event_id), почему дубликаты допустимы на этом слое
- Silver Layer: deduplication через ROW_NUMBER() + Window, нормализация схемы, бизнес-ключи, CDC handling, гарантии уникальности без PK
- Gold Layer: serving marts, агрегаты для BI, KPI-таблицы, SLA-гарантии качества, когда пушить в ClickHouse vs оставлять в Iceberg
- Dimensional Modeling (Kimball) на Spark: star schema, snowflake schema, факты и измерения, surrogate keys, grain таблицы фактов
- SCD (Slowly Changing Dimensions) в PySpark: Type 1 (overwrite), Type 2 (history rows) и Type 3 (prev/curr columns) - реализация через Iceberg MERGE INTO
- Data Vault 2.0 на Spark: Hubs, Links, Satellites - почему DV хорошо параллелится, surrogate keys через SHA-256 хеш бизнес-ключей, load patterns
- Data Vault + Medallion: Raw Vault → Business Vault (PIT-таблицы, Bridge-таблицы, вычисляемые сущности) → Data Marts с Kimball-моделью
- Constraints и дубликаты в Lakehouse: почему Iceberg не enforce PK, idempotent pipelines как architectural principle, MERGE INTO vs append + dedup, snapshot isolation для rollback
- Event-driven modeling: immutable events vs mutable rows, append logs как source of truth, CDC-семантика (sequence numbers, offsets, event ids), когда event-log лучше SCD
- Lambda-архитектура в Lakehouse: Batch Layer как source of truth, Speed Layer на Spark Structured Streaming, foreachBatch + MERGE INTO, exactly-once через Iceberg, Serving Layer без двойного счёта
- Kappa-архитектура на Spark Structured Streaming: единый стриминговый пайплайн для истории и real-time, Replay через startingOffsets=earliest, watermark и State Store, Blue-Green migration без downtime
- Lambda vs Kappa - критерии выбора: Replay Factor, TCO (S3 vs Kafka retention), паттерн Incremental Batch (trigger=availableNow), decision framework и три production-кейса
Self-host open-source стек
Никаких зарубежных SaaS-сервисов. Все технологии из лекций разворачиваются локально или on-prem.
8 практических проектов
Практические задания к лекциям связывают весь стек — от оптимизации JOIN до end-to-end CDC pipeline.
Ежедневный batch-пайплайн bronze→silver→gold с Iceberg, ACID-транзакциями и Airflow оркестрацией.
Потоковая репликация PostgreSQL через Debezium → Kafka → Spark Structured Streaming → Iceberg MERGE.
Воспроизвести skewed join, измерить baseline, устранить через salting + AQE. Сравнить Spark UI до/после.
Real-time агрегация событий с watermarks, stateful операциями и записью результата в ClickHouse.
Prometheus + Grafana dashboard для Spark: GC time, shuffle bytes, stage duration, executor utilization.
TPC-H бенчмарк: сравнить vanilla Spark, Apache Comet (2.2x) и RAPIDS GPU (9x+) на одном датасете.
Курс по PySpark для инженеров данных
Это лучший бесплатный курс по PySpark 2026 на русском языке — Free PySpark Certification Course, охватывающий путь от введения в PySpark до уровня Senior Data Engineer. 301 лекций разбиты на 27 модулей: от фундаментальных концепций Apache Spark до Catalyst Optimizer, Tungsten, Apache Iceberg, Kafka Structured Streaming и native execution engines.
Мастерство PySpark: от новичка до продвинутого уровня — именно так выстроена программа курса. В отличие от большинства tutorial'ов, здесь каждая лекция объясняет не только «как», но и «почему» — через internals движка, реальные примеры из Spark UI и практические задания. Весь стек — self-host open-source без зарубежных SaaS.