Created by AIImproved by people

Riqli · living documents · updated continuously

Where AI knowledge meets human practice.

Share with friends
image
Обработка и подготовка данных (ETL/ELT) Hard skill. (Data preprocessing. ETL. ELT. Data cleaning. Data transformation. Apache Spark. Hadoop. Self-study. Q&A. Tutorials. Documentation.)
sections

Введение

contents

Аннотация. Современный бизнес генерирует колоссальные объёмы данных, но их ценность раскрывается только после качественной обработки. Курс посвящён фундаментальным процессам обработки и подготовки данных, а именно методологиям ETL и ELT. Вы изучите полный цикл: от извлечения из разрозненных источников до загрузки в хранилища, а также освоите современные инструменты, включая Apache Spark и Hadoop. Эти знания помогут вам строить надёжные и масштабируемые конвейеры данных, устраняя главную боль инженеров — хаос и неконтролируемое качество информации.

contents

Цель курса. После прохождения курса вы сможете самостоятельно проектировать, реализовывать и оптимизировать процессы извлечения, трансформации и загрузки данных, выбирая между подходами ETL и ELT для решения конкретных бизнес-задач, а также эффективно использовать распределённые вычислительные системы для работы с большими данными.

contents

Результаты обучения.

  • Знать: архитектуру ETL/ELT-процессов, методы очистки и валидации данных, принципы работы с Apache Spark и Hadoop, критерии выбора инструментов.
  • Уметь: настраивать конвейеры данных, писать скрипты трансформации с использованием PySpark, применять паттерны обработки ошибок и оптимизировать производительность запросов.
  • Владеть: навыками профилирования данных, управления качеством данных, построения Data Lakehouse и автоматизации ETL-процессов с помощью современных фреймворков.
contents

Для кого этот курс.

Курс создан для инженеров данных, аналитиков и разработчиков, которые хотят систематизировать знания в области подготовки данных. Он будет полезен всем, кто участвует в построении DWH-решений или мигрирует на облачные платформы. Вы получите не только теорию, но и практические примеры с реальными инструментами.

Этот курс не для начинающих пользователей Excel или менеджеров без технического бэкграунда. Предполагается базовое понимание SQL, Python и концепций баз данных. Если вы ищете общее введение в Big Data без погружения в детали кода и архитектуры, этот курс может быть слишком технически насыщенным.

sections

Модуль 1. Основы ETL и ELT

contents

Определение ETL. Аббревиатура ETL расшифровывается как Extract, Transform, Load. Это классический подход к интеграции данных, при котором сначала происходит извлечение (Extract) сырых данных из различных источников (базы данных, файлы, API). Затем следует этап трансформации (Transform), где данные очищаются, обогащаются, агрегируются и переформатируются в соответствии с бизнес-правилами. И только после этого данные загружаются (Load) в финальное хранилище, часто в виде размерных таблиц. Основная философия ETL — подготовка данных до загрузки, что гарантирует высокое качество информации в целевом хранилище. Исторически ETL является доминирующим подходом в классических корпоративных DWH.

contents

Определение ELT. ELT (Extract, Load, Transform) — это современная парадигма, ставшая популярной с развитием облачных хранилищ и мощных вычислительных движков. В этом подходе данные сначала извлекаются и загружаются в целевую систему в сыром, необработанном виде. Трансформация происходит непосредственно внутри хранилища или вычислительной платформы (например, Snowflake, Google BigQuery, Apache Spark). Это позволяет использовать вычислительные мощности целевой системы для манипуляции огромными объёмами данных. Главное преимущество ELT — гибкость: вы загружаете все данные и решаете, как их трансформировать, уже когда возникнет бизнес-вопрос. Этот подход идеально подходит для Big Data и Data Lake архитектур.

contents

Сравнение ETL и ELT. Критерии выбора между ETL и ELT зависят от нескольких факторов. ETL лучше подходит для структурированных данных и строгих требований к качеству данных на этапе загрузки. Он позволяет сэкономить место в хранилище, так как загружаются только нужные данные. Однако ETL требует сложных предварительных моделей данных и негибок к изменениям. ELT, напротив, быстрее на этапе разработки, так как не требует сложной трансформации до загрузки. Он лучше масштабируется и позволяет использовать необработанные данные для различных целей. Основной недостаток ELT — высокие требования к вычислительным ресурсам хранилища, а также необходимость тщательно управлять стоимостью трансформаций.

contents

Когда использовать ETL. Выбирайте ETL, если у вас строгие нормативные требования к конфиденциальности и вы должны контролировать данные до того, как они попадут в общее хранилище. ETL предпочтителен при интеграции с устаревшими системами, где сложно выполнять тяжёлые трансформации. Если ваша бизнес-отчётность требует фиксированных, неизменных агрегатов и вы используете традиционные DWH на основе SQL, ETL будет надёжным выбором. Также ETL полезен, когда стоимость хранения в целевом хранилище высока, и вы стремитесь минимизировать объём загружаемых данных, очищая их на промежуточном уровне.

contents

Когда использовать ELT. ELT становится идеальным выбором в среде Cloud Data Warehouses (Snowflake, Redshift, BigQuery), где вычисления масштабируются независимо от хранения. Применяйте ELT, если вы работаете с большим количеством неструктурированных или полуструктурированных данных (JSON, лог-файлы). Этот подход обеспечивает максимальную гибкость: аналитики могут исследовать сырые данные без участия инженеров. Если вы строите Data Lakehouse и хотите избежать многократных ETL-пайплайнов для разных сценариев использования, ELT значительно упрощает архитектуру.

contents

История эволюции подходов. ETL-архитектура доминировала в 90-х и 2000-х годах с распространением реляционных DWH от Oracle и Teradata. Она была детерминированной и строгой. С появлением Apache Hadoop в 2006 году начался сдвиг: концепция Data Lake позволила хранить сырые данные дёшево. ELT получил второе дыхание с облачными платформами, где машины для вычислений (compute) отделены от хранилищ. Сегодня индустрия движется к гибридным моделям, но понимание эволюции критически важно для выбора правильной стратегии.

sections

Модуль 2. Архитектура конвейеров данных

contents

Источники данных. Первый шаг любого ETL/ELT процесса — извлечение данных из источников. К ним относятся реляционные СУБД (PostgreSQL, MySQL), NoSQL (MongoDB, Cassandra), файловые системы (CSV, Parquet, Avro), потоковые данные (Kafka, Kinesis) и внешние API (REST, SOAP). Выбор протокола подключения и метода извлечения зависит от объёма: для больших таблиц используют инкрементальную загрузку, для малых — полный снэпшот. Важно учитывать влияние на исходную систему: для тяжёлых запросов использовать реплики, чтобы не нагружать production.

contents

Способы извлечения данных. Существует два основных способа извлечения: логическое и физическое. Логическое извлечение подразумевает чтение данных через SQL-запросы или API. Физическое — чтение файлов логов транзакций (CDC — Change Data Capture) или копирование файлов баз данных. CDC является предпочтительным для репликации в реальном времени, например, с помощью инструментов Debezium, которые читают binlog MySQL. Периодическое извлечение (batch) проще в реализации и подходит для аналитических задач, не требующих секундной актуальности.

contents

Инкрементальная загрузка. Вместо того чтобы каждый раз выгружать все данные (полная загрузка), инкрементальная загрузка переносит только новые или изменённые записи. Это снижает нагрузку на сеть и источники. Технически инкремент реализуется через поля временных меток (updated_at), идентификаторы транзакций или механизмы CDC. Ключевой вызов — обработка «мягких» удалений и дубликатов. Рекомендуется всегда иметь стратегию upsert (обновление или вставка) в целевой таблице, чтобы сохранять историю изменений.

contents

Промежуточный слой (Staging Area). Staging Area — это временное хранилище, куда данные попадают сразу после извлечения. Это буферная зона между источниками и финальной трансформацией. Здесь данные могут быть просто скопированы в «сыром» виде для дальнейшей проверки. Staging критичен для отказоустойчивости: если процесс трансформации упал, вы не потеряете данные и сможете перезапустить конвейер с чистого листа. Часто Staging реализуется как отдельная схема в облачном DWH или как каталог в Data Lake.

contents

Слой трансформации (Transform Layer). Это «сердце» ETL/ELT, где происходят все преобразования. В ETL трансформации выполняются на отдельном сервере или в кластере Spark до загрузки. В ELT — это SQL-скрипты или dbt-модели внутри хранилища. Трансформации делятся на «легкие» (приведение типов, переименование) и «тяжелые» (агрегации, денормализация, аналитика). Современный подход — использовать DAG (Directed Acyclic Graph) для управления зависимостями между трансформациями, что позволяет параллельно выполнять независимые шаги.

contents

Целевые хранилища (Target Layer). Финальная загрузка направляет данные в целевое хранилище. Это может быть Data Warehouse (DWH) для структурированных отчётов, Data Lake для сырых данных, или Data Mart для конкретного бизнес-подразделения. В облачных архитектурах часто используется гибрид: холодные данные хранятся в S3 (Lake), горячие — в Redshift или Snowflake. В целевом слое должны быть настроены права доступа, жизненные циклы данных (партиционирование) и механизмы бэкапа.

contents

Оркестрация конвейеров. Оркестрация управляет последовательностью выполнения задач. Она отвечает за запуск по расписанию, обработку ошибок и повторные попытки. Популярные инструменты: Apache Airflow, Prefect, Dagster. Airflow, написанный на Python, позволяет описывать пайплайны как код (DAG). Настройка зависимостей (например, Task A зависит от Task B) и управление ресурсами — ключевые функции. Современная практика — использовать Kubernetes для масштабирования воркеров оркестратора.

contents

Мониторинг и логирование. Без мониторинга ETL/ELT процессы превращаются в «чёрный ящик». Необходимо логировать время выполнения, объём обработанных строк, статусы шагов и ошибки. Инструменты вроде Datadog или Prometheus + Grafana позволяют визуализировать метрики. Рекомендуется настроить алерты на отклонения: например, если загрузка данных длится на 50% дольше обычного или количество записей упало ниже порога. Подробное логирование помогает быстро локализовать проблему.

sections

Модуль 3. Трансформация данных

contents

Очистка данных (Data Cleaning). Это фундаментальный этап подготовки данных. Он включает удаление дубликатов, исправление опечаток, обработку выбросов и стандартизацию форматов (дат, валют). Например, в PySpark для удаления дубликатов используется dropDuplicates(). Для обработки пропусков применяют стратегии: удаление строк (dropna()), заполнение средним/медианой или использование алгоритмов интерполяции. Ошибка на этом этапе приводит к эффекту «мусор на входе — мусор на выходе» (GIGO).

contents

Нормализация и денормализация. Нормализация — это процесс организации данных для устранения избыточности (обычно до 3-й нормальной формы). Используется в OLTP-системах. Денормализация, наоборот, объединяет таблицы для ускорения чтения, что характерно для DWH (звёздные схемы). В трансформациях вы часто переходите от нормализованных источников к денормализованным витринам. Применяйте денормализацию обдуманно: она ускоряет SELECT, но замедляет UPDATE и увеличивает место.

contents

Агрегация данных. Агрегация — это суммирование или группировка данных для получения итоговых метрик. Примеры: вычисление ежедневных продаж, среднего чека, суммы выручки по категориям. В SQL это GROUP BY, в Spark — groupBy().agg(sum(...)). При агрегации важно правильно выбрать ключи группировки и учитывать иерархии (Roll-up, Drill-down). Создание предварительных агрегатов (summaries) — ключевой приём для ускорения отчётов в BI-системах.

contents

Преобразование типов данных. Корректное приведение типов критично для производительности и точности. Частая проблема: даты хранятся как строки. Их необходимо конвертировать в тип Timestamp: totimestamp(column,yyyyMMdd)to_timestamp(column, 'yyyy-MM-dd'). Также преобразуйте строки в числа для арифметических операций. В PySpark используйте col.cast('double'). Ошибки типа данных — одна из главных причин падения пайплайнов, поэтому валидация на этапе Staging помогает их отловить заранее.

contents

Обогащение данных (Data Enrichment). Обогащение добавляет недостающую информацию из внешних источников. Например, к транзакциям вы добавляете географические координаты по IP или данные о курсах валют. Это улучшает качество аналитики. В пайплайнах обогащение часто выполняется через JOIN с референсными данными. Технический вызов — версионность справочников: если курсы меняются ежедневно, вы должны использовать снэпшоты на момент транзакции.

contents

Применение бизнес-логики. Самый сложный аспект — внедрение предметных правил. Пример: расчёт зарплаты с учётом коэффициентов, KPI или скидок. Бизнес-логика часто меняется, поэтому важно отделять её от инфраструктурного кода. Рекомендуется использовать инструменты типа dbt для описания логики на SQL с использованием макросов, а в ETL-пайплайнах выделять эту логику в отдельные микросервисы или функции.

sections

Модуль 4. Инструменты и технологии

contents

Apache Spark. Это распределённый движок для обработки больших данных, являющийся стандартом де-факто для ETL. Spark позволяет выполнять трансформации в памяти, что значительно быстрее Hadoop MapReduce. Основной интерфейс — DataFrame API, напоминающий панды. Spark использует DAG (направленный ациклический граф) для оптимизации запросов. Пример кода:

df = spark.read.parquet('s3://data') processed = df.filter(df['age'] > 18).groupBy('city').count() processed.write.parquet('s3://output')

contents

Hadoop и HDFS. Hadoop — это экосистема, включающая HDFS (распределённую файловую систему) и YARN (менеджер ресурсов). Раньше Hadoop был синонимом Big Data. Сейчас он часто используется как распределённое хранилище, а Spark выступает вычислительным слоем. HDFS разбивает файлы на блоки (128 МБ) и реплицирует их на узлы кластера, обеспечивая отказоустойчивость. Важно знать команды HDFS: hdfs dfs -put, hdfs dfs -get для управления файлами.

contents

PySpark и Python. PySpark — это API для работы со Spark из Python. Он позволяет использовать преимущества Python (простоту, библиотеки) с мощью распределённых вычислений. Для высоконагруженных ETL-задач PySpark часто предпочтительнее Scala из-за скорости разработки. Однако важно помнить о стоимости: в PySpark сериализация Python объектов может быть медленнее, чем нативные JVM-коллекции. Используйте встроенные методы Spark вместо пользовательских UDF, когда это возможно, для повышения производительности.

contents

dbt (Data Build Tool). dbt — это инструмент для трансформации данных, работающий по принципу ELT. Он не извлекает и не загружает данные, а только трансформирует их прямо в облачном DWH (Snowflake, BigQuery) с помощью SQL. dbt управляет зависимостями между моделями, тестирует данные (уникальность, непустые значения) и генерирует документацию. Это мощное решение для Data Engineer, позволяющее версионировать трансформации в Git и применять практики CI/CD.

contents

Apache Airflow. Как уже упоминалось, Airflow — ведущий инструмент оркестрации. Пайплайн в Airflow — это DAG (набор задач). Пример задачи на Python:

from airflow.operators.python_operator import PythonOperator def extract_data(): ... extract_task = PythonOperator(task_id='extract', python_callable=extract_data, dag=dag)
Airflow позволяет планировать задачи (например, ежедневно в 3:00) и имеет веб-интерфейс для мониторинга. Ключевая концепция — «SLA» (Service Level Agreement) для гарантии времени выполнения.

contents

Стриминговые платформы (Kafka). Apache Kafka используется для построения потоковых конвейеров данных (real-time ETL). Данные публикуются в топики, откуда их читают потребители. В отличие от пакетной обработки, потоковая обеспечивает задержку в миллисекунды. Для обработки потоков в Spark используется Spark Streaming (или Structured Streaming), который разбивает поток на микро-батчи. Kafka гарантирует доставку сообщений (at-least-once, exactly-once), что критично для конвейеров с финансовыми данными.

contents

Облачные сервисы. Крупные провайдеры предлагают managed-решения: AWS Glue (серверный Spark), GCP Dataflow (Apache Beam), Azure Data Factory. Они избавляют от администрирования кластеров. AWS Glue интегрируется с S3 и Redshift, имеет встроенный каталог данных (Glue Data Catalog). Выбор облачного сервиса ускоряет развёртывание ETL-инфраструктуры, но требует понимания моделей ценообразования (например, стоимость за vCPU-час или сканируемый объём данных).

sections

Модуль 5. Качество и управление данными

contents

Профилирование данных (Data Profiling). Профилирование — это анализ структуры, содержания и качества данных. Вы проверяете типы данных, распределение значений, количество пропусков, уникальность ключей. Инструменты (например, Great Expectations, Pandas Profiling) помогают автоматизировать этот процесс. Профилирование должно выполняться до начала активной трансформации. Выводы профилирования влияют на выбор методов очистки (например, если в колонке 90% пропусков, её, скорее всего, стоит исключить).

contents

Тестирование качества данных. Качество данных — это не опция, а обязательство. Используйте фреймворки для тестирования: Great Expectations (GX) или dbt tests. Тесты бывают: проверка на пустоту (not_null), уникальность (unique), соответствие шаблону (регулярное выражение для email). В GX тесты описываются как «Expectations». Если тест провален, пайплайн может остановиться (fail fast) или отправить алерт. Это предотвращает попадание некачественных данных в витрины.

contents

Управление метаданными. Метаданные — это данные о данных. Они включают информацию о происхождении (lineage), бизнес-глоссарий и технические свойства (тип, длина). Хорошее управление метаданными упрощает навигацию по хранилищу. Инструменты: Apache Atlas, DataHub (LinkedIn). Ведение каталога данных — критично для крупных организаций, чтобы аналитики могли находить и интерпретировать наборы данных, не спрашивая инженеров.

contents

Data Lineage (Происхождение данных). Lineage показывает путь данных от источника до потребления. Это «карта» потока данных, которая помогает оценить влияние изменений. Если вы меняете колонку в исходной таблице, lineage покажет, какие витрины, отчёты или модели dbt будут затронуты. Ручное отслеживание сложно; современные платформы (например, Informatica, Marquez) собирают lineage автоматически через парсинг логов или SQL-запросов.

contents

Обработка ошибок и восстановление. ETL-пайплайны должны быть устойчивы к сбоям. Стратегии включают: Idempotency (повторный запуск не создаёт дубликатов), использование транзакций (атомарность загрузки), и мёртвые очереди (dead letter queues) для «битых» записей. Важно проектировать конвейер так, чтобы он мог быть перезапущен с момента сбоя, а не сначала. Это достигается за счёт сохранения «чекпоинтов» (checkpoints) в потоковых системах.

contents

Безопасность данных. При извлечении, трансформации и загрузке данные должны быть защищены. Используйте шифрование в покое (S3 SSE) и в движении (TLS). Управляйте доступом на основе ролей (RBAC). В облаках применяйте IAM-роли. Также важно маскировать чувствительные данные (PII) на этапе трансформации: заменять реальные имена на хешированные значения или использовать псевдонимизацию. Это особенно актуально для GDPR и HIPAA.

sections

Модуль 6. Оптимизация производительности

contents

Партиционирование данных. Партиционирование — это разделение таблицы на части по ключу (например, по дате). В Spark это ускоряет чтение, так как запрос читает только нужные партиции, а не всю таблицу. При загрузке используйте partitionBy('date'). Однако слишком мелкое партиционирование (тысячи папок) создаёт оверхед. Оптимально выбирать ключи с высокой кардинальностью и частым использованием в фильтрах WHERE.

contents

Сортировка и бакетирование. Сортировка внутри партиций позволяет использовать оптимизации пропуска данных. Бакетирование (bucket) — это хэширование данных по ключу для улучшения джойнов (bucket join). В Hive/Spark: CLUSTER BY (id) INTO 16 BUCKETS. Правильная настройка бакетов минимизирует перемешивание данных (shuffle) при объединении больших таблиц. Это одна из самых эффективных оптимизаций для сложных трансформаций.

contents

Настройка Spark. Оптимизация Spark конфигураций: spark.sql.shuffle.partitions влияет на параллелизм при агрегациях. Установите его в 23×количествоядер2-3 × количество ядер. spark.executor.memory и spark.executor.cores подбираются под ресурсы кластера. Используйте Broadcast Join для маленьких таблиц, кэшируйте часто используемые промежуточные результаты. Следите за «узким горлышком» — большими перемешиваниями (shuffle).

contents

Оптимизация SQL запросов. В контексте ELT (где трансформации на SQL) оптимизация включает: использование индексов, избегание SELECT *, применение оконных функций вместо коррелированных подзапросов. В облачных DWH важно знать особенности ворклоадов: например, в Snowflake используйте «clustering keys», в BigQuery — «partition and cluster». Анализ плана выполнения (EXPLAIN) обязателен перед запуском тяжёлых трансформаций.

contents

Масштабирование кластеров. Автоматическое масштабирование позволяет адаптировать ресурсы под нагрузку. В облаках это автоскейлинг (AWS Auto Scaling). Для ETL с неравномерной нагрузкой (например, высокая в начале месяца) — это экономит бюджет. Важно настроить политику масштабирования на основе метрик очереди или загрузки CPU, чтобы избежать перерасхода средств в пиковые часы.

contents

Управление памятью и GC. В Spark управляйте памятью: разделите её на Storage (кэш) и Execution (shuffle). Если Execution не хватает, происходит запись на диск, что резко замедляет работу. Для минимизации GC (Garbage Collection) используйте сериализацию Kryo: conf.set('spark.serializer', 'org.apache.spark.serializer.KryoSerializer'). Для больших строковых данных это даёт прирост скорости до 20%.