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

    --conf mapreduce.fileoutputcommitter.algorithm.version=2
    ```

  - Используйте форматы **Parquet/ORC** с отсечением столбцов (column pruning) и проталкиванием предикатов (predicate pushdown). Управляйте мелкими файлами с помощью уплотнения (compaction), стремясь к размеру файла 128512 МиБ для эффективного сканирования.
- 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)
  - Входное партиционирование: достаточное количество партиций для утилизации всех ядер; начните с 24x от общего числа ядер исполнителей (executor cores). Управляется через spark.default.parallelism (для RDD) и опции чтения (для DataFrames).
  - Партиции для перемешивания (shuffle partitions): значение по умолчанию 200 часто бывает недостаточным или избыточным. Настройте:
--conf spark.sql.shuffle.partitions= {total_executor_cores * 2 to 3}
```
      --conf spark.sql.autoBroadcastJoinThreshold=64m
      ```

    - «Подсаливание» ключей для «горячих» партиций; применяйте предварительную агрегацию на стороне map; выполняйте фильтрацию как можно раньше.
    - Включите Adaptive Query Execution (AQE) для объединения партиций после перемешивания и обработки перекошенных соединений (skewed joins):
  --conf spark.sql.adaptive.enabled=true
  ```
      --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 (...)
```

Безопасность, наблюдаемость и затраты

    --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
 ```
  1. Усиление безопасности и сетевых настроек

    • Запускайте кластеры с выделенными сервисными аккаунтами, предоставляя только роли, необходимые для путей GCS, хранилища метаданных и наборов данных BigQuery. Создавайте кластеры с частными IP-адресами в ограниченной подсети с Private Google Access и ограничивайте доступ к UI с помощью правил брандмауэра.
    • Обоснование: Принцип минимальных привилегий и сетевая изоляция уменьшают поверхность атаки; частный исходящий трафик для плоскости управления позволяет избежать раскрытия в публичном интернете.
  2. Настройка логирования, истории и оповещений

    • Включите запись логов событий Spark в GCS и разверните History Server; направляйте логи драйвера/YARN в Cloud Logging с настроенным сроком хранения. Добавьте оповещения в Monitoring для длительно ожидающих контейнеров, повторяющихся сбоев задач или чрезмерной продолжительности задания.
    • Обоснование: Централизованные логи помогают в поиске первопричин; проактивные оповещения позволяют на ранней стадии выявлять перекос данных, ошибки OOM или снижение производительности ввода-вывода.
  3. Выборочная модернизация с помощью Dataproc Serverless для специальных задач и эластичных всплесков нагрузки

    • Перенесите спорадические или исследовательские рабочие нагрузки Spark SQL в Dataproc Serverless; продолжайте выполнять ночные конвейеры на эфемерных кластерах до полной проверки их работы в serverless-режиме.
    • Обоснование: Serverless-режим устраняет необходимость в управлении кластерами и масштабируется автоматически, что идеально подходит для непредсказуемых нагрузок; существующие рабочие процессы продолжают работать с минимальными изменениями в коде.
  4. Проверка механизмов фиксации в объектном хранилище и управление мелкими файлами

    • Установите алгоритм 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.

Сдайте экзамен →

Просмотреть Google →

Related guides

Все включено

Одна подписка. Каждый экзамен.

Каждый план открывает неограниченный поиск ответов, практические тесты, объяснения AI и полную библиотеку ресурсов — на более чем 20 языках.

Ежемесячно
24.87
Just €0.83/day
Все включено:
  • Неограниченный поиск ответов
  • Неограниченные практические тесты
  • Объяснения на основе AI
  • Полная библиотека ресурсов
  • 20+ языков
  • Еженедельные обновления контента
  • Награды и рефералы
  • Приоритетная поддержка
Начать бесплатную пробную версию

Кредитная карта не требуется*

Лучшая цена
12 месяцев
179.87
Just €0.49/daySave 40%
Все включено:
  • Неограниченный поиск ответов
  • Неограниченные практические тесты
  • Объяснения на основе AI
  • Полная библиотека ресурсов
  • 20+ языков
  • Еженедельные обновления контента
  • Награды и рефералы
  • Приоритетная поддержка
Начать бесплатную пробную версию

Кредитная карта не требуется*

✓ Включен бесплатный план · ✓ Отмена в любое время · ✓ Все планы открывают полный продукт