Google PDE: Обмен сообщениями, прием событий и сервисы реального времени — Руководство по подготовке
Часть Google Professional Data Engineer — Руководство по подготовке. Практикуйтесь с проверенными ответами в центре экзаменов Google, или пройдите тесты на время на ExamRoll.io.
Интеграция на основе событий, Eventarc и совместимость с Kafka
- Eventarc и CloudEvents
- Eventarc направляет события от сервисов Google Cloud (и пользовательских источников через Pub/Sub) в Cloud Run, GKE или Workflows, используя спецификацию CloudEvents. Атрибуты, такие как type, source, subject, обеспечивают точную фильтрацию и возможность аудита.
- Используйте фильтры по атрибутам для минимизации фанаута (fan-out) и снижения нагрузки на последующие системы.
- Доставка осуществляется по принципу «как минимум один раз» (at-least-once); делайте обработчики идемпотентными и, по возможности, не сохраняющими состояние (stateless).
- Пример триггера Eventarc:
gcloud eventarc triggers create gcs-finalize-to-run
–destination-run-service=ingestor
–event-filters=“type=google.cloud.storage.object.v1.finalized”
–event-filters=“bucket=my-data-bucket”
–service-account=eventarc-sa@PROJECT_ID.iam.gserviceaccount.com - Совместимость с Kafka и управляемая миграция
- Шаблоны Dataflow соединяют Kafka и Pub/Sub для поэтапной миграции. Зеркалируйте топики с сохранением ключей; сначала переключайте потребителей (consumers), затем производителей (producers), или используйте двойную запись на время перехода.
- Pub/Sub Lite предлагает потоковую передачу с партициями и выделенной пропускной способностью, маршрутизацией по ключу и более низкой стоимостью; он является региональным/зональным и подходит для нагрузок, подобных Kafka, где важны предсказуемая производительность и порядок сообщений в пределах партиции.
- Аспекты миграции:
- Порядок (Ordering): сопоставляйте ключи Kafka с ключами упорядочивания (ordering keys) Pub/Sub или партициями Lite.
- Смещения (Offsets): передавайте смещения как атрибуты сообщений для диагностики; после миграции потребители не могут полагаться на смещения Kafka.
- Доставка (Delivery): примите принцип «как минимум один раз» (at-least-once); обеспечьте идемпотентность на последующих этапах.
- Схемы (Schemas): мигрируйте определения из Confluent Schema Registry в схемы Pub/Sub или стандартизируйте на Protobuf/Avro с совместимыми правилами эволюции схем.
Шаблоны потоковой загрузки данных, пропускная способность, масштабирование, безопасность и эксплуатация
- Шаблоны загрузки в реальном времени
- Pub/Sub -> Dataflow -> BigQuery: используйте приемник BigQuery Storage Write API для высокой пропускной способности и идемпотентности со смещениями потока; направляйте сбои в таблицу недоставленных сообщений для анализа.
- Pub/Sub -> Dataflow -> Cloud Storage: архивируйте необработанные события для повторной обработки; используйте оконные, сжатые записи для баланса между стоимостью и задержкой.
- Pub/Sub -> операционные хранилища: записывайте в Bigtable для поиска с низкой задержкой, в Spanner для транзакций со строгой согласованностью или в Cloud SQL/Firestore в зависимости от потребностей рабочей нагрузки. Обеспечьте идемпотентные операции upsert с ключом по уникальному ID события.
- Доставка «как минимум один раз», предотвращение дубликатов и идемпотентность
- Передавайте уникальные event_id и event_time в каждом сообщении; используйте UUID, сгенерированные на стороне отправителя.
- Дедупликация при потоковой передаче в BigQuery: установите insertId или используйте Storage Write API с упорядоченными потоками; все равно защищайте запросы логикой дедупликации.
- Пример дедупликации на стороне запроса: WITH ranked AS ( SELECT t.*, ROW_NUMBER() OVER (PARTITION BY event_id ORDER BY event_time DESC) AS rn FROM dataset.events t ) SELECT * EXCEPT(rn) FROM ranked WHERE rn = 1;
- Для push-эндпоинтов возвращайте код 2xx только после успешной обработки; в противном случае ожидайте повторной доставки.
- Пропускная способность сообщений, квоты и масштабирование
- Издатели (Publishers): группируйте сообщения в пакеты и повторно используйте соединения; распараллеливайте нагрузку между несколькими клиентами. Используйте много ключей упорядочивания для масштабирования упорядоченных рабочих нагрузок.
- Подписчики (Subscribers): предпочитайте потоковую pull-подписку с управлением потоком (максимальное количество необработанных байт/сообщений). Устанавливайте дедлайны подтверждения в соответствии со временем обработки и продлевайте их при необходимости.
- Отслеживайте и запрашивайте увеличение квот на пропускную способность публикации и подписки по мере роста объемов; проектируйте с запасом (например, 2x от ожидаемого пика) для поглощения всплесков нагрузки.
- Согласованность и доступность
- Потоковая передача в BigQuery обеспечивает согласованность в конечном счете для видимости в запросах; для интерактивных запросов, которые должны включать потоковые строки, делайте задержку на основе наблюдаемой латентности (например, 2x от P50 задержки доступности) или проектируйте с использованием агрегаций, выровненных по водяным знакам (watermarks) в Dataflow, и запрашивайте материализованные результаты.
- Безопасность
- IAM: предоставляйте роли с минимальными привилегиями (pubsub.publisher для отправителей в топик; pubsub.subscriber для потребителей в подписку). Используйте выделенные сервисные аккаунты для каждой рабочей нагрузки.
- Аутентификация push-подписок: настройте push-подписки для прикрепления токенов OIDC от сервисного аккаунта; проверяйте аудиторию (audience) на эндпоинте. Предпочитайте приватные эндпоинты Cloud Run для встроенной аутентификации и TLS.
- Шифрование: Pub/Sub шифрует данные при передаче и хранении; используйте CMEK в топиках для ключей, управляемых клиентом. Применяйте VPC Service Controls для снижения риска эксфильтрации данных. При необходимости используйте шифрование на стороне клиента для конфиденциальных полей данных.
- Операционная диагностика отставания, повторных доставок и сбоев подписчиков
- Мониторинг с помощью Cloud Monitoring:
- subscription/num_undelivered_messages и oldest_unacked_message_age для отслеживания отставания.
- expired_ack_deadline_count для обнаружения пропущенных подтверждений, вызывающих дубликаты.
- publish_request_count и pull_request_count для отслеживания пропускной способности.
- Расследуйте пропажу событий на дашборде, воспроизводя известный набор данных через конвейер и сравнивая результаты на каждом этапе, чтобы изолировать неисправное преобразование или приемник.
- Для потоковой обработки в Dataflow:
- Используйте автомасштабирование с подходящим значением maxWorkers для поглощения нагрузки от множества источников.
- Осушайте (drain) конвейеры перед несовместимыми обновлениями, чтобы позволить завершиться текущей работе и предотвратить потерю данных.
- Для уведомлений о вставке в BigQuery, направляйте аудиторские записи Cloud Logging через приемник, отфильтрованный по конкретным таблицам, в топик Pub/Sub для оповещений.
- Мониторинг с помощью Cloud Monitoring:
Практический сценарий
Компании Contoso Freight нужна глобальная платформа для обработки событий в реальном времени для загрузки 10 000 телеметрических сообщений в минуту от грузовиков, обогащения событий, обеспечения интерактивной аналитики и запуска рабочих процессов при получении файлов от внешних партнеров. Некоторые CSV-файлы от партнеров содержат некорректные строки, и аналитическая команда должна иметь возможность проверять ошибки, не блокируя поток данных.
- Создание основного слоя обмена сообщениями и схемы
- Действие: Определите схему Avro для телеметрии и прикрепите ее к топику Pub/Sub
telemetryс принудительной проверкой схемы (require). Включите упорядочивание сообщений и публикуйте сordering_key = hash(device_id). - Обоснование: Принудительная проверка схемы на уровне топика позволяет на раннем этапе отбраковывать некорректно сформированные события. Упорядочивание по устройствам поддерживает последовательную обработку, когда это необходимо, а хеширование распределяет ключи для сохранения пропускной способности.
- Подготовка подписок с изоляцией и очередью недоставленных сообщений
- Действие: Создайте pull-подписку
telemetry-stream-subдля Dataflow с топиком для недоставленных сообщенийtelemetry-dltиmax_delivery_attempts=10. Добавьте подписку BigQuerytelemetry-raw-bqдля сохранения необработанных событий в таблицу, партиционированную по времени, для отслеживания происхождения данных и возможности повторной обработки. - Обоснование: DLQ (очередь недоставленных сообщений) изолирует проблемные сообщения для расследования. Отдельная подписка BigQuery обеспечивает простой в эксплуатации путь экспорта для сохранения необработанных событий независимо от конвейера обработки.
- Построение потокового конвейера Dataflow для обогащения и записи данных
- Действие: Загружайте данные из
telemetry-stream-subс помощью потоковой pull-подписки с управлением потоком. Проверяйте по схеме, обогащайте справочными данными и вычисляйте оконные агрегаты. Записывайте в BigQuery с помощью Storage Write API с именованным потоком иinsertId = event_id; записывайте необработанные резервные копии в Cloud Storage ежечасно; перенаправляйте некорректные/неудачные записи в таблицу BigQuery для недоставленных сообщений. - Обоснование: Storage Write API обеспечивает высокопроизводительную запись с низкой задержкой и идемпотентностью через
insertId/смещения потока. Таблица недоставленных сообщений позволяет анализировать их, не блокируя поток, а архивы в Cloud Storage дают возможность повторной обработки.
- Обработка дубликатов и согласованности в конечном счете в аналитике
- Действие: Для интерактивных запросов, которые должны исключать дубликаты, публикуйте
event_idиevent_timeв каждой записи и используйте представление для дедупликации: CREATE OR REPLACE VIEW analytics.latest_events AS SELECT * EXCEPT(rn) FROM ( SELECT e.*, ROW_NUMBER() OVER (PARTITION BY event_id ORDER BY event_time DESC) rn FROM analytics.events e ) WHERE rn = 1; Внедрите небольшую задержку в запросах на основе наблюдаемой доступности потоковой передачи BigQuery (например, удвоенная медианная задержка). - Обоснование: Доставка «как минимум один раз» требует идемпотентных записей и дедупликации на стороне запроса. Ожидание уменьшает количество пропущенных строк, находящихся в обработке, учитывая задержку видимости при потоковой передаче.
- Интеграция файлов от партнеров с помощью Eventarc
- Действие: Настройте Eventarc для маршрутизации событий
object.finalizedиз Cloud Storage для бакетаpartner-dropsв сервис Cloud Run, который запускает пакетное задание Dataflow для загрузки CSV в BigQuery, отправляя ошибки парсинга в таблицу недоставленных сообщений. - Обоснование: Eventarc обеспечивает оркестрацию на основе событий с фильтрацией CloudEvents по бакету и префиксу объекта. Пакетное задание Dataflow отделяет некорректные строки для анализа, оперативно загружая корректные данные.
- Обеспечение безопасности платформы
- Действие: Используйте разные сервисные аккаунты: отправители получают роль
pubsub.publisherдля топикаtelemetry; SA воркера Dataflow получаетpubsub.subscriberдля подпискиtelemetry-stream-subи права на запись в целевые датасеты BigQuery и Cloud Storage; триггер Eventarc использует выделенный SA с рольюinvokerдля Cloud Run. Включите CMEK для топикаtelemetryи датасетов BigQuery. Настройте push-эндпоинты, если они есть, с OIDC и проверкой аудитории (audience). - Обоснование: Принцип минимальных привилегий в IAM и CMEK отвечают требованиям безопасности и соответствия; аутентифицированная доставка предотвращает спуфинг.
- Надежная эксплуатация и масштабирование
- Действие: Настройте автомасштабирование Dataflow с большим значением
maxWorkersдля поглощения пиков. Отслеживайтеsubscription/oldest_unacked_message_ageиexpired_ack_deadline_count; настройте оповещения при превышении пороговых значений. Для изменений в конвейере, нарушающих совместимость, выполняйте развертывание с осушением (drain), чтобы избежать потери сообщений. Если отставание растет, увеличьте параллелизм подписчиков и продлите дедлайны подтверждения пропорционально времени обработки. - Обоснование: Проактивный мониторинг позволяет на ранней стадии обнаруживать отставание и повторные доставки. Автомасштабирование и настроенные дедлайны подтверждения предотвращают лавинообразное появление дубликатов. Осушение (draining) сохраняет сообщения в обработке во время обновлений.
Такая архитектура обеспечивает отказоустойчивую, безопасную и наблюдаемую потоковую загрузку данных в реальном времени с пакетной интеграцией на основе событий, поддерживает устойчивость к дубликатам и эволюцию схемы, а также предлагает быструю аналитику, изолируя некорректные данные для целенаправленного исправления.
← Потоковая обработка с Dataflow и Apache Beam · Все домены · Spark →
Отработать эти вопросы → · Тесты на время на 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.
Сдайте экзамен →