Google PDE: Прием, интеграция и миграция данных — Руководство по подготовке
Часть Google Professional Data Engineer — Руководство по подготовке. Практикуйтесь с проверенными ответами в центре экзаменов Google, или пройдите тесты на время на ExamRoll.io.
Обзор
Приём, интеграция и миграция данных в Google Cloud охватывают повторяемые шаблоны, управляемые сервисы и средства операционного контроля, которые превращают разнообразные исходные системы в надёжные, доступные для запросов наборы данных. Эффективные архитектуры отделяют транспортировку от преобразования, разделяют производителей и потребителей и отдают предпочтение идемпотентным конвейерам с контрольными точками, чётким отслеживанием происхождения данных и верификацией. В этом разделе рассматриваются шаблоны приёма данных, сервисы Google Cloud для перемещения данных и CDC, контроль схем и качества данных, возможности подключения и гибридная интеграция, а также стратегии переключения. На протяжении всего раздела освещаются компромиссы проектирования и режимы отказа.
Шаблоны и рабочие нагрузки для приёма данных
- Пакетный приём данных: Периодические извлечения или загрузки файлов через определённые интервалы. Подходит для предсказуемых затрат и дозагрузки исторических данных. Режим отказа: большие, редкие пакеты вызывают всплески потребления ресурсов, длительные окна для наверстывания и срыв SLA. Меры по смягчению: подбирайте правильный размер окон для пакетов, сегментируйте по времени или ключу и используйте параллелизм.
- Массовая загрузка: Единовременные или крупномасштабные загрузки (например, начальная загрузка исторических данных). Предпочтительны колоночные или самоописывающиеся форматы (Parquet, Avro) и прямая загрузка в аналитическое хранилище (BigQuery) или промежуточная в Cloud Storage. Компромисс: запросы к внешним таблицам позволяют избежать этапов загрузки, но переносят затраты на сканирование во время выполнения запроса.
- Инкрементальная загрузка: Периодическая загрузка дельт по временным меткам или меткам верхнего уровня. Требует надёжной дедупликации и идемпотентных операций upsert. Режим отказа: рассинхронизация часов или поздно поступающие записи. Используйте серверные временные метки коммитов и водяные знаки.
- Захват изменённых данных (CDC): Непрерывная репликация операций вставки, обновления и удаления из операционных баз данных. Лучше всего подходит для аналитики в режиме, близком к реальному времени, и миграций с минимальным временем простоя. Компромиссы:
- Упорядочивание: Большинство инструментов CDC сохраняют порядок в рамках транзакций и обычно в пределах одного шарда, но не гарантируют глобальный порядок между шардами. Используйте временные метки коммитов транзакций и первичные ключи для восстановления последовательности.
- Семантика доставки: Обычно используется семантика «как минимум один раз»; создавайте идемпотентные приёмники или выполняйте дедупликацию с помощью уникальных идентификаторов изменений.
- Снимок + CDC: Начните с консистентного снимка, затем примените изменения из точной последовательности логов, чтобы достичь паритета без простоя.
Реляционные, SaaS, локальные источники и файлы:
- Реляционные источники: Используйте нативные средства CDC или столбцы с временными метками. Для массовой загрузки экспортируйте данные в Avro/Parquet и размещайте их в Cloud Storage.
- SaaS-источники: Предпочитайте API поставщиков с инкрементальными токенами; интегрируйтесь через управляемые коннекторы (например, в Data Fusion). Регулируйте скорость для соблюдения лимитов и обрабатывайте дрейф схемы.
- Локальные источники: Выбирайте между передачей на основе агентов, VPN/Interconnect + Private Google Access или офлайн-загрузкой с помощью Transfer Appliance.
- Приём файлов: При наличии множества мелких файлов объединяйте их (например, с помощью tar), чтобы уменьшить накладные расходы на RPC. Используйте
gsutil -mили параллелизированные клиенты; объединяйте или преобразуйте их в более крупные колоночные файлы для аналитики.
Сервисы Google Cloud для приёма, интеграции и миграции данных
- Datastream (бессерверный CDC): Захватывает изменения из MySQL, PostgreSQL и Oracle и передаёт их в Cloud Storage, BigQuery (через шаблоны) или Pub/Sub. Он сохраняет границы транзакций и метаданные коммитов; глобальный порядок не гарантируется. Применяйте упорядочивание на последующих этапах по ключу и временной метке коммита. Ожидайте семантику доставки «как минимум один раз»; проектируйте идемпотентные потребители (например, BigQuery MERGE с идентификаторами изменений).
- Database Migration Service (DMS): Для миграции баз данных с минимальным временем простоя с использованием нативной репликации. DMS создаёт консистентный снимок, а затем непрерывно реплицирует изменения, используя GTID/LSN/SCN. Он специально создан для миграций по типу lift-and-shift, а не для произвольных преобразований. Для аналитики при необходимости дополняйте DMS с помощью Dataflow или Data Fusion.
- Cloud Data Fusion: Управляемый сервис интеграции с коннекторами к реляционным, SaaS, файловым и системам обмена сообщениями. Создавайте конвейеры с этапами преобразования (соединения, агрегации, конвертация форматов, пользовательские рецепты Wrangler) и отслеживайте происхождение данных по источникам и полям. В операционном плане он выполняет планирование, повторные попытки и отправляет метрики. Используйте Data Fusion для ELT/ETL без кода или с минимальным его использованием и для централизованного управления коннекторами.
- Storage Transfer Service (STS): Управляемая передача данных по расписанию из AWS S3, Azure Blob, локальных сред (с помощью агентов), SFTP и списков URL в Cloud Storage. Поддерживает манифесты, инкрементальную синхронизацию, контроль пропускной способности и проверку целостности по контрольным суммам. Режимы отказа включают неэффективность при работе с мелкими файлами и троттлинг API; смягчайте последствия с помощью пакетирования и настраиваемого параллелизма.
- Transfer Appliance: Автономное зашифрованное устройство для начальной загрузки данных объёмом от нескольких терабайт до петабайта, когда пропускная способность сети ограничена или данные слишком чувствительны для длительной передачи. Встроены функции цепочки ответственности и шифрования. После начальной загрузки используйте STS или CDC для передачи дельт.
- Cloud Pub/Sub + Dataflow: Pub/Sub разделяет производителей и потребителей для потоковых или микропакетных шаблонов. Dataflow предлагает автомасштабируемую обработку потоков/пакетов с сохранением состояния, контрольными точками и водяными знаками. Используйте BigQuery Storage Write API для потоковой передачи с низкой задержкой и гарантиями «ровно один раз» для потока по умолчанию; в противном случае полагайтесь на семантику дедупликации по
insertId.
При миграции с Hadoop на Dataproc минимизируйте использование Persistent Disk, храня данные в Cloud Storage с помощью коннектора GCS, и используйте эфемерные или автомасштабируемые кластеры. Это позволяет избежать больших затрат на блочное хранилище, сохраняя при этом семантику, совместимую с HDFS, для обработки данных.
Схема, валидация и качество данных на границе
- Сопоставление схем и преобразование типов: стандартизируйте данные с использованием строго типизированных схем на ранних этапах. Avro или Parquet сохраняют схему и обеспечивают её простое развитие. В BigQuery отдавайте предпочтение секционированным и кластеризованным таблицам для снижения стоимости сканирования. Пример: создание секционированной таблицы для ежедневного анализа
undefined
- Обработка некорректных записей: направляйте отклонённые записи в очередь недоставленных сообщений (dead-letter queue) в Pub/Sub или в карантинный бакет в Cloud Storage. Используйте побочные выводы (side outputs) в Dataflow или сборщики ошибок (error collectors) в Data Fusion. Логируйте ошибки парсинга с примерами данных (payloads) и версиями схем для анализа.
- Валидация: выполняйте проверки на границе системы перед сохранением данных:
- Структурная: соответствие схеме, наличие обязательных полей, типы данных, домены перечислений (enum).
- Ссылочная: существование внешних ключей через поиск в кэшированных таблицах измерений.
- Разумность: диапазоны для временных меток, геозоны, неотрицательные значения.
- Уникальность: коллизии первичных или составных ключей.
- Идемпотентная загрузка: используйте детерминированные ключи и операции upsert. В BigQuery реализуйте операцию MERGE с использованием естественного или суррогатного ключа изменений. Пример:
undefined
- Водяные знаки (watermarks) и обработка опоздавших данных: в потоковых конвейерах настраивайте водяные знаки по времени события (event-time watermarks) и допустимое опоздание (allowed lateness) для баланса между полнотой данных и задержкой. Опоздавшие данные направляются по корректирующим маршрутам или инициируют процессы заполнения пропусков (backfills).
- Сверка данных: отслеживайте количество строк и контрольные суммы для каждой секции/окна от источника до приёмника. Фиксируйте позиции в логах CDC (LSN/SCN) и временные метки коммитов; сохраняйте их в контрольной таблице для подтверждения непрерывности и выявления пробелов.
Подключение, надёжность и эксплуатация
Сетевое подключение и частный доступ:
- Гибридная среда: используйте Cloud VPN или Dedicated/Partner Interconnect для частного подключения. Включите Private Google Access или Private Service Connect для частного доступа к API Google, таким как Cloud Storage.
- Безопасность: используйте сервисные аккаунты для идентификации рабочих нагрузок, IAM с минимальными привилегиями, VPC Service Controls для предотвращения утечки данных и CMEK, где это необходимо.
- Пропускная способность: масштабируйте параллелизм на стороне клиента, но в конечном счёте пропускная способность определяется полосой пропускания. Для массовой передачи данных предпочтительно использовать Transfer Appliance для начальной загрузки, а затем STS или CDC для инкрементальных обновлений.
Контрольные точки и противодавление (backpressure): Dataflow управляет контрольными точками и автомасштабированием; проектируйте приёмники, способные поглощать всплески нагрузки (буферизация в Cloud Storage, пакетная запись в BigQuery). Для Pub/Sub настраивайте управление потоком (flow control) и крайние сроки подтверждения (ack deadlines), чтобы предотвратить лавинообразные повторные доставки сообщений.
Упорядоченность и согласованность при использовании CDC:
- Datastream сохраняет порядок внутри транзакций и выдаёт метаданные коммитов; потребители восстанавливают порядок для каждого ключа, используя временные метки коммитов. Ожидайте семантику доставки «как минимум один раз» (at-least-once); обеспечьте идемпотентность.
- DMS обеспечивает согласованность базы данных при переключении со снимка на репликацию, используя нативные логи. Используйте реплики чтения или стратегии двойной записи для поэтапного переключения.
Стратегия работы с файлами для аналитики: для доступа из нескольких систем к большим объёмам данных храните канонические данные в Cloud Storage и, если это экономически целесообразно, предоставляйте к ним доступ через постоянные внешние таблицы для специальных запросов (ad hoc). Для производственной аналитики загружайте данные в секционированные таблицы BigQuery, чтобы минимизировать стоимость сканирования для каждого запроса.
Оптимизация для мелких файлов: объединяйте мелкие файлы (например, ~1000 в один tar-архив) перед передачей, а затем распаковывайте их в облаке. Используйте параллельное выполнение gsutil и правила жизненного цикла для переноса на другие уровни хранения и удаления промежуточных артефактов.
Эксплуатационные риски и способы их устранения:
- Дрейф схемы от SaaS-источников: включите эволюцию схемы в Data Fusion и обеспечьте совместимость. Настройте оповещения о несовместимых изменениях.
- Часовые пояса и кодировки: нормализуйте данные к UTC и UTF-8 на входе.
- Пробелы в CDC: отслеживайте срок хранения логов источника; настройте оповещения, когда отставание реплики приближается к пределу хранения.
- Квоты: потоковая вставка в BigQuery, лимиты на частоту вызовов API; переходите на пакетную обработку при приближении к лимитам.
Переключение, обратная загрузка и проверка
- Планирование переключения:
- Большой взрыв (Big bang): короткая заморозка, одномоментное переключение. Самая низкая операционная сложность; самый высокий риск, если потребуется откат.
- Поэтапный или сине-зеленый (blue/green): параллельная работа с зеркалированием записей, постепенным переключением трафика и теневыми чтениями. Более высокая стоимость; более безопасный откат.
- Обратная загрузка:
- Выполните начальную массовую загрузку (Transfer Appliance или STS) с использованием Avro/Parquet для сохранения схемы. Секционируйте и кластеризуйте данные во время загрузки, чтобы избежать переделок.
- Запустите CDC с известной позиции в логе одновременно со созданием снимка, чтобы захватить изменения (дельты) во время массовой передачи. Выполните сверку по общей временной метке (watermark) перед открытием для производственной эксплуатации.
- Откат:
- Поддерживайте унаследованную систему в режиме только для чтения во время проверки. В сценариях с двойной записью управляйте записью через функциональный флаг (feature flag) для быстрого отката. Сохраняйте согласованную контрольную точку для повторного применения или отмены изменений CDC при необходимости.
- Проверка миграции:
- Структурная: количество строк и контрольные суммы по разделам совпадают; схема и ограничения эквивалентны.
- Временная: нет пробелов от момента создания снимка до переключения; позиции CDC непрерывны.
- Бизнес-паритет: сравните агрегаты и KPI за определенные периоды; выполните приемочные запросы.
- Производительность: проверьте пропускную способность приема данных, задержку запросов и затраты на соответствие бюджетам.
Практический сценарий
Компания Northstar Retail должна консолидировать в Google Cloud глобальную смесь локальных транзакционных систем Oracle и MySQL, события из SaaS CRM и ежедневные выгрузки в формате CSV для аналитики и машинного обучения в режиме, близком к реальному времени. Им также необходимо перенести унаследованный кластер Hadoop, не неся больших затрат на блочное хранилище, и выполнить переключение с нулевым или минимальным временем простоя.
- Создание безопасного гибридного подключения
- Используйте Partner Interconnect для основной пропускной способности и Cloud VPN в качестве резервного канала. Включите Private Google Access, чтобы локальные рабочие нагрузки могли получать частный доступ к Cloud Storage и Pub/Sub. Обоснование: Частные пути минимизируют риски, связанные с исходящим трафиком, и задержку, а Private Google Access позволяет избежать использования публичных IP-адресов, соблюдая при этом политику безопасности.
- Эффективная начальная загрузка исторических данных
- Для 800 ТБ исторических данных HDFS скопируйте их в Cloud Storage с помощью Transfer Appliance (начальная массовая загрузка). После начальной загрузки ежедневно запускайте Storage Transfer Service для экспортированного локального ресурса NFS, чтобы забирать изменения до момента переключения. Обоснование: Transfer Appliance позволяет избежать длительной перегрузки сети; STS обеспечивает инкрементальную синхронизацию по расписанию с проверкой контрольных сумм. Хранение данных в Cloud Storage с использованием GCS connector позволяет обрабатывать их в Dataproc без необходимости выделять 50 ТБ Persistent Disk на каждый узел.
- Миграция операционных баз данных с помощью CDC
- Используйте DMS для миграции MySQL и PostgreSQL с минимальным временем простоя. Для CDC из Oracle в аналитическую систему используйте Datastream для выгрузки в промежуточную зону в Cloud Storage, а затем предоставленный Google шаблон Dataflow для загрузки в BigQuery. Обоснование: DMS использует нативные механизмы репликации для надежной комбинации снимка и непрерывной синхронизации; Datastream предоставляет бессерверный CDC с метаданными коммитов, а шаблон Dataflow обеспечивает упорядоченные, идемпотентные записи в BigQuery.
- Прием данных из SaaS и файловых источников
- Создайте конвейеры Cloud Data Fusion, используя коннекторы SaaS для событий CRM с инкрементальными токенами, и файловый конвейер для приема ежедневных CSV-файлов с SFTP-сервера поставщика через STS. Нормализуйте данные в формат Avro в специально подготовленном бакете Cloud Storage, а затем загрузите их в секционированные таблицы BigQuery. Обоснование: Cloud Data Fusion централизует коннекторы, преобразование и отслеживание происхождения данных (lineage). Стандартизация на Avro сохраняет схему и упрощает ее развитие; секционированные таблицы BigQuery снижают стоимость запросов.
- Потоковая передача событий в реальном времени
- Публикуйте события с веб-сайтов и из магазинов в Pub/Sub. Обрабатывайте их с помощью Dataflow для парсинга, валидации, обогащения и установки временных меток (watermarking); записывайте в BigQuery через Storage Write API и архивируйте необработанные данные в формате Avro в Cloud Storage. Обоснование: Pub/Sub разделяет производителей и потребителей; Dataflow обеспечивает автомасштабирование, обработку с сохранением состояния, контрольные точки и обработку запаздывающих данных; двойная запись гарантирует как аналитику с низкой задержкой, так и надежное хранение необработанных данных.
- Обеспечение контроля качества данных и схем на входе
- Реализуйте реестр схем и их валидацию в Dataflow/Data Fusion. Направляйте некорректные записи в карантинный бакет GCS и в топик недоставленных сообщений (dead-letter topic) Pub/Sub. Применяйте проверки доменных значений (например, коды валют, временные метки в UTC) и выполняйте дедупликацию с использованием составных ключей. Обоснование: Раннее отклонение и помещение в карантин предотвращают распространение некачественных данных; идемпотентность и дедупликация защищают от дубликатов при доставке по принципу «хотя бы один раз» (at-least-once) из источников CDC и потоковой передачи.
- Оптимизация хранения и доступа к аналитическим данным
- Загружайте подготовленные наборы данных в секционированные и кластеризованные таблицы BigQuery. Предоставьте доступ к архивам необработанных данных как к постоянным внешним таблицам для нечастого анализа. Для OLTP-нагрузок, которые остаются транзакционными, сохраните Cloud SQL с репликами чтения. Обоснование: Секционирование и кластеризация минимизируют стоимость сканирования; внешние таблицы позволяют избежать ненужных загрузок для редкого доступа; Cloud SQL сохраняет семантику ACID для транзакционных приложений.
- Планирование переключения, обратной загрузки и отката
- Выполните создание снимка + CDC для каждой СУБД; достигните точки сверки, где количество строк и контрольные суммы совпадают. Запустите сине-зеленое развертывание с двойной записью на 48 часов, постепенно переключая операции чтения на BigQuery. Поддерживайте функциональный флаг для отмены записей при обнаружении расхождений. Обоснование: Сине-зеленое развертывание снижает риски; проверка по известной временной метке (watermark) гарантирует полноту данных; флаги обеспечивают быстрый откат.
- Проверка и наблюдаемость
- Создайте контрольные таблицы, фиксирующие LSN/SCN источника, временные метки коммитов, количество строк и контрольные суммы по разделам. Отслеживайте задержку Datastream, состояние репликации DMS, временные метки (watermarks) Dataflow, очередь сообщений в Pub/Sub, статус заданий STS и метрики потоковых вставок в BigQuery. Обоснование: Сквозное отслеживание происхождения данных и количественный контроль обеспечивают проверяемое подтверждение корректности и своевременное оповещение о пробелах или задержках.
Разделяя уровни начальной загрузки (landing), подготовки (curation) и предоставления данных (serving); используя Cloud Storage в качестве надежного и недорогого промежуточного хранилища и архива; применяя DMS/Datastream для CDC с идемпотентными потребителями; и обеспечивая контроль схем и качества на входе, Northstar Retail достигает безопасного, масштабируемого приема данных и низкорисковой, проверяемой миграции с предсказуемыми затратами.
← 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.
Сдайте экзамен →