PySpark Roadmap
PySpark Roadmap Middle → Senior DE
Модули Стек Практика Источники CV
Начать →
Модули Стек Практика Источники Начать → CV ↗
⚡ Главная
00. Фундамент
Что такое Apache Spark: место в экосистеме Big Data Архитектура кластера: Driver, Executor и Cluster Manager Spark vs Hadoop MapReduce: in-memory вычисления и Lazy Evaluation RDD, DataFrame и Dataset: эволюция абстракций Spark Transformations vs Actions: анатомия ленивых вычислений Narrow vs Wide трансформации: shuffle изнутри SparkSession: точка входа, конфигурация и master URL spark-submit: упаковка зависимостей и запуск на кластере Локальная установка PySpark и первый job с нуля Глоссарий PySpark
00. START: Введение в PySpark
Введение в PySpark Настройка SparkSession Чтение и запись файлов в PySpark Чтение и запись CSV-файлов в PySpark Чтение и запись JSON-файлов в PySpark Обращение к колонкам в PySpark DataFrame Выборка колонок в PySpark 📋 Фильтрация данных в PySpark Группировка данных в PySpark Джоины в PySpark Пивот данных в PySpark Работа с NULL-значениями в PySpark Функции работы с датами и временем в PySpark Математические функции в PySpark Строковые функции в PySpark Оконные функции в PySpark Функции lead() и lag() в PySpark Оконные функции с rowsBetween в PySpark
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: UnsafeRow, off-heap память и Whole-Stage CodeGen Unified Memory Model: execution vs storage регионы и граница Spill на диск: причины, диагностика в Spark UI и как избежать GC Pressure: G1GC настройки, off-heap память и сигналы в логах Py4J Gateway и Python Workers: как Python-код переходит в JVM Data Locality: перемести код к данным, а не данные к коду Жизненный цикл Spark-приложения: от SparkContext до Block Manager 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: Spark SQL против Python UDF Date/Time Functions: от строк до таймзон 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 Контроль Ingestion: inferSchema, режимы чтения и изоляция брака Actions в деталях: foreachPartition, checkpoint для длинных DAG, toPandas safety Тестирование PySpark: от хаоса к надёжности Pandas API on Spark (pyspark.pandas): когда и как DataFrame Wrangling: filtering joins, set operations и wide↔long
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 и совместимость 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, Delta OPTIMIZE, Iceberg rewrite_data_files MinIO: установка, настройка bucket policy и интеграция со Spark Boto3: Python SDK для S3-операций - листинг, загрузка, Presigned URLs и управление метаданными 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: рентген вашего приложения Skew Detection: как увидеть skew в Task Duration и shuffle read AQE: Adaptive Query Execution - включение, параметры и что оно реально делает 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(), persist() и управление памятью Все 5 стратегий JOIN: физика алгоритмов и управление Catalyst 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 Профилирование PySpark: Java Stack Traces, PySpark Profiler и WholeStageCodegen
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: реализация через 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 YARN: Resource Manager, Node Manager и режимы запуска Spark CI/CD для Spark: тестирование, линтинг и автоматический деплой 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
Проблема 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: установка, операторы, ограничения и 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-нативный движок с полным Spark API
21. Подготовка к собеседованию
Блок 1 - Архитектура и основы Блок 2 - Оптимизация и производительность Блок 3 - DataFrame API и SQL Блок 4 - Python и расширенные возможности Блок 5 - Практические кейсы Блок 6 - S3, HDFS и Apache Iceberg Блок 7 - Kafka, Streaming и мониторинг Блок 8 - Практические вопросы 101 практическое упражнение по PySpark
22. Anti-Patterns: Как убить кластер за 5 минут
collect() и toPandas(): когда весь датасет едет в Driver explode() misuse: N×M взрыв строк и когда использовать flatten repartition() abuse и coalesce(): когда shuffle нужен, а когда нет Кэш без стратегии: persist до фильтрации, нет unpersist и cache в цикле UDF Pitfalls: closure, null без защиты и non-determinism Shuffle Explosion: COUNT DISTINCT, множественный GROUP BY и sort без причины Schema Anti-Patterns: inferSchema в продакшне, вложенные структуры и from_json без схемы Cartesian Join Trap: случайный cross join, BNLJ на TB и non-equi join ловушки
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 AI-агент для Data Engineering Настройка среды: VSCode + OpenCode + Dockerized Spark + Local Lakehouse Agent Context: как дать AI-агенту знания Senior Data Engineer Custom Skills для DE: SCD2, CDC, Bronze Ingestion и MERGE-паттерны Custom Skills для DE (часть 2): Spark Reviewer, SQL Reviewer, DAG Generator, 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: готовые PySpark и Spark SQL skills для AI-агента
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: дедупликация, CDC и гарантии уникальности без Primary Keys Gold Layer: serving marts, агрегаты для BI, KPI и дилемма Iceberg vs ClickHouse Dimensional Modeling (Kimball) на Spark: Star Schema, Snowflake, Grain, Surrogate Keys SCD в PySpark: Type 1, Type 2 и Type 3 через Iceberg MERGE INTO Data Vault 2.0 на Spark: Hubs, Links, Satellites и параллельный инжект Data Vault + Medallion: Raw Vault → Business Vault → Kimball Data Marts Constraints в Lakehouse: почему нет PK, идемпотентность, MERGE vs Dedup, Snapshot 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 и Real-time потоков на Spark Kappa-архитектура на Spark Structured Streaming: стриминг как единый источник правды Сравнение Lambda vs Kappa: критерии выбора, стоимостная модель и эволюция в Lakehouse
Главная › Kafka + Structured Streaming › Kafka Internals: producer routing, partition leader, ISR и acknowledgment

Kafka Internals: producer routing, partition leader, ISR и acknowledgment

Kafka Internals: producer routing, partition leader, ISR и acknowledgment

streaming

Урок в разработке - материал появится в ближайшее время.

Что будет в этом уроке¶

Kafka Internals: producer routing, partition leader, ISR и acknowledgment

Следующий → Log Segment Lifecycle: retention.ms, retention.bytes и segment rolling

На этой странице

  • Что будет в этом уроке
⚡ PySpark Advanced Roadmap Бесплатный курс по PySpark, Spark SQL, Apache Iceberg, Kafka, ClickHouse и dbt. Roadmap не включает зарубежные SaaS — S3 здесь совместимый API (MinIO, Ceph, Ozone).
Модули Self-host stack Практика Источники Sitemap