Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
На этой странице объясняется, как вы используете конвейеры Lakeflow на протяжении всего срока службы конвейера данных — от первых проектных решений до масштабной работы и компромиссов, стоящих за каждым этапом. Каждый раздел ссылается на статьи, которые объясняют, как это сделать.
Это руководство предполагает знакомство с основными концепциями инженерии данных. Если вы новичок в конвейерах, начните с Apache Spark Declarative Pipelines, чтобы узнать, что это за продукт и декларативная модель, а затем пройдитесь по учебнику: Постройте ETL-пайплайн с использованием захвата изменений.
Обзор жизненного цикла конвейера
Конвейер проходит через шесть этапов:
- Планируйте и проектируйте: Определите, что вы создаёте, и выберите инструменты, язык и вычисления, которые подходят.
- Вводите данные: надёжно и постепенно вводите исходные данные в конвейер.
- Трансформируйте и моделируйте: очищайте, проверяйте, объединяйте и превращайте данные в таблицы, которым пользователи могут доверять.
- Ввод в эксплуатацию: Поместите конвейер в систему контроля версий, протестируйте его, настройте его выполнение по расписанию и перемещайте его между средами.
- Запуск в продакшене: мониторинг, оповещение, отладка, заполнение, безопасность и отслеживание линии по мере работы конвейера без присмотра.
- Зрелость и масштабирование: Подтвердите готовность к производству и поддерживайте здоровье конвейера по мере роста объёмов и размера команды.
Уровни не строго последовательны, но соответствуют порядку возникновения вопросов. Поскольку конвейеры Lakeflow занимаются оркестрацией, контрольными точками, повторениями и инкрементальной обработкой, ваша работа на каждом этапе в основном является проектным решением, а не реализацией.
Планирование и проектирование
Ваши первые решения формируют всё впоследствии. Для того, как декларативная модель сравнивается с написанием процедурных шагов самостоятельно, см. раздел «Процедурная и декларативная обработка данных» в Azure Databricks.
Несколько вариантов для настройки старта:
- Отдельный набор данных или конвейер. Одно материализованное представление или одна потоковая таблица могут быть определены с помощью SQL как автономный набор данных, а Azure Databricks управляет связанным с ними конвейером обновления. Создавайте конвейер Lakeflow и управляйте им как единым целым, когда вам нужны разработка на Python, приемники или многоэтапная оркестрация. См. Автономные трубопроводы против Lakeflow.
- SQL или Python (или оба варианта). SQL подходит для трансформаций, которые в основном состоят из фильтров, объединений и агрегирований. Python подходит для пользовательской логики, внешних библиотек или для программного генерирования множества похожих таблиц. Выбор делается для каждого файла, а не для всего конвейера, так что можно комбинировать оба варианта и не нужно определяться заранее.
- Serverless или классическая модель вычислений. Serverless рекомендуется по умолчанию и убирает конфигурацию кластера. Выбирайте Classic, если нужны конкретные типы экземпляров, пользовательские политики кластера или скрипт для инициатора. См. Настройка бессерверного конвейера и Настройка классических вычислительных ресурсов для конвейеров.
- Триггерное или непрерывное выполнение. Запуск по триггеру, так как он потребляет вычислительные ресурсы только во время выполнения. Непрерывный режим позволяет вычислениям работать для обработки новых данных с минимальной задержкой, что обычно является самым высоким фактором затрат, поэтому резервируйте его для подтвержденного требования к задержке. См. раздел "Триггерный и непрерывный режимы конвейера".
Конвейер выводит граф выполнения из наборов данных, на которые ссылается ваш код, поэтому проектирование в основном сводится к именованию и секвенированию наборов данных. Основное решение заключается в том, каким типом должен быть каждый выходной набор данных: потоковые таблицы для инкрементальных данных с частым добавлением записей или материализованные представления для пересчитываемых агрегатов и соединений. Этот выбор определяет затраты и корректность, потому что инкрементальная обработка масштабируется со скоростью поступления новых данных, тогда как полный пересчёт масштабируется со всем объёмом исторических данных. Для того, какой тип подходит для какой задачи, см. Что такое пайплайны?
Поскольку конвейерный код — это обычный Python и SQL, вы можете писать, использовать и проверять его в собственном редакторе перед развертыванием в общем рабочем пространстве.
На этом этапе
Вопросы, над которыми стоит подумать на этом этапе:
- Как выбрать между отдельным набором данных и полным конвейером?
- Как определить источники данных и понять, как к ним подключиться?
- Как спроектировать архитектуру моего пайплайна, прежде чем писать какой-либо код?
- Как выбрать формат файла и слой хранения?
- Как настроить локальную среду разработки?
- Как спланировать масштаб и оценку стоимости до начала строительства?
Прием данных
Ключевой вопрос при проектировании заключается в том, является ли источник данных доступным только для добавления или изменяется на месте. Это определяет, как вы моделируете цель:
- Источники только для добавления, такие как файлы, поступающие в облачное хранилище, или события на шине сообщений, поступают в потоковую таблицу, которая сохраняет контрольные точки своего прогресса, чтобы перезапуск не обрабатывал данные повторно и не приводил к их потере. Auto Loader обрабатывает файлы, открывает новые и делает выводы и развивает схемы по мере их поступления. Шины сообщений, такие как Apache Kafka, Центры событий Azure, Amazon Kinesis и Google Pub/Sub, считываются напрямую в потоковую таблицу. Дедупликация вниз по потоку, поскольку шина может передавать одно и то же событие более одного раза. Что касается Центры событий Azure, см. Использование Центры событий Azure в качестве источника данных для конвейера.
-
Источники, которые обновляют и удаляют строки, такие как большинство баз данных и многие системы программного обеспечения как услуги (SaaS), используют систему захвата изменений данных (CDC). Полное копирование при каждом запуске неэкономично и выполняется всё медленнее по мере роста объёма исходных данных, поэтому CDC считывает только те записи, которые изменились с момента последнего запуска. API
AUTO CDCприменяет эти изменения без написанной вручную логики слияния; см. API AUTO CDC: упрощение отслеживания изменений данных с помощью конвейеров. Поток применяет CDC к потоковой таблице, и несколько потоков могут направлять данные в одну таблицу — именно так несколько источников объединяются в одну целевую таблицу.
Создание контрольных точек и повторные попытки выполняются автоматически, поэтому конвейер продолжает работу с последнего обработанного смещения, а не обрабатывает всё заново. Две меры защиты нужно включить вручную:
- Столбец спасённых данных фиксирует записи, не соответствующие ожидаемой схеме.
- Ожидания применяют действие на уровне ряда, которое вы определяете.
Если контрольная точка потоковой обработки становится недействительной, выбирайте наименее затратный способ восстановления, который позволяет сохранить данные в таблице.
На этом этапе
Вопросы, над которыми стоит подумать на этом этапе:
- Как принимать данные из базы данных и выбирать между полной нагрузкой и CDC?
- Как принимать данные из API?
- Как загружать потоковые данные или данные о событиях?
- Как надёжно вводить файлы?
- Как мне справляться с сбоями при поглощении, не теряя данные?
Преобразование и моделирование
Трансформация превращает поглощённые данные в чистые таблицы, которым люди и инструменты могут доверять. Именно здесь узор медальона (от бронзы к серебру, затем к золоту) приобретает конкретную форму.
На первом месте — чистка и подтверждение. Ожидания — это встроенная функция конвейеров Lakeflow: ограничения качества данных, которые конвейер проверяет для каждой строки при каждом запуске, с подсчётом количества успешных и неуспешных проверок, поэтому контроль качества выполняется непрерывно, а не только как разовая проверка. Решите, что должно происходить, если строка не проходит проверку (выдать предупреждение и оставить её, удалить её или считать обновление неуспешным), и где должна выполняться эта проверка. Ворота обычно расположены на границе бронзы и серебра, так что всему, что ниже по течению, можно доверять без повторной проверки.
Объединение и агрегирование формируют этап перехода от серебра к золоту. Материализированный вид подходит для пакетного объединения или агрегации по существующим таблицам, потому что он сохраняет согласованность результатов с исходными данными: он обновляется постепенно, когда запрос и источники позволяют, и в противном случае полностью пересчитывается, давая тот же результат. Это делает его правильным выбором, когда корректность важнее, чем задержка, поскольку он пересчитывает соединения при изменении размерности. См. Как обновляются конвейеры? Присоединение к прямым трансляциям увеличивает состояние неограниченности, поэтому для соединений и агрегирований требуется водяной знак, чтобы ограничить, сколько времени конвейер ожидает поздно поступающих данных.
Две идеи правильности проходят через этот этап:
-
Идемпотентность означает, что конвейер даёт один и тот же результат, сколько бы раз он ни выполнялся на одном и том же входе. Пайплайны Lakeflow идемпотентны в отношении компонентов, которыми они управляют, таких как чтение с контрольными точками и апсерты по ключу
AUTO CDC; вы обеспечиваете идемпотентность собственной логики, избегая недетерминированных функций в повторно вычисляемых представлениях. - Обработка как минимум один раз и ровно один раз. Управляемые таблицы Delta-to-Delta фиксируют входные и выходные данные каждой микропартии вместе, обеспечивая семантику exactly-once по умолчанию. Это не распространяется на пограничные случаи, например на пользовательский приёмник, целевой объект, не являющийся Delta, или непроверенный пользовательский источник, где операцию записи следует рассматривать как at-least-once и делать её идемпотентной, например выполняя upsert по ключу.
Здесь также поддерживаются медленно меняющиеся размерности (SCD): AUTO CDC напрямую реализует SCD типа 1 и типа 2, так что вам достаточно выбрать тип, а не писать логику для отслеживания истории.
На этом этапе
Вопросы, над которыми стоит подумать на этом этапе:
- Как мне очистить и проверить входящие данные?
- Как отслеживать историю изменений с течением времени с помощью медленно меняющихся измерений (SCD)?Что такое SCD?
- Как объединить потоковые и статические данные?Как эффективно агрегировать данные?
- Как смоделировать свои данные для дальнейшего использования?
- Как обеспечить гарантии обработки в трубопроводах Lakeflow?
- Обработка как минимум один раз и обработка ровно один раз: в чём разница и какой вариант мне нужен?
- Как обрабатывать данные, поступающие с задержкой или не по порядку?
Ввод в эксплуатацию
Операционализация переносит конвейер из чего-то, что работает за вас, в то, что команда может построить, протестировать и многократно отправлять. Конвейер — это исходный код плюс конфигурация, поэтому применяются обычные практики программной инженерии.
Тестирование охватывает сразу две вещи: логику трансформации и постоянное качество данных, проходящих через неё. Объекты ожидания непрерывно обрабатывают сторону данных. Для логики вынесите преобразования в обычные функции и тестируйте их модульными тестами вне среды выполнения, затем проверьте граф конвейера с помощью пробного запуска, прежде чем что-либо материализовать. См . модульное тестирование конвейеров.
Храните код конвейера в Git и упаковывайте его для развертывания, чтобы его можно было регулярно проверять, отменять и развёртывать в разных средах. Этот пакет не является альтернативой трубопроводам Lakeflow. Это проектная и CI/CD-обвязка вокруг вашего конвейера, при этом логика работы с данными остаётся декларативной. Параметризуйте значения, специфичные для окружающей среды, такие как имена каталогов и пути, чтобы один и тот же код работал без изменений в каждой среде. См. раздел "Использование параметров с конвейерами".
Чтобы запускать конвейер по расписанию, включите его в Запуск конвейеров в рабочем процессе: Databricks рекомендует планировать и оркестрировать конвейеры с помощью заданий, что также позволяет координировать конвейер с другими задачами, например добавлять в цепочку последующий отчёт или несколько конвейеров. В рамках одного запуска конвейер упорядочивает и параллелизирует собственные наборы данных, поэтому оркестрация координирует только задачи вне конвейера.
На этом этапе
Вопросы, над которыми стоит подумать на этом этапе:
- Как протестировать конвейер данных и чем это отличается от тестирования обычного программного обеспечения?
- Как управлять версиями и сотрудничать в команде над конвейерным кодом?
- Как мне запланировать или организовать автоматический запуск моего конвейера?
- Как безопасно перенести мой пайплайн от разработки к стадированию и затем в продакшен?
- Как настроить CI/CD для моего конвейера?
Запуск в промышленной среде
Когда конвейер начинает работать в автономном режиме с реальными данными, работа заключается в том, чтобы понимать, исправно ли он работает, и устранять проблемы, если нет.
Мониторинг осуществляется на трёх уровнях глубины. Список Jobs & Pipelines позволяет сразу увидеть статус недавних запусков. Интерфейс мониторинга конвейера отображает все таблицы и потоки с цветовой маркировкой по статусу, с количеством строк, метриками качества данных и метриками отставания для потоковых таблиц. Журнал событий под обоими — источник истины для всего, что связано с программой или историей. Настройте уведомления о сбоях, чтобы узнать о неисправном запуске до того, как заинтересованные стороны об этом сообщат. Для обзора мониторинговых поверхностей см. раздел Мониторинг конвейеров.
Отлаживайте, двигаясь в обратном направлении от сбоя, отмеченного на графике, к полной информации об ошибке в журнале событий, а затем повторно запустите только то, что завершилось сбоем. Поведение повторных попыток различается в зависимости от триггера: ручные обновления отключают автоматические повторные попытки, чтобы сразу видеть ошибки, а запланированные обновления повторяют восстановительные неудачи. Таким образом, оповещение в продакшене может исчезнуть после повторной попытки, тогда как при интерактивной разработке тот же сбой не исчезнет. Во время разработки Genie Code помогает диагностировать и исправлять ошибки на уровне кода по мере итераций, хотя сегодня он ориентирован на создание конвейеров, а не на диагностику производственных запусков.
Моделируйте обратную загрузку как отдельный, явно заданный одноразовый поток, направленный в тот же целевой объект, что и ваш обычный инкрементальный поток. Если хранить это отдельно, можно зафиксировать, когда и как была загружена история, и сохранить логику установившегося состояния простой.
Обеспечьте безопасность конвейера, контролируя, кто может им управлять, запуская его как отдельный сервис, а не личный аккаунт, и храня учетные данные в секретном формате, а не в исходном коде. Происхождение данных отслеживается автоматически, вплоть до уровня столбцов. Конвейер записывает данные во внешнюю систему через приемники в конвейерах Lakeflow — это та граница, где применяется описанный выше принцип «как минимум один раз».
На этом этапе
Вопросы, над которыми стоит подумать на этом этапе:
- Как мне отслеживать, что мой конвейер работал успешно?
- Как мне получать уведомления, если что-то выходит из строя?
- Как отладить неудачный запуск конвейера?
- Как заполнить исторические данные?
- Как контролировать и прогнозировать стоимость запуска своего конвейера?
- Как защитить свой конвейер, включая учетные данные, контроль доступа и персональные данные?
- Как мне документировать свой конвейер и отслеживать происхождение данных?
Развивайтесь и масштабируйтесь
Зрелый конвейер работает без присмотра и растёт без переписывания. Подтверждение готовности и планирование масштабирования определяют этот этап.
Готовность к производству — это контрольный список по качеству данных, надёжности, наблюдаемости, развертыванию, стоимости и управлению. Относитесь к каждому неотмеченному пункту как к известному пробелу: определены ли ожидаемые значения для каждого набора данных, который может получать некачественные данные, запланирован ли конвейер, а не запускается ли он вручную, настроены ли уведомления о сбоях, выполняется ли он от имени сервисного субъекта, развёрнут ли он из системы контроля версий как минимум в среды dev и prod. Контроль качества данных и уведомления дешевле всего внедрить, и именно они с наибольшей вероятностью помогут выявить незамеченный проблемный запуск.
Масштабируйте в ответ на конкретные сигналы того, что состояние конвейера ухудшается:
- Длительность обновления растёт.
- Автомасштабирование раз за разом упирается в свой предел.
- Стоимость растет быстрее, чем базовый бизнес.
- Материализованные представления переключаются на полный пересчёт.
Сначала попробуйте вычислительные рычаги, например, перейти на серверлесс или подстроить режим производительности под ваши требования к задержке. Помимо этого, важнее всего то, как вы организуете наборы данных между конвейерами:
- У конвейера есть ограничение на параллельность: он обновляет только определённое количество наборов данных одновременно. Когда в конвейере больше наборов данных, чем это ограничение, дополнительные обновления ждут в очереди, поэтому общее время обновления конвейера увеличивается.
- Группировать связанные наборы данных и разделять несвязанные наборы. Группируйте по домену, общей частоте обновления и зависимостям; разделяйте на границах ответственности, слоёв и латентности. Например, отделение загрузки данных от трансформации не позволяет медленной загрузке задерживать все последующие этапы и позволяет каждому конвейеру оставаться достаточно небольшим, чтобы не превышать ограничение параллелизма.
Объединить два небольших конвейера позже проще, чем разделить один крупный, уже находящийся в производстве. Для того, как группировать и разделять наборы данных, см . раздел «Организация наборов данных между конвейерами Lakeflow».
На этом этапе
Вопросы, над которыми стоит подумать на этом этапе:
- Как узнать, что мой пайплайн готов к производству?
- Как масштабировать свой пайплайн по мере роста объема данных?
- Как организовать наборы данных между конвейерами?