Лучший курс по PySpark 2026 · Middle → Senior Data Engineer

Введение в PySpark и Spark SQL —
от архитектуры до production

Мастерство PySpark: от новичка до продвинутого уровня. Free PySpark course на русском — 301 лекций, Catalyst, AQE, Iceberg, Kafka, ClickHouse и self-host платформа без SaaS-привязки.

Бесплатно · На русском · Apache Spark 4.x · 2026

Целевая архитектура batch + streaming + serving
SOURCES PostgreSQL Transactional DB · WAL Apache Kafka Event streams · topics S3 / HDFS Object store · raw files INGESTION Debezium CDC · WAL replication Kafka Connect · Structured Streaming APACHE SPARK 4.x — UNIFIED PROCESSING ENGINE PySpark DataFrame · Dataset API pandas UDF · Arrow Spark SQL Catalyst · AQE · DPP CBO · Tungsten codegen Struct. Streaming micro-batch · continuous watermarks · stateful APACHE ICEBERG v2 · LAKEHOUSE · ACID SNAPSHOTS · TIME TRAVEL Bronze raw · append-only Silver cleansed · validated Gold aggregated · business-ready SERVING ClickHouse OLAP · BI · sub-second dbt + Spark SQL models · incremental Grafana + Spark UI metrics · observability
301
лекций — от введения
до продвинутого уровня
8
практических
проектов
Free
PySpark certification course
без оплаты и регистрации
4.x
Apache Spark 2026
актуальная версия 4.x.x
Написано уроков Курс в активной разработке · уроки выходят каждую неделю
218 / 301
72%
Программа курса

21 этап обучения — от новичка до продвинутого уровня

Лучший бесплатный PySpark tutorial на русском: от введения в PySpark и Spark internals до оптимизации, Lakehouse и стриминга. Каждый этап — отдельный модуль с лекциями и практикой.

PySpark Roadmap — 21 этап обучения
Готов
00
Фундамент
  • Что такое 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
Готов
00
START: Введение в PySpark
  • Введение в 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: скользящие суммы и средние, управление диапазоном окна
Готов
01
Spark Internals
  • 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 и архитектура протокола
Готов
02
Продвинутый PySpark API
  • 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
Готов
03
Spark SQL Advanced
  • 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-специфика
Готов
04
Partitioning
  • 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
Готов
05
S3-compatible Storage
  • 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
Готов
06
HDFS
  • 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 в метастере
Готов
07
Optimization
  • Диагностика: чтение 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
08
Monitoring
  • 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
Готов
09
Batch + Streaming Patterns
  • 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
Готов
10
Open Table Formats
  • 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 - матрица по кейсам
11
PostgreSQL
  • 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 поля и обработка операций
12
ClickHouse
  • 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: разные роли и не путать
13
Kafka + Structured Streaming
  • 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
Готов
14
dbt + Spark
  • 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: матрица выбора по сложности логики и объёму данных
15
Polars + DuckDB
  • 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 - матрица по сценариям
16
Data Quality
  • 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 - установка и базовая настройка
17
MLflow + MLOps
  • 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
18
Self-host Platform
  • 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 и автодеплой
19
Capstone Projects
  • Проект 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 решений и защита технических решений
Готов
20
Native Execution Engines NEW
  • Проблема 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× быстрее
Готов
21
Подготовка к собеседованию
  • Блок 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, работа с датами
22
Anti-Patterns: Как убить кластер за 5 минут NEW
  • 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
23
AI-Assisted PySpark Development
  • 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
24
OpenCode & Agentic Workflows
  • Что такое 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 - подключаем как базу знаний агента
Готов
25
Data Modeling в Lakehouse
  • 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.

⚙️
Compute
Apache Spark 4.x, PySpark, Spark SQL, Spark Connect, YARN / Kubernetes
💾
Storage
MinIO, Ceph RGW, Apache Ozone, HDFS — S3-compatible API
🏔️
Lakehouse
Apache Iceberg v2, Delta Lake 3.x (UniForm), Apache Hudi, Hive Metastore, Iceberg REST Catalog, Project Nessie
🗄️
Databases
PostgreSQL, ClickHouse MergeTree
📡
Streaming
Apache Kafka (KRaft, без ZooKeeper), Debezium CDC, Spark Structured Streaming
🔄
Orchestration & Transform
Apache Airflow, dbt Core + dbt-spark, KubernetesPodOperator, SparkSubmitOperator
🔍
Quality & Catalog
Great Expectations, Soda Core, OpenLineage, OpenMetadata / DataHub, Unity Catalog OSS
📊
Observability
Prometheus, Grafana, Spark History Server, SparkMeasure, Alertmanager
In-process Analytics
Polars, DuckDB — single-node альтернативы Spark для малых объёмов
🧪
ML Platform
MLflow Tracking, MLflow Model Registry, Feature Store на Iceberg
🚀
Native Engines
Apache Comet, Gluten + Velox, RAPIDS cuDF — ускорение без смены кода
Практика курса

8 практических проектов

Практические задания к лекциям связывают весь стек — от оптимизации JOIN до end-to-end CDC pipeline.

Intermediate
Batch Lakehouse Pipeline

Ежедневный batch-пайплайн bronze→silver→gold с Iceberg, ACID-транзакциями и Airflow оркестрацией.

Intermediate
CDC → Iceberg с Debezium

Потоковая репликация PostgreSQL через Debezium → Kafka → Spark Structured Streaming → Iceberg MERGE.

Advanced
Skew Join Optimization Lab

Воспроизвести skewed join, измерить baseline, устранить через salting + AQE. Сравнить Spark UI до/после.

Advanced
Kafka Streaming Analytics

Real-time агрегация событий с watermarks, stateful операциями и записью результата в ClickHouse.

Intermediate
Monitoring Stack Setup

Prometheus + Grafana dashboard для Spark: GC time, shuffle bytes, stage duration, executor utilization.

Advanced
Native Engine Benchmark

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.

Курс по PySpark для инженеров данных
Материалы

Ключевые источники

Book Learning Spark, 2nd ed. — Damji et al. (O'Reilly)
Book Spark: The Definitive Guide — Chambers & Zaharia
Docs Apache Spark Official Documentation — spark.apache.org
Docs Apache Iceberg Documentation — iceberg.apache.org
Book Fundamentals of Data Engineering — Reis & Housley
Course DataTalks.Club DE Zoomcamp — бесплатный онлайн-курс
Blog Databricks Engineering Blog — advanced Spark internals
Blog Martin Kleppmann — Designing Data-Intensive Applications
Paper Lakehouse: A New Generation of Open Platforms (CIDR 2021)
Paper Apache Iceberg: An Architectural Look Under the Covers