Google PDE: Потоковая обработка с Dataflow и Apache Beam — Руководство по подготовке

Часть Google Professional Data Engineer — Руководство по подготовке. Практикуйтесь с проверенными ответами в центре экзаменов Google, или пройдите тесты на время на ExamRoll.io.

Обзор

Потоковая обработка в Google Cloud основана на единой модели программирования Apache Beam, выполняемой исполнителем Dataflow. Beam предоставляет логическую абстракцию — конвейеры преобразований над PCollections, — которая отделяет ваш код от деталей выполнения, таких как параллелизм, автомасштабирование и отказоустойчивость. В потоковой обработке корректность зависит от временной семантики (время события и время обработки), оконной обработки (фиксированные, скользящие, сессионные, глобальные окна), ватермарок, триггеров и обработки опоздавших данных. Для эффективной эксплуатации Dataflow требуется правильный подбор размера воркеров, политика автомасштабирования, выбор движка потоковой обработки и механизма shuffle, идемпотентная архитектура приёмников, обработка недоставленных сообщений и надёжная наблюдаемость.

Модель Apache Beam и временная семантика

Сценарии сбоев и компромиссы:

Эксплуатация Dataflow для потоковых рабочих нагрузок

Развертывание, шаблоны и стратегии обновления

undefined

undefined

Практический сценарий

Компания NovaTrack Inc. принимает глобальную телеметрию с 50 000 датчиков температуры и должна предоставлять поминутные агрегаты, сохранять необработанные данные и отображать их на дашборде в реальном времени. Ожидаются периодические некорректные сообщения и доставка с нарушением порядка. Решение должно автоматически масштабироваться, предоставлять некорректные записи для анализа и поддерживать обновления без простоя.

Подход:

  1. Прием данных и семантика времени

    • Создать региональный топик Pub/Sub и издателей для каждого региона с атрибутами deviceId и eventTs (RFC3339). По возможности включить ключи упорядочивания (ordering keys) по deviceId.
    • Обоснование: Pub/Sub обеспечивает надежный, эластичный прием данных с доставкой «как минимум один раз» (at-least-once). Прикрепление временных меток события на входе сохраняет истинное время события; упорядочивание по устройствам уменьшает нарушение порядка внутри потока от одного устройства без создания централизованных узких мест.
  2. Потоковый конвейер Dataflow с окнами по времени события

    • Читать из выделенной подписки через PubSubIO, извлекая eventTs в качестве временной метки Beam, с откатом на publishTime в случае отсутствия.
    • Применить FixedWindows размером в 1 минуту с ранним триггером через 30 секунд и поздними срабатываниями на каждый опоздавший элемент; установить допустимое опоздание (allowed lateness) в 10 минут и накапливать панели (accumulating panes).
    • Обоснование: Окна по времени события обеспечивают точные поминутные агрегаты; ранние срабатывания поставляют данные на дашборд со свежестью менее минуты; поздние срабатывания корректируют агрегаты по мере поступления запоздавших данных. Ограничение на опоздание лимитирует размер состояния и затраты.
  3. Валидация, обогащение и маршрутизация в очередь недоставленных сообщений (dead-letter)

    • Реализовать ParDo, который парсит JSON, проверяет схему и диапазоны значений, и обогащает данные небольшим статическим справочником через дополнительный вход (side input), загружаемый из BigQuery при запуске задания.
    • Использовать TupleTags для отправки валидных записей в основной выходной поток, а ошибочных — в PCollection для недоставленных сообщений (dead-letter), содержащую исходные данные, ошибку, deviceId и временную метку парсинга; записывать DLQ в партиционированную таблицу BigQuery.
    • Обоснование: Дополнительные входы (side inputs) хранят справочные данные в памяти для низкой задержки. Сбор недоставленных сообщений позволяет анализировать и целенаправленно повторно обрабатывать некорректные строки, не блокируя основной поток.
  4. Агрегация и предотвращение «горячих» ключей (hot-key)

    • Ключировать по deviceId и вычислять поминутные avg/min/max с помощью CombineFns. Для метрик топ-N по регионам, шардировать по region#N, чтобы избежать «горячих» ключей, а затем повторно агрегировать.
    • Обоснование: Комбайнеры (Combiners) минимизируют объем перемешиваемых данных (shuffle) и затраты; шардирование ключей предотвращает узкие места на одном ключе при сборе данных по регионам.
  5. Приемники (sinks) и семантика «ровно один раз» (exactly-once)

    • Записывать необработанные проверенные события и поминутные агрегаты в BigQuery с помощью BigQueryIO и Storage Write API. Установить стабильный ID вставки (insert id) на основе deviceId + eventTs для идемпотентности при любых пользовательских повторных попытках.
    • Обоснование: Storage Write API обеспечивает высокопроизводительный прием данных с низкой задержкой и семантикой «ровно один раз» в рамках одного потока. Стабильные ID обеспечивают дедупликацию на стороне получателя в случае повторных отправок.
  6. Стратегия консистентности для дашборда

    • Дашборд запрашивает партиционированные таблицы с агрегатами с окном просмотра в 2 минуты относительно водяного знака или с фиксированной задержкой в 2 раза превышающей наблюдаемую задержку доступности для потоковых данных.
    • Обоснование: Видимость потоковых данных в BigQuery является согласованной в конечном счете (eventually consistent); небольшая задержка чтения предотвращает пропуск строк, находящихся в обработке, сохраняя при этом поведение, близкое к реальному времени.
  7. Эксплуатация: автомасштабирование и Streaming Engine

    • Включить Streaming Engine; установить maxWorkers на основе ожидаемого пика (например, 3x от среднего), выбрать тип машины, подходящий для парсинга и шифрования (CPU-bound), и увеличить загрузочный диск для временных данных shuffle.
    • Мониторить отставание водяного знака (watermark lag), объем необработанных данных в секундах (backlog seconds), CPU и пропускную способность для каждого шага; настроить оповещения на устойчивое отставание и всплески в скорости поступления сообщений в DLQ.
    • Обоснование: Streaming Engine выносит состояние/shuffle за пределы воркеров, что обеспечивает эластичность и упрощает обновления; правильный подбор ресурсов и мониторинг предотвращают скрытое нарушение SLO.
  8. Развертывание и обновления с помощью Flex Templates

    • Упаковать конвейер как Flex Template с параметрами: входная подписка, выходные таблицы, таблица DLQ, maxWorkers и регион. При несовместимом изменении запустить новый конвейер, нацеленный на тот же топик, но с новой подпиской, проверить выходные данные, а затем осушить (drain) старое задание. Опционально можно создать снимок Pub/Sub и настроить чтение новой подписки из этого снимка, чтобы гарантировать отсутствие пробелов в данных.
    • Обоснование: Flex Templates обеспечивают повторяемые, параметризованные развертывания. Проверенное «сине-зеленое» (blue/green) переключение с осушением (drain) позволяет достичь нулевой потери данных и минимального простоя.
  9. Повторная обработка и пакетное заполнение исторических данных

    • Хранить сжатые Avro-файлы с необработанными событиями в Cloud Storage через дополнительный выход (side output); запускать пакетный конвейер Dataflow для заполнения исторических данных или повторной обработки в BigQuery при изменении моделей или схем.
    • Обоснование: Надежные архивы необработанных данных поддерживают воспроизводимость и эволюцию схемы, не влияя на «горячий» путь обработки данных.

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


Аналитика с BigQuery и инженерия хранилищ · Все домены · Обмен сообщениями

Отработать эти вопросы → · Тесты на время на 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+ языков
  • Еженедельные обновления контента
  • Награды и рефералы
  • Приоритетная поддержка
Начать бесплатную пробную версию

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

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