Google PDE: Spark, Dataproc и распределенная обработка данных — Руководство по подготовке
Часть Google Professional Data Engineer — Руководство по подготовке. Практикуйтесь с проверенными ответами в центре экзаменов Google, или пройдите тесты на время на ExamRoll.io.
Обзор
Apache Spark в Google Cloud Dataproc предоставляет управляемую эластичную платформу для распределенной обработки данных. В зависимости от требований к контролю, изменчивости времени выполнения и накладных расходов на управление вы можете выбирать между долгоживущими или эфемерными кластерами Dataproc, а также Dataproc Serverless для Spark. Spark предлагает устойчивые абстракции (RDD), реляционные API (DataFrames и Spark SQL) и отказоустойчивый механизм выполнения DAG, оптимизированный для итеративных и пакетных ETL-процессов в большом масштабе. В Google Cloud сервис Cloud Storage заменяет HDFS, обеспечивая долговечное и недорогое хранилище; коннектор BigQuery позволяет напрямую выгружать данные для анализа; а Dataproc Metastore централизует управление схемами. Эффективные решения согласовывают жизненные циклы хранилища и вычислительных ресурсов, настраивают Spark под рабочую нагрузку, внедряют средства наблюдаемости и применяют безопасность с принципом наименьших привилегий и сетевой изоляцией.
Архитектура Dataproc: кластеры, Serverless, хранилище и Metastore
- Типы кластеров и роли узлов
- Основные (master) узлы размещают YARN, HDFS NameNode (если используется) и UI драйверов Spark; режим высокой доступности (HA) использует несколько основных узлов.
- Рабочие узлы запускают исполнители (executors) и HDFS DataNodes (если используется).
- Вторичные/вспомогательные рабочие узлы обычно являются прерываемыми/spot-узлами для обеспечения эластичной и более дешевой емкости без ролей HDFS.
- Образы объединяют версии ОС и компонентов (например, 2.1-debian11, 2.2-ubuntu20); закрепляйте версии образов, чтобы контролировать совместимость Spark/Hadoop и выполнять обновления осознанно.
- Component Gateway публикует пользовательские интерфейсы (Spark History Server, YARN RM) безопасно через HTTPS.
- Dataproc Serverless для Spark
- Отсутствие необходимости в развертывании кластера, автоматическое автомасштабирование и посекундная тарификация для исполнителей и драйверов. Идеально подходит для спорадических или пиковых заданий, а также для минимизации операционных накладных расходов.
- Компромиссы: меньше низкоуровневых настроек, чем в кластерах; задержка при запуске задания может быть выше, чем у «прогретых» кластеров; используйте бессерверные метрики и журналы событий для устранения неполадок.
- Автомасштабирование
- Политики автомасштабирования кластера добавляют/удаляют рабочие узлы на основе метрик YARN/Spark и периодов охлаждения (cooldowns), позволяя отдельно настраивать группы основных и вторичных рабочих узлов.
- В режиме Serverless автомасштабирование управляется сервисом; проектируйте задания так, чтобы они были параллельны по партициям и избегали узких мест, связанных с сериализацией, для наилучшего масштабирования.
- Хранилище и коннекторы
- Предпочитайте использовать Google Cloud Storage (GCS) в качестве основной системы хранения данных; это отделяет вычислительные ресурсы от хранилища, снижает затраты на постоянные диски и позволяет данным сохраняться после удаления кластеров.
- Коннектор GCS (gs://) интегрируется с Hadoop/Spark. Запись в объектные хранилища использует протоколы фиксации (commit protocols); установите алгоритм FileOutputCommitter версии 2, чтобы уменьшить накладные расходы на переименование и ускорить фиксацию заданий в GCS:
--conf mapreduce.fileoutputcommitter.algorithm.version=2
```
- Используйте форматы **Parquet/ORC** с отсечением столбцов (column pruning) и проталкиванием предикатов (predicate pushdown). Управляйте мелкими файлами с помощью уплотнения (compaction), стремясь к размеру файла 128–512 МиБ для эффективного сканирования.
- Hive metastore
- Централизуйте схемы и метаданные таблиц в **Dataproc Metastore** (управляемый Apache Hive Metastore) или в хранилище метаданных на базе Cloud SQL для совместного использования каталогов между кластерами.
- Используйте **внешние таблицы**, указывающие на GCS, для обеспечения долговечности данных; партиционируйте по дате/часу, чтобы ограничить стоимость сканирования.
- Задания, инициализация и рабочие процессы
- Отправляйте задания `spark`, `pyspark`, `spark-sql` или `hadoop`. **Скрипты инициализации** устанавливают дополнительные библиотеки или агенты при создании кластера (например, коннекторы, библиотеки Python).
- **Шаблоны рабочих процессов (Workflow templates)** параметризуют многоэтапные конвейеры; они могут создавать эфемерные кластеры для каждого рабочего процесса, а затем удалять их. Это улучшает изоляцию и сокращает затраты на простой.
- **Эфемерные кластеры** рекомендуются для пакетных ETL-процессов; данные и хранилище метаданных находятся вне кластера (GCS, Dataproc Metastore, BigQuery).
- Интеграция с BigQuery
- **Коннектор Spark для BigQuery** читает и пишет данные напрямую в BigQuery; рассмотрите использование BigQuery Storage Read API для высокой пропускной способности и Write API для потоковой вставки с низкой задержкой и семантикой «ровно один раз» (exactly-once).
- Для обслуживания таблиц выполняйте последующие операции **MERGE** или перезапись партиций в BigQuery, чтобы атомарно завершить загрузку.
### Модель Spark, настройка производительности и надежность
- API и выполнение
- RDD: низкоуровневые, неизменяемые, типобезопасные в Scala/Java; вы контролируете партиционирование и персистентность.
- DataFrames/Datasets: реляционные, оптимизированные с помощью Catalyst; предпочтительны для ETL благодаря оптимизации запросов и генерации кода.
- Трансформации ленивые (map, filter, join); действия запускают выполнение (count, collect, save). Spark строит DAG из стадий, разделенных перемешиваниями (shuffles); задачи выполняются для каждой партиции.
- Партиционирование и перемешивание (shuffle)
- Входное партиционирование: достаточное количество партиций для утилизации всех ядер; начните с 2–4x от общего числа ядер исполнителей (executor cores). Управляется через spark.default.parallelism (для RDD) и опции чтения (для DataFrames).
- Партиции для перемешивания (shuffle partitions): значение по умолчанию 200 часто бывает недостаточным или избыточным. Настройте:
--conf spark.sql.shuffle.partitions= {total_executor_cores * 2 to 3}
```
- Целевой размер партиции после широких трансформаций — ~100–256 МиБ; слишком маленький размер вызывает накладные расходы планировщика, а слишком большой — риски OOM у исполнителя (executor).
- Перемешивание (shuffle) — основная статья расходов для операций joins, groupBy и orderBy. Обеспечьте достаточный объем памяти и дискового пространства для исполнителей; для кластеров с интенсивным перемешиванием рассмотрите использование локальных SSD.
- Перекос (skew) и стратегия соединений (join)
- Обнаружение перекоса (задачи с долгим временем выполнения, большие размеры партиций). Способы устранения:
- Широковещательная рассылка (broadcast) маленьких таблиц, чтобы избежать перемешиваний:
- Обнаружение перекоса (задачи с долгим временем выполнения, большие размеры партиций). Способы устранения:
--conf spark.sql.autoBroadcastJoinThreshold=64m
```
- «Подсаливание» ключей для «горячих» партиций; применяйте предварительную агрегацию на стороне map; выполняйте фильтрацию как можно раньше.
- Включите Adaptive Query Execution (AQE) для объединения партиций после перемешивания и обработки перекошенных соединений (skewed joins):
--conf spark.sql.adaptive.enabled=true
```
- Кэширование, контрольные точки и происхождение данных (lineage)
- Кэшируйте «горячие» промежуточные DataFrames экономно, только при повторном использовании; предпочитайте MEMORY_AND_DISK, чтобы избежать OOM.
- Сохраняйте контрольные точки (checkpoint) для длинных цепочек зависимостей в GCS или HDFS, чтобы ограничить пересчет в случае сбоев.
- Исполнители (executors) и динамическое выделение ресурсов
- Правильно подбирайте размер исполнителей, чтобы сбалансировать параллелизм и накладные расходы на сборку мусора (GC):
- Ядер на исполнителя: 2–5 для сбалансированных задач I/O и CPU; меньшее количество ядер сокращает паузы на GC.
- Накладные расходы на память: установите spark.yarn.executor.memoryOverhead для широких перемешиваний.
- Включите динамическое выделение ресурсов (dynamic allocation) с внешним сервисом перемешивания (external shuffle service) на кластерах, чтобы масштабировать исполнителей в соответствии с нагрузкой:
- Правильно подбирайте размер исполнителей, чтобы сбалансировать параллелизм и накладные расходы на сборку мусора (GC):
--conf spark.dynamicAllocation.enabled=true
--conf spark.shuffle.service.enabled=true
--conf spark.dynamicAllocation.minExecutors=0
--conf spark.dynamicAllocation.maxExecutors=200
```
- Паттерны отказоустойчивости для пакетных ETL-задач
- Идемпотентные записи: записывайте во временный/промежуточный (staging) путь, затем атомарно повышайте с помощью коммита на уровне каталога; для BigQuery записывайте в промежуточную таблицу и используйте MERGE:
MERGE target t USING staging s
ON t.id = s.id
WHEN MATCHED THEN UPDATE SET ...
WHEN NOT MATCHED THEN INSERT (...)
```
- Инкрементальная обработка: используйте фильтрацию на основе водяных знаков (watermark) по партициям ingestion_date; ведите манифест обработанных данных в GCS, чтобы избежать повторной обработки.
- Обработка «мертвых писем» (dead-letter): при ошибках парсинга/валидации направляйте некорректные записи на карантинный путь/таблицу с диагностикой. Для строгого контроля схемы и встроенных DLQ рассмотрите использование Dataflow; в Spark реализуйте try/catch для каждой записи и отдельный приемник (sink).
Безопасность, наблюдаемость и затраты
- Идентификация и доступ
- Запускайте кластеры и задания от имени выделенных сервисных аккаунтов с минимальными привилегиями IAM. Предоставляйте только необходимые роли, например:
- roles/dataproc.worker для сервисных аккаунтов инстансов
- roles/storage.objectViewer или objectAdmin для путей ввода-вывода в GCS
- roles/bigquery.dataEditor для целевых наборов данных
- Для Dataproc Serverless используйте отдельные сервисные аккаунты для каждого задания, чтобы ограничить область доступа.
- Запускайте кластеры и задания от имени выделенных сервисных аккаунтов с минимальными привилегиями IAM. Предоставляйте только необходимые роли, например:
- Сетевая изоляция и шифрование
- Используйте кластеры с частными IP-адресами в подсети VPC, ограничивайте доступ к UI мастер-узлов с помощью правил брандмауэра и включите Private Google Access для доступа к GCS/BigQuery без исходящего трафика в интернет.
- Размещайте кластеры в проектах Shared VPC для централизованного управления. При необходимости включите Kerberos в Dataproc для аутентификации внутри кластера.
- Шифруйте неактивные данные с помощью CMEK: настройте CMEK для бакетов GCS, Persistent Disks, Dataproc Metastore и BigQuery; по умолчанию для передачи данных используется TLS.
- Логирование, история и метрики
- Включите запись логов событий Spark в GCS и разверните History Server:
--conf spark.eventLog.enabled=true
--conf spark.eventLog.dir=gs://bucket/spark-events/
```
- Dataproc передает логи драйвера и YARN в Cloud Logging; экспортируйте их в приемники (sinks) для хранения и анализа инцидентов.
- Выполняйте мониторинг с помощью метрик Cloud Monitoring: ожидающие контейнеры YARN, CPU, память, состояние HDFS (если используется), пропускная способность GCS. Настройте оповещения о длительных повторных попытках этапов (stage retries), потере исполнителей (executor loss) и всплесках спекулятивного выполнения (speculative execution).
- Анализ сбоев: распространенные причины включают отстающие задачи (stragglers) из-за перекоса данных (skew), нехватку памяти (OOM) у исполнителей во время shuffle, сбои при фиксации данных в объектном хранилище и потерю прерываемых/spot-узлов. Увеличивайте количество повторных попыток обдуманно; чрезмерные повторы могут увеличить затраты и задержки.
- Оптимизация затрат
- Используйте эфемерные кластеры или Dataproc Serverless, чтобы избежать затрат на простой; храните данные в GCS, чтобы минимизировать использование постоянных дисков.
- Добавляйте прерываемые/spot-узлы в качестве вторичных рабочих узлов для поглощения пиковой нагрузки; проектируйте систему с учетом повторных вычислений, так как задачи на потерянных узлах будут перезапущены. Не размещайте мастер-узлы на прерываемых инстансах.
- Правильно подбирайте типы машин и используйте автомасштабирование для сокращения мощностей, когда очереди пусты. Отдавайте предпочтение форматам Parquet/ORC с отсечением партиций (partition pruning) для снижения затрат на сканирование и CPU.
- Избегайте мелких файлов путем уплотнения выходных данных; меньшее количество более крупных файлов сокращает накладные расходы на метаданные и время выполнения задания.
- Для коротких периодических заданий (например, еженедельный 30-минутный Spark ETL) прерываемые рабочие узлы или serverless-режим часто обеспечивают наилучший профиль затрат.
#### Практический сценарий
Компания Acme Retail переносит локальный кластер Hadoop из 30 узлов, на котором выполняются ночные ETL-процессы на Spark и Hive для последующей аналитики. Они хотят повторно использовать существующие задания с минимальными изменениями, избежать постоянного управления кластерами, сохранять данные после удаления кластеров и сократить расходы на хранение.
Подход:
1) Размещение данных и метаданных в управляемых сервисах
- Храните все необработанные и подготовленные данные в Cloud Storage в формате Parquet с партиционированием (например, dt=YYYY-MM-DD).
- Обоснование: GCS — это надежное и недорогое хранилище, которое отделяет вычисления от хранения данных, что позволяет эфемерным кластерам и serverless-заданиям работать без постоянных дисков. Партиционированный Parquet обеспечивает проталкивание предикатов (predicate pushdown) и эффективное сканирование.
2) Централизация каталога с помощью Dataproc Metastore
- Перенесите хранилище метаданных Hive в Dataproc Metastore. Создайте внешние таблицы Hive, ссылающиеся на пути в GCS, и сохраните существующую логику схем и партиций.
- Обоснование: Управляемое хранилище метаданных позволяет нескольким эфемерным кластерам и serverless-заданиям совместно использовать определения таблиц без необходимости запускать высокодоступный инстанс MySQL/PostgreSQL.
3) Использование эфемерных кластеров Dataproc для пакетных ETL и шаблонов рабочих процессов для оркестрации
- Определите шаблон рабочего процесса, который создает кластер с нужным образом (например, 2.1-debian11), запускает задания Spark (spark-sql и pyspark) и удаляет кластер по завершении. Добавьте действия по инициализации для установки любых пользовательских библиотек.
- Обоснование: Эфемерные кластеры устраняют затраты на простой и изолируют зависимости заданий. Шаблоны рабочих процессов обеспечивают повторяемость и параметризацию (даты, входные пути).
4) Включение автомасштабирования и прерываемых рабочих узлов
- Примените политику автомасштабирования с небольшой основной группой рабочих узлов и большим пулом прерываемых вторичных рабочих узлов; настройте периоды охлаждения (cooldowns) для быстрого уменьшения масштаба после выполнения.
- Обоснование: Основные рабочие узлы поддерживают стабильность кластера; прерываемые рабочие узлы поглощают нагрузку от shuffle и широких преобразований (wide transformations) по более низкой цене. Механизмы повторных попыток в Spark/YARN обрабатывают задачи, потерянные из-за прерывания узлов.
5) Интеграция с BigQuery через коннектор Spark BigQuery
- Для загрузки измерений/фактов записывайте результаты Spark во временные (staging) таблицы BigQuery, а затем выполняйте операторы MERGE для атомарного обновления целевых таблиц. Там, где прямая перезапись безопасна, записывайте данные в партиционированные таблицы, используя режим перезаписи партиций (partition overwrite mode).
- Обоснование: BigQuery обслуживает аналитику и BI в больших масштабах; подход staging+MERGE обеспечивает транзакционно-подобные операции upsert из пакетных заданий Spark, уменьшая несогласованность данных в последующих системах.
6) Настройка Spark для производительности и надежности
- Установите количество shuffle-партиций относительно ядер исполнителей и включите AQE:
--conf spark.sql.shuffle.partitions=600
--conf spark.sql.adaptive.enabled=true
```
- Используйте broadcast joins для небольших таблиц-измерений и сохраняйте контрольные точки (checkpoint) для длинных цепочек преобразований (lineages) в GCS для повышения стабильности.
- Обоснование: Правильное партиционирование уменьшает перекос данных и накладные расходы планировщика; AQE адаптируется к профилям данных во время выполнения; контрольные точки ограничивают объем повторных вычислений после сбоев.
Усиление безопасности и сетевых настроек
- Запускайте кластеры с выделенными сервисными аккаунтами, предоставляя только роли, необходимые для путей GCS, хранилища метаданных и наборов данных BigQuery. Создавайте кластеры с частными IP-адресами в ограниченной подсети с Private Google Access и ограничивайте доступ к UI с помощью правил брандмауэра.
- Обоснование: Принцип минимальных привилегий и сетевая изоляция уменьшают поверхность атаки; частный исходящий трафик для плоскости управления позволяет избежать раскрытия в публичном интернете.
Настройка логирования, истории и оповещений
- Включите запись логов событий Spark в GCS и разверните History Server; направляйте логи драйвера/YARN в Cloud Logging с настроенным сроком хранения. Добавьте оповещения в Monitoring для длительно ожидающих контейнеров, повторяющихся сбоев задач или чрезмерной продолжительности задания.
- Обоснование: Централизованные логи помогают в поиске первопричин; проактивные оповещения позволяют на ранней стадии выявлять перекос данных, ошибки OOM или снижение производительности ввода-вывода.
Выборочная модернизация с помощью Dataproc Serverless для специальных задач и эластичных всплесков нагрузки
- Перенесите спорадические или исследовательские рабочие нагрузки Spark SQL в Dataproc Serverless; продолжайте выполнять ночные конвейеры на эфемерных кластерах до полной проверки их работы в serverless-режиме.
- Обоснование: Serverless-режим устраняет необходимость в управлении кластерами и масштабируется автоматически, что идеально подходит для непредсказуемых нагрузок; существующие рабочие процессы продолжают работать с минимальными изменениями в коде.
Проверка механизмов фиксации в объектном хранилище и управление мелкими файлами
- Установите алгоритм FileOutputCommitter v2 и уплотняйте выходные данные до 256–512 МиБ на файл с помощью repartition/coalesce перед записью.
- Обоснование: В объектных хранилищах отсутствует атомарное переименование; оптимизированные механизмы фиксации (committers) сокращают накладные расходы на копирование/переименование. Уплотнение решает проблему мелких файлов, улучшая производительность и снижая затраты.
Такая архитектура позволяет повторно использовать существующие задания Spark и Hive с минимальным рефакторингом, обеспечивает долговечность данных в GCS, централизует схемы, ограничивает радиус поражения при инцидентах безопасности, предоставляет надежные средства наблюдаемости и оптимизирует затраты за счет эфемерных кластеров, автомасштабирования, прерываемых мощностей и целевого использования serverless-выполнения.
← Обмен сообщениями · Все домены · Прием →
Отработать эти вопросы → · Тесты на время на ExamRoll.io →
Pass the whole exam — not just this question
You found this answer. Get every verified question and explanation in one place, and save hours of prep. Free to start.
Сдайте экзамен →