Google PDE: Оркестрация рабочих процессов и автоматизация пайплайнов — Руководство по подготовке
Часть Google Professional Data Engineer — Руководство по подготовке. Практикуйтесь с проверенными ответами в центре экзаменов Google, или пройдите тесты на время на ExamRoll.io.
Обзор
Оркестрация рабочих процессов и автоматизация конвейеров координируют задачи обработки данных между сервисами, чтобы прием, преобразование, проверка качества и публикация данных происходили надежно, безопасно и экономически эффективно. В Google Cloud оркестрация должна соответствовать модели выполнения каждой рабочей нагрузки: пакетные задания по расписанию, потоковая обработка на основе событий, специальные (ad-hoc) или длительные задания. Цели проектирования — повторяемость, идемпотентность, наблюдаемость, принцип наименьших привилегий и безопасное продвижение по средам.
Ключевые варианты:
- Оркестрация пакетных заданий, ориентированная на код, с помощью Cloud Composer (Apache Airflow) для DAG, зависимостей между задачами и расширенного планирования.
- Бессерверная хореография API с помощью Cloud Workflows для легковесных, управляемых событиями последовательностей между сервисами.
- Конечные точки выполнения, такие как задания Cloud Run или Dataproc, запускаемые Cloud Scheduler для cron-задач или Eventarc для событий.
- Нативная SQL-оркестрация с помощью Dataform для преобразований в BigQuery, утверждений (assertions) и управления релизами.
Операционная модель делает акцент на повторных попытках с ограниченной экспоненциальной задержкой, тайм-аутах, SLA, догоняющих запусках (catchup) и заполнении пропущенных данных (backfills), идемпотентной архитектуре задач для безопасных повторных запусков и надежной обработке сбоев с перехватом в очередь недоставленных сообщений (dead-letter). Безопасность обеспечивается через сервисные аккаунты для каждого конвейера, изоляцию секретов, параметризацию и IAM с принципом наименьших привилегий. CI/CD, инфраструктура как код и комплексная телеметрия дополняют подход, готовый к использованию в производственной среде.
Оркестрация в Google Cloud: инструменты и паттерны
Cloud Composer (Airflow)
DAG (направленные ациклические графы) определяют графы выполнения с явными зависимостями. Используйте TaskFlow API или операторы (например, BigQuery, Dataflow, Dataproc, Cloud Run) для описания задач. Сенсоры и откладываемые (deferrable) операторы снижают нагрузку на планировщик при условиях ожидания (например, финализация объекта в Cloud Storage или появление партиции в BigQuery).
Планирование: cron-выражения, start_date, end_date и catchup управляют историческими запусками. Используйте catchup для заполнения пропущенных данных (backfills); отключайте для задач, смежных с потоковой обработкой, или для неидемпотентных целей. Ограничивайте параллелизм с помощью max_active_runs и пулов для защиты нижестоящих систем.
Зависимости: set_upstream/set_downstream или зависимости taskflow. Для оркестрации на основе метаданных генерируйте задачи динамически из управляющей таблицы BigQuery (например, список клиентов/партиций), используя динамическое сопоставление задач (dynamic task mapping), что сохраняет время парсинга DAG стабильным и делает задачи управляемыми данными.
Пример (краткий) фрагмента DAG: from airflow import DAG from datetime import datetime, timedelta from airflow.providers.google.cloud.operators.dataflow import DataflowStartFlexTemplateOperator
default_args = dict(retries=3, retry_delay=timedelta(minutes=5), execution_timeout=timedelta(hours=2), sla=timedelta(hours=3)) with DAG(‘daily_csv_import’, start_date=datetime(2023,1,1), schedule_interval=‘0 2 * * *’, catchup=True, max_active_runs=1, default_args=default_args) as dag: import_job = DataflowStartFlexTemplateOperator( task_id=‘import’, body={’launchParameter’: {‘jobName’: ‘csv-import-{{ ds_nodash }}’, ‘parameters’: {‘dlq_table’: ‘bqproj.dlq.bad_rows’}}} )
Cloud Workflows, Cloud Scheduler, задания Cloud Run и выполнение на основе событий
- Cloud Workflows оркестрирует API Google и HTTP-эндпоинты со встроенными повторными попытками, циклами, параллельными ветвями и компенсирующей логикой. Он идеально подходит для легковесного управления потоком выполнения между сервисами, такими как BigQuery, Dataflow, Batch и задания Cloud Run.
- Cloud Scheduler запускает Workflows, темы Pub/Sub или HTTP-сервисы для автоматизации в стиле cron. Для ежедневной пакетной обработки в 02:00 запланируйте Workflow, который запускает задание Dataflow или Dataproc.
- Задания Cloud Run выполняют контейнеризированные шаги пакетной обработки с автоматическими повторами и минимальными операционными затратами. Они хорошо сочетаются с Workflows для многоэтапных задач обработки данных или для предварительной/постобработки вокруг Dataflow или BigQuery.
- Управляемый событиями подход: используйте Eventarc для маршрутизации событий финализации объектов Cloud Storage, сообщений Pub/Sub или Audit Logs в Cloud Run или Workflows. Для получения уведомлений о заданиях вставки (insert-job) в BigQuery для одной таблицы создайте приемник (sink) Cloud Logging с расширенным фильтром для Pub/Sub, а затем запускайте ваш потребитель (consumer) из этой темы.
Dataform: SQL-воркфлоу для BigQuery
- Моделируйте графы зависимостей с помощью ref(), определяйте таблицы/представления/инкрементальные обновления и оркестрируйте сборки по тегам или расписаниям. Dataform компилирует SQLX в упорядоченные планы выполнения, обеспечивая оркестрацию на основе метаданных из декларативных определений.
- Утверждения (assertions) обеспечивают качество данных. Утверждение — это запрос, который должен вернуть ноль строк для успешного прохождения. Пример утверждения: – definitions/assert_non_negative_prices.sqlx config { type: “assertion” } SELECT 1 FROM ${ref(‘prices_daily’)} WHERE price < 0 LIMIT 1
- Релизы и контроль репозитория: храните код в репозитории, используйте ветки и ревью, а также продвигайте релизы с тегами по средам (например, dev, test, prod) с переменными, специфичными для каждой среды. Контролируйте развертывания с помощью проверок CI/CD и результатов утверждений.
Dataproc, Dataflow и паттерны хранения данных
- Для повторного использования Hadoop/Spark с минимальными операционными затратами используйте Dataproc с коннектором GCS, чтобы сохранять данные дольше жизненного цикла кластера и минимизировать затраты на постоянные диски. Создавайте эфемерные кластеры для каждого задания для изоляции и контроля затрат; оркестрируйте с помощью Composer или Workflows.
- Для пакетного приема данных с некорректными строками запустите Dataflow, чтобы записывать корректные записи в BigQuery, а ошибки парсинга/валидации направлять в таблицу BigQuery для недоставленных сообщений (dead-letter) для анализа.
Надежность, обработка сбоев и идемпотентность
Повторные попытки, тайм-ауты и экспоненциальная задержка
- Используйте ограниченную экспоненциальную задержку (exponential backoff) для временных сбоев и ограничивайте общее окно повторных попыток в соответствии с SLA задания. Например, фронтенд или задача, опрашивающая базу данных каждые 15 минут, должна повторять попытки с экспоненциальной задержкой в течение 15 минут, а затем выдавать контролируемую ошибку.
- В Airflow настраивайте
execution_timeoutдля каждой задачи и глобальные SLA для DAG; в Workflows устанавливайте тайм-ауты для каждого шага и политики повтора сmax_doublingsиmax_retry_duration. Для заданий Cloud Run задавайте количество повторных попыток и задержку (backoff).
Заполнение пропущенных данных (backfill), наверстывание (catchup) и обработка сбоев
- Включайте наверстывание (catchup) для пересчета исторических данных, когда задачи идемпотентны, а источники партиционированы по дате. Для недетерминированных результатов или внешних побочных эффектов рассмотрите использование DAG только для заполнения (backfill) или таблиц аудита записи для отслеживания созданных данных.
- Используйте топики/таблицы недоставленных сообщений (dead-letter) для сбоев на уровне записей при потоковой/пакетной обработке. Для пакетной обработки в Dataflow отлавливайте некорректные строки с помощью тегов ошибок и агрегируйте метрики ошибок; для потоковой — используйте DLQ в Pub/Sub.
Проектирование идемпотентных задач и повторные запуски
- BigQuery: предпочитайте
MERGEилиINSERTс ключами для дедупликации; используйтеinsertIdдля дедупликации потоковых вставок. Для пакетной обработки записывайте данные во временную (staging) таблицу, а затем выполняйтеMERGEв целевую таблицу в рамках транзакционно безопасного шага, чтобы обеспечить возможность полного перезапуска. - Cloud Storage: используйте предварительные условия по поколению (generation preconditions) и детерминированные имена объектов (например, префикс/дата/хеш), чтобы повторные запуски безопасно выполняли перезапись только в ожидаемых случаях.
- Pub/Sub и Dataflow: проектируйте с расчетом на доставку «хотя бы один раз» (at-least-once). Включайте в сообщения идентификаторы (например, ID пакета, логическая временная метка события), чтобы последующие системы могли выполнять дедупликацию и анализировать задержки. Если бизнес-правила допускают семантику «побеждает первое обработанное событие», задокументируйте этот компромисс и отслеживайте расхождения; в противном случае определяйте «победителя» по времени события, используя дополнительные критерии для разрешения неоднозначностей.
- Восстановление после частичного сбоя: партиционируйте выходные данные по
run_idили дате, записывайте маркеры завершения и настраивайте зависимость последующих задач от этих маркеров. Повторно обрабатывайте только те партиции, которые помечены как незавершенные.
Устранение неполадок и масштабируемость
- Если на потоковой панели мониторинга отсутствуют события, но в Pub/Sub они есть, прогоните через конвейер Dataflow известный фиксированный набор данных, чтобы выявить дефекты в преобразованиях. Проверьте настройки окон (windowing), триггеров и допустимой задержки (allowed lateness).
- Распространенный сценарий сбоя: создание потокового конвейера без подходящих окон/триггеров для неограниченных источников или некорректное использование разделенного окна (sharded window) может привести к сбою при создании конвейера или к неконтролируемому росту состояния (state blowups).
- Масштабируйте Dataflow с помощью параметра
max workersи алгоритма автомасштабирования; при всплесках нагрузки (например, 50 000 установок) увеличьте максимальное количество воркеров, чтобы обеспечить горизонтальное масштабирование в пиковые моменты.
Безопасность, параметризация, среды и CI/CD
Параметризация и управление конфигурацией
- Выносите конфигурацию для каждой среды вовне. В Composer используйте Variables, Connections и переменные среды; шаблонизируйте параметры DAG по дате выполнения или партиции. В Workflows используйте аргументы времени выполнения и отдельные воркфлоу для каждой среды или считывайте конфигурацию из Secret Manager.
- Используйте оркестрацию на основе метаданных, считывая управляющую таблицу (например, датасет с конфигурацией в BigQuery), которая содержит списки клиентов, источников или партиций. Генерируйте задачи динамически, чтобы изменения в коде были отделены от изменений, управляемых данными.
Секреты, сервисные аккаунты и принцип минимальных привилегий
- Храните учетные данные в Secret Manager и обращайтесь к ним во время выполнения. Избегайте встраивания секретов в код или в Airflow Variables.
- Назначайте каждому конвейеру отдельный сервисный аккаунт с минимально необходимым набором IAM-ролей. Для регулируемого доступа к BigQuery изолируйте клиентские данные в отдельные датасеты, предоставляйте роли для конкретных датасетов только утвержденным пользователям и ограничивайте доступ к BigQuery API только для утвержденных участников (principals). Для мультитенантности создавайте по датасету на каждого клиента и привязывайте только соответствующие роли.
CI/CD и инфраструктура как код
- Управляйте инфраструктурой (средами Composer, Workflows, заданиями Scheduler, топиками Pub/Sub, приемниками логов) с помощью Terraform. Используйте модули для стандартизации проектов/сред, секретов и сервисных аккаунтов.
- Собирайте и тестируйте код конвейеров с помощью Cloud Build или GitHub Actions. Автоматизируйте юнит-тесты, проверку SQL-кода (linting), «сухие запуски» (dry-runs) Dataform и валидацию DAG в Airflow. Продвигайте артефакты с помощью тегов; для Composer упаковывайте DAG в развертываемые пакеты; для Dataform используйте релизные ветки, которые продвигаются после успешного прохождения проверок (assertions).
- Продвижение развертывания: dev → test → prod через отдельные проекты и параметризованные конфигурации. Используйте непрерывную доставку (continuous delivery) с этапами ручного подтверждения и окнами изменений для высокорисковых продвижений.
Наблюдаемость, оповещения и регламенты (Runbooks)
Телеметрия и оповещения
- Направляйте все журналы оркестрации в Cloud Logging со структурированными полями (pipeline, dag_id, run_id, task_id, partition). Экспортируйте журналы ошибок в Monitoring с помощью метрик на основе журналов (log-based metrics). Настройте оповещения для следующих случаев:
- Пропущенные запуски по расписанию или нарушения SLA
- Несколько последовательных сбоев задач
- Рост бэклога (например, неподтвержденные сообщения в Pub/Sub, системное отставание в Dataflow)
- Сбои проверок качества данных
- Cloud Composer: отслеживайте длительность выполнения DAG/задач, процент успешных выполнений, глубину очереди и работоспособность планировщика. Настройте on_failure_callback для отправки уведомлений и запуска регламентов по устранению сбоев.
- Cloud Workflows: анализируйте журналы выполнения (Execution logs) и задержки на шагах; добавляйте явные повторные попытки и обработчики ошибок; создавайте пользовательские журналы с идентификаторами корреляции.
- Уведомления об изменении таблиц BigQuery: создайте приемник (sink) журналов на уровне проекта с расширенным фильтром для заданий вставки (insert jobs), нацеленных на конкретную таблицу, и экспортируйте их в Pub/Sub; ваш инструмент мониторинга подписывается на этот топик для получения мгновенных оповещений без шума от других таблиц.
Разработка регламентов (Runbook)
- Для каждого конвейера документируйте триггеры, зависимости, SLA, процедуры отката/повторного выполнения и безопасные шаги для повторной обработки данных (backfill). Включите «воспроизведение на фиксированном наборе данных» (fixed dataset replay) для Dataflow, инструкции по остановке потокового задания (draining), повторной обработке сбойных партиций и устранению проблем с сообщениями из DLQ.
- Фиксируйте типичные признаки сбоев (например, отказано в доступе, превышена квота, несоответствие схемы) с помощью деревьев принятия решений и путей эскалации.
Практический сценарий
Компании Acme Retail Analytics необходимо ежедневно загружать CSV-файлы от партнеров, которые иногда содержат некорректные строки, преобразовывать и загружать корректные данные в BigQuery, а также предоставлять некорректные строки для расследования. Они также хотят реализовать обогащение данных на основе событий для обновления цен почти в реальном времени и обеспечить безопасное развертывание из среды разработки в производственную.
Подход:
Хранилище и триггеры на основе событий
- Создайте выделенный бакет Cloud Storage с версионированием объектов и унифицированным доступом на уровне бакета. Включите уведомления о финализации объектов (object finalize) в Pub/Sub через Eventarc.
- Обоснование: Финализация объекта — это надежное событие для запуска последующей загрузки; версионирование обеспечивает возможность повторных запусков и аудита.
Пакетная загрузка с обработкой неисправных сообщений (dead-letter)
- Используйте Cloud Composer для запуска ежедневного DAG в Airflow в 02:00 с включенным параметром catchup. DAG запускает пакетное задание Dataflow, которое разбирает CSV-файлы, проверяет схему и записывает корректные записи в BigQuery, используя детерминированные промежуточные таблицы, а затем выполняет MERGE в партиционированные целевые таблицы. Некорректные/сбойные записи направляйте в таблицу для неисправных данных (dead-letter table) в BigQuery.
- Обоснование: Dataflow масштабирует парсинг/валидацию; MERGE обеспечивает идемпотентность; сбор неисправных данных позволяет анализировать их, не блокируя конвейер, что соответствует рекомендуемому шаблону для обработки некорректных строк.
Обогащение данных на основе событий
- Разверните задание Cloud Run для выполнения легковесного обогащения данных при инкрементальных обновлениях цен. Запускайте его через Cloud Workflows, который прослушивает сообщения Pub/Sub из Eventarc при поступлении небольших файлов с обновлениями в течение дня.
- Обоснование: Бессерверные контейнеры с Workflows обеспечивают оркестрацию небольших событий с низкой задержкой и минимальными операционными затратами, в то время как тяжелые преобразования остаются в пакетной обработке.
Средства обеспечения надежности
- Настройте повторные попытки с экспоненциальной задержкой (exponential backoff) для временных сбоев в заданиях Dataflow и Cloud Run, ограничив общее время повторных попыток в рамках SLA для DAG. Установите тайм-ауты выполнения для каждой задачи и колбэки on_failure в Airflow; в Workflows установите max_doublings и max_retry_duration.
- Обоснование: Ограниченная задержка между повторами помогает соблюдать SLA и предотвращает бесконечные повторные попытки.
Безопасность и принцип наименьших привилегий
- Запускайте каждый компонент под выделенным сервисным аккаунтом: SA для оркестратора Composer, SA для воркеров Dataflow, SA для задания Cloud Run. Предоставляйте только необходимые роли: GCS read для бакета с входящими данными для Dataflow, BigQuery dataEditor для целевых наборов данных и Viewer для журналов. Храните секреты в Secret Manager и обращайтесь к ним во время выполнения.
- Обоснование: Это обеспечивает соблюдение принципа наименьших привилегий и изолирует радиус поражения при сбое.
Оркестрация на основе метаданных
- Ведите управляющую таблицу в BigQuery, содержащую список источников-партнеров, шаблоны файлов и целевые наборы данных. Во время выполнения DAG Airflow запрашивает эту таблицу и использует динамическое сопоставление задач (dynamic task mapping) для создания задач для каждого партнера.
- Обоснование: Добавление нового партнера становится изменением данных, а не кода, что снижает риск при развертывании.
Наблюдаемость и оповещения
- Создавайте структурированные журналы с run_id и partner_id. Создайте политики оповещения о нарушении SLA для DAG, системном отставании в Dataflow и ненулевом количестве записей в таблице неисправных данных. Для вставок в целевую таблицу BigQuery настройте приемник (sink) Cloud Logging с расширенным фильтром для этой таблицы, отправляющий данные в топик Pub/Sub, который используется инструментом мониторинга Acme.
- Обоснование: Гранулированные оповещения позволяют быстро выявлять и устранять проблемы без лишнего шума.
CI/CD и развертывание по средам
- Управляйте инфраструктурой (бакеты, Pub/Sub, Eventarc, Composer, Workflows, наборы данных BigQuery) с помощью Terraform. Используйте Cloud Build для проверки синтаксиса DAG в Airflow, запуска модульных тестов и развертывания в среду разработки Composer. Продвигайте изменения в тестовую и производственную среды с помощью параметризованных конфигураций и этапов ручного утверждения после прохождения проверок в Dataform и интеграционных тестов.
- Обоснование: Декларативные, повторяемые развертывания и безопасное продвижение изменений между средами.
Регламент (Runbook) и восстановление
- Задокументируйте шаги для повторной обработки данных за определенную дату: восстановление CSV из версионированных объектов, повторный запуск задания Dataflow для этой партиции, выполнение MERGE результатов и просмотр записей в DLQ. Включите процедуру «воспроизведения на фиксированном наборе данных» (fixed dataset replay) для изоляции ошибок преобразования при возникновении расхождений.
- Обоснование: Идемпотентный дизайн и задокументированные процедуры восстановления упрощают устранение частичных сбоев.
← Прием · Все домены · Машинное обучение →
Отработать эти вопросы → · Тесты на время на 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.
Сдайте экзамен →