Google PDE: Потоковая обработка с Dataflow и Apache Beam — Руководство по подготовке
Часть Google Professional Data Engineer — Руководство по подготовке. Практикуйтесь с проверенными ответами в центре экзаменов Google, или пройдите тесты на время на ExamRoll.io.
Обзор
Потоковая обработка в Google Cloud основана на единой модели программирования Apache Beam, выполняемой исполнителем Dataflow. Beam предоставляет логическую абстракцию — конвейеры преобразований над PCollections, — которая отделяет ваш код от деталей выполнения, таких как параллелизм, автомасштабирование и отказоустойчивость. В потоковой обработке корректность зависит от временной семантики (время события и время обработки), оконной обработки (фиксированные, скользящие, сессионные, глобальные окна), ватермарок, триггеров и обработки опоздавших данных. Для эффективной эксплуатации Dataflow требуется правильный подбор размера воркеров, политика автомасштабирования, выбор движка потоковой обработки и механизма shuffle, идемпотентная архитектура приёмников, обработка недоставленных сообщений и надёжная наблюдаемость.
Модель Apache Beam и временная семантика
Конвейеры, преобразования, PCollections, исполнители:
- Конвейер Beam применяет направленный ациклический граф преобразований (PTransforms) к коллекциям PCollections (ограниченным или неограниченным).
- Исполнители (Dataflow, Spark, Flink, Direct) запускают конвейер; Dataflow обеспечивает управляемое автомасштабирование, создание контрольных точек и операционную видимость.
- Преобразования включают поэлементные (ParDo), группировку и объединение (GroupByKey, Combine), соединения (CoGroupByKey) и операции ввода-вывода (PubSubIO, BigQueryIO, FileIO).
Окна:
- Фиксированные окна: непересекающиеся временные срезы (например, минутные “кувыркающиеся” окна) для периодических агрегаций.
- Скользящие окна: пересекающиеся окна для сглаженных “скользящих” метрик (например, 5-минутные окна, сдвигающиеся каждую минуту).
- Сессионные окна: динамические окна, которые закрываются после периода неактивности; идеально подходят для пользовательских сессий или всплесков активности устройств.
- Глобальное окно: представление всего неограниченного потока без разделения на окна по умолчанию; часто используется в паре с триггерами для периодической материализации.
Время события и время обработки:
- Время события: момент, когда событие произошло в источнике; позволяет выполнять логически согласованные агрегации, несмотря на переменные задержки при передаче.
- Время обработки: момент, когда событие наблюдается конвейером; полезно для операционных триггеров, но не для семантической корректности.
Ватермарки:
- Ватермарка оценивает полноту данных по времени события (предположение исполнителя о том, что он увидел все события до момента времени T).
- Ватермарки могут продвигаться неравномерно или останавливаться из-за обратного давления (backpressure) или задержек в источнике; опоздавшие данные — это все данные, приходящие с временной меткой < ватермарки.
Триггеры и опоздания:
- По умолчанию: триггер AfterWatermark, который срабатывает, когда ватермарка проходит конец окна; при нулевой допустимой задержке (allowed lateness = 0) опоздавшие данные отбрасываются.
- Ранние срабатывания (по времени обработки или по количеству элементов) позволяют получать предварительные результаты с низкой задержкой.
- Поздние срабатывания позволяют вносить исправления при поступлении опоздавших данных; режим накопления (accumulation mode) определяет, будут ли панели (panes) накапливать результаты или отбрасывать предыдущий вывод.
- Выбирайте допустимую задержку исходя из требований бизнеса и компромиссов между хранением и вычислениями; большая задержка увеличивает объём сохраняемого состояния и стоимость.
Обработка с состоянием, таймеры, сессионизация, дедупликация:
- DoFn с состоянием (Stateful DoFns) хранят состояние для каждого ключа (например, последнее виденное событие, текущие агрегаты) и устанавливают таймеры для вывода или очистки состояния.
- Сессионизация естественным образом выражается через сессионные окна (SessionWindows); для пользовательской логики используйте состояние по ключу и таймеры по времени обработки/события.
- Дедупликация: используйте стабильный идентификатор для каждого события и либо Distinct/Combine для каждого окна, либо состояние по ключу (например, фильтр Блума или множество с TTL). Это компромисс между использованием памяти и ложноположительными срабатываниями с одной стороны и строгой точностью с другой.
Сценарии сбоев и компромиссы:
- Использование окон по времени обработки для бизнес-метрик приводит к расхождениям при всплесках нагрузки или повторных попытках; предпочитайте окна по времени события.
- Слишком маленькие окна с частыми ранними срабатываниями триггеров вызывают избыточный вывод панелей (panes) и усиление записи в приёмник (write amplification).
- Неограниченная допустимая задержка может привести к разрастанию состояния; всегда ограничивайте TTL состояния и устанавливайте таймеры для очистки неактивных ключей.
Эксплуатация Dataflow для потоковых рабочих нагрузок
Определение размера и автомасштабирование воркеров:
- Горизонтальное автомасштабирование добавляет/удаляет воркеры на основе размера бэклога, отставания водяного знака, загрузки ЦП и пропускной способности; установите разумное значение maxWorkers для поглощения пиковых нагрузок.
- Выбирайте типы машин в зависимости от узких мест: для задач, ограниченных ЦП (больше vCPU), для задач, ограниченных памятью (типы с большим объемом памяти), для задач, ограниченных сетью (более крупные ВМ уменьшают накладные расходы на shuffle).
- Увеличивайте загрузочный диск при интенсивных операциях shuffle или при использовании файловых приемников. Отслеживайте системное отставание (system lag) и бэклог в секундах (backlog seconds).
Streaming Engine и shuffle:
- Streaming Engine выносит состояние и shuffle во внутреннюю службу (бэкенд), повышая эластичность, снижая нагрузку на память воркеров и обеспечивая более быстрые обновления.
- Для этапов с интенсивной пакетной обработкой или массовой группировкой по ключам используйте Dataflow Shuffle, чтобы перенести операции ввода-вывода shuffle с воркеров. Оба подхода уменьшают сбои «горячих» воркеров и интенсивное использование диска (disk thrashing).
Противодавление, горячие ключи и перекос данных:
- Dataflow управляет противодавлением через динамическую перебалансировку работы; тем не менее, при необходимости настраивайте управление потоком данных в источнике (например, количество ожидающих сообщений/байт в Pub/Sub).
- Горячие ключи (например, популярные идентификаторы) создают «отстающих» (stragglers). Для решения этой проблемы используйте шардирование ключей (key#N), частичную предварительную агрегацию с последующей сменой ключа или аппроксимации на основе эскизов (sketch-based).
- Перекос из-за записей-выбросов (огромные полезные нагрузки) или всплесков активности от издателей может потребовать создания разделов для каждого издателя, пакетирования или сжатия.
Интеграция с Pub/Sub:
- Используйте топики Pub/Sub для приема данных; включите атрибуты сообщений для метаданных (например, deviceId, метка времени события).
- Выполняйте прием данных с помощью PubSubIO; извлекайте метки времени событий из атрибутов или из полезной нагрузки, в противном случае используйте время публикации.
- Ключи упорядочивания (ordering keys) обеспечивают порядок для каждого ключа; Dataflow все равно требует идемпотентного поведения на последующих этапах из-за гарантии доставки «как минимум один раз» (at-least-once).
Паттерны потоковой передачи в BigQuery:
- Предпочитайте BigQueryIO с Storage Write API для высокой пропускной способности, низкой задержки и семантики «ровно один раз» (exactly-once) в рамках потока за счет смещений потока и автоматических повторных попыток.
- Для простых конвейеров с низкой скоростью потоковая вставка (streaming inserts) является приемлемым вариантом; установите insertId для дедупликации повторных попыток клиента.
- Запросы к потоковым буферам являются согласованными в конечном счете; для аналитики, критичной ко времени, выполняйте запросы после задержки буфера (например, подождите ~2x от наблюдаемой задержки доступности) или материализуйте данные через окна микро-батчей и режим committed в Storage Write API.
Эффекты «ровно один раз», идемпотентность, повторное воспроизведение и приемники:
- Beam гарантирует обработку «как минимум один раз»; семантика «ровно один раз» должна быть достигнута на стороне приемника с помощью идемпотентных записей, транзакций или ключей дедупликации.
- BigQuery: используйте потоки по умолчанию (default streams) или зафиксированные потоки (committed streams) в Storage Write API для семантики «ровно один раз» в рамках потока; при потоковой вставке установите стабильный insertId.
- Файлы: записывайте временные файлы с уникальными именами, финализируйте их по завершении окна и обеспечивайте атомарное переименование; избегайте перезаписи, чтобы предотвратить частичные дубликаты.
- Внешние базы данных: используйте операции upsert с ключом на основе стабильного идентификатора или реализуйте окна дедупликации.
- Проектируйте с учетом возможности повторного воспроизведения: поддерживайте детерминированные преобразования; убедитесь, что приемники выполняют дедупликацию при повторных попытках.
Обработка нежелательных сообщений, маршрутизация ошибок, наблюдаемость:
- Оборачивайте рискованные операции парсинга/обогащения в блок try/catch внутри ParDo и отправляйте сбои в PCollection для нежелательных сообщений (dead-letter) через TupleTag; включайте полезную нагрузку, код ошибки и контекст.
- Направляйте очереди нежелательных сообщений (DLQ) в BigQuery или Cloud Storage для анализа; рассмотрите возможность использования отдельного топика Pub/Sub для повторной обработки.
- Наблюдаемость: используйте метрики заданий Dataflow (отставание водяного знака, системное отставание, пропускная способность), пользовательские счетчики, метрики распределения и логи для каждого шага в Cloud Logging. Создайте оповещения об отставании и частоте ошибок в Cloud Monitoring. Используйте Error Reporting для агрегации исключений.
Паттерны для настройки производительности:
- Эффективное чтение: для источников BigQuery предпочитайте Storage Read API или чтение на основе запросов, которые выбирают только необходимые поля и применяют фильтры.
- Оптимизация слияния (Combine lifting): используйте комбинаторы (combiners) для уменьшения объема данных для shuffle перед GroupByKey.
- Дополнительные входы (Side inputs): кэшируйте небольшие справочные данные в памяти; следите за коэффициентом ветвления (fanout) и частотой обновлений.
- Сериализация: используйте компактные схемы (Avro/Proto) и избегайте избыточного парсинга JSON на «горячих» путях выполнения.
Развертывание, шаблоны и стратегии обновления
Flex Templates:
- Упаковка конвейеров в контейнеризированные, параметризованные шаблоны для воспроизводимых развертываний. Flex Templates поддерживают пользовательские зависимости, образы с GPU и изоляцию окружений.
- Вынесение параметров времени выполнения (например, входная подписка, выходная таблица, приемник для недоставленных сообщений, maxWorkers) для обеспечения развертываний под конкретное окружение.
Обновления и совместимость конвейеров:
- Dataflow поддерживает обновление «на месте» (in-place) для многих потоковых конвейеров, если имена преобразований, спецификации состояний и типы выходных данных остаются совместимыми. Используйте стабильные имена PTransform.
- При несовместимых изменениях графа или состояния выполните контролируемое переключение: запустите новое задание, затем выполните «осушение» (drain) старого, чтобы завершить обработку текущих данных и прекратить чтение новых элементов.
Осушение (draining) и снимки (snapshots):
- Осушение (drain) корректно завершает обработку, записывает оставшиеся выходные данные и прекращает работу; координируйте его с хранением данных в Pub/Sub или снимками, чтобы избежать пробелов.
- Для обеспечения непрерывности можно создать снимок Pub/Sub, запустить новый конвейер с чтением из снимка или с соответствующей временной метки, проверить выходные данные, а затем осушить старое задание.
Примеры конфигурации:
- Пример оконной обработки с ранними/поздними триггерами и накоплением:
undefined
- Пример BigQueryIO с использованием Storage Write API:
undefined
- Распространенные ошибки:
- Запись в файловые приемники в потоковом режиме без использования оконной записи может привести к зависанию финализации; включите оконную запись и триггеры.
- Неограниченный рост: если забыть ограничить состояние или допустимое опоздание (allowed lateness), это может привести к утечкам памяти и сбоям масштабирования.
- Отсутствие временных меток: если не назначать временные метки событий, конвейер по умолчанию переключается на время обработки, что приводит к потере корректности при переменных задержках.
Практический сценарий
Компания NovaTrack Inc. принимает глобальную телеметрию с 50 000 датчиков температуры и должна предоставлять поминутные агрегаты, сохранять необработанные данные и отображать их на дашборде в реальном времени. Ожидаются периодические некорректные сообщения и доставка с нарушением порядка. Решение должно автоматически масштабироваться, предоставлять некорректные записи для анализа и поддерживать обновления без простоя.
Подход:
Прием данных и семантика времени
- Создать региональный топик Pub/Sub и издателей для каждого региона с атрибутами deviceId и eventTs (RFC3339). По возможности включить ключи упорядочивания (ordering keys) по deviceId.
- Обоснование: Pub/Sub обеспечивает надежный, эластичный прием данных с доставкой «как минимум один раз» (at-least-once). Прикрепление временных меток события на входе сохраняет истинное время события; упорядочивание по устройствам уменьшает нарушение порядка внутри потока от одного устройства без создания централизованных узких мест.
Потоковый конвейер Dataflow с окнами по времени события
- Читать из выделенной подписки через PubSubIO, извлекая eventTs в качестве временной метки Beam, с откатом на publishTime в случае отсутствия.
- Применить FixedWindows размером в 1 минуту с ранним триггером через 30 секунд и поздними срабатываниями на каждый опоздавший элемент; установить допустимое опоздание (allowed lateness) в 10 минут и накапливать панели (accumulating panes).
- Обоснование: Окна по времени события обеспечивают точные поминутные агрегаты; ранние срабатывания поставляют данные на дашборд со свежестью менее минуты; поздние срабатывания корректируют агрегаты по мере поступления запоздавших данных. Ограничение на опоздание лимитирует размер состояния и затраты.
Валидация, обогащение и маршрутизация в очередь недоставленных сообщений (dead-letter)
- Реализовать ParDo, который парсит JSON, проверяет схему и диапазоны значений, и обогащает данные небольшим статическим справочником через дополнительный вход (side input), загружаемый из BigQuery при запуске задания.
- Использовать TupleTags для отправки валидных записей в основной выходной поток, а ошибочных — в PCollection для недоставленных сообщений (dead-letter), содержащую исходные данные, ошибку, deviceId и временную метку парсинга; записывать DLQ в партиционированную таблицу BigQuery.
- Обоснование: Дополнительные входы (side inputs) хранят справочные данные в памяти для низкой задержки. Сбор недоставленных сообщений позволяет анализировать и целенаправленно повторно обрабатывать некорректные строки, не блокируя основной поток.
Агрегация и предотвращение «горячих» ключей (hot-key)
- Ключировать по deviceId и вычислять поминутные avg/min/max с помощью CombineFns. Для метрик топ-N по регионам, шардировать по region#N, чтобы избежать «горячих» ключей, а затем повторно агрегировать.
- Обоснование: Комбайнеры (Combiners) минимизируют объем перемешиваемых данных (shuffle) и затраты; шардирование ключей предотвращает узкие места на одном ключе при сборе данных по регионам.
Приемники (sinks) и семантика «ровно один раз» (exactly-once)
- Записывать необработанные проверенные события и поминутные агрегаты в BigQuery с помощью BigQueryIO и Storage Write API. Установить стабильный ID вставки (insert id) на основе deviceId + eventTs для идемпотентности при любых пользовательских повторных попытках.
- Обоснование: Storage Write API обеспечивает высокопроизводительный прием данных с низкой задержкой и семантикой «ровно один раз» в рамках одного потока. Стабильные ID обеспечивают дедупликацию на стороне получателя в случае повторных отправок.
Стратегия консистентности для дашборда
- Дашборд запрашивает партиционированные таблицы с агрегатами с окном просмотра в 2 минуты относительно водяного знака или с фиксированной задержкой в 2 раза превышающей наблюдаемую задержку доступности для потоковых данных.
- Обоснование: Видимость потоковых данных в BigQuery является согласованной в конечном счете (eventually consistent); небольшая задержка чтения предотвращает пропуск строк, находящихся в обработке, сохраняя при этом поведение, близкое к реальному времени.
Эксплуатация: автомасштабирование и Streaming Engine
- Включить Streaming Engine; установить maxWorkers на основе ожидаемого пика (например, 3x от среднего), выбрать тип машины, подходящий для парсинга и шифрования (CPU-bound), и увеличить загрузочный диск для временных данных shuffle.
- Мониторить отставание водяного знака (watermark lag), объем необработанных данных в секундах (backlog seconds), CPU и пропускную способность для каждого шага; настроить оповещения на устойчивое отставание и всплески в скорости поступления сообщений в DLQ.
- Обоснование: Streaming Engine выносит состояние/shuffle за пределы воркеров, что обеспечивает эластичность и упрощает обновления; правильный подбор ресурсов и мониторинг предотвращают скрытое нарушение SLO.
Развертывание и обновления с помощью Flex Templates
- Упаковать конвейер как Flex Template с параметрами: входная подписка, выходные таблицы, таблица DLQ, maxWorkers и регион. При несовместимом изменении запустить новый конвейер, нацеленный на тот же топик, но с новой подпиской, проверить выходные данные, а затем осушить (drain) старое задание. Опционально можно создать снимок Pub/Sub и настроить чтение новой подписки из этого снимка, чтобы гарантировать отсутствие пробелов в данных.
- Обоснование: Flex Templates обеспечивают повторяемые, параметризованные развертывания. Проверенное «сине-зеленое» (blue/green) переключение с осушением (drain) позволяет достичь нулевой потери данных и минимального простоя.
Повторная обработка и пакетное заполнение исторических данных
- Хранить сжатые 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.
Сдайте экзамен →