Завод данных Azure Трубопроводы для автоматизированных данных Рабочие процессы

Azure Data Factory (ADF) - это полностью управляемый облачный сервис интеграции данных Microsoft. Он позволяет организациям создавать, планировать и организовывать рабочие процессы данных в масштабе, перемещая и преобразуя данные в различных источниках и местах назначения. Два фундаментальных строительных блока ADF - это трубопроводы и триггеры. Трубопроводы определяют работу - последовательность действий, которые перемещают, обрабатывают или анализируют данные. Триггеры определяют, когда эта работа происходит - будь то по повторяющемуся графику, в ответ на внешние события или в течение конкретных временных окон. Вместе они образуют основу автоматизированных, готовых к производству конвейеров данных. Эта статья предоставляет подробное руководство по производству триггеров и трубопроводов Azure Data Factory, охватывая их типы, создание, лучшие практики и интеграцию с более широкой экосистемой Azure.

Обсуждение Azure Data Factory Pipelines

Что такое трубопроводы?

Трубопровод — это логическая группировка действий, которые выполняют единицу работы.Деятельность может быть простой — например, копирование данных из Azure Blob Storage в Azure SQL Database — или сложной — например, запуск ноутбука Databricks, выполнение хранимой в SQL процедуры или вызов пользовательского REST API. Трубопроводы позволяют определять поток выполнения, включая условное ветвление, петлю и параллельную обработку. Каждый трубопровод может иметь одно или несколько действий, и действия могут быть связаны через контрольный поток зависимости, такие как «На успехе», «На отказе» или «Скип». Это позволяет создавать сложную, ветвящуюся логику без написания какого-либо кода.

Ключевые виды деятельности

Azure Data Factory подразделяет деятельность на три основные группы:

  • Действия по перемещению данных — Эти данные копируются между поддерживаемыми хранилищами данных.Основной деятельностью является Копирование активности, которая поддерживает более 90 встроенных разъемов (например, Amazon S3, Google BigQuery, Snowflake, SAP HANA).
  • Деятельность по преобразованию данных — Эти данные трансформируются с использованием вычислительных ресурсов. Примеры включают HDInsight Hive, Azure Databricks (Python, Scala, или R), Stored Procedure и SSIS Integration Runtime.
  • Контрольная деятельность — Они организуют поток трубопровода. Примерами являются ForEach (петля над коллекцией), If Condition (ветвление), Wait (выполнение паузы), Execute Pipeline (вызов другого трубопровода) и Validation Activity (проверка существования файла).

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

Зависимость от деятельности и выполнение трубопроводов

Действия внутри трубопровода выполняются на основе их зависимостей. По умолчанию действия выполняются последовательно. Для параллельного выполнения действий можно опустить зависимости. Мощной особенностью является возможность использования динамических выражений и параметров. Например, вы можете передать имя набора данных, параметр даты или строку соединения в качестве переменной, делая трубопроводы многоразовыми в разных средах. Деятельность также поддерживает политики повторного использования (количество повторных запросов, интервал повторного использования) и тайм-ауты, которые необходимы для надежности производства.

Триггеры Azure Data Factory: запланированное исполнение

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

Расписание триггеров

Триггеры расписания запускают трубопроводы по фиксированному календарному графику - например, каждые 15 минут, почасово в верхней части часа или ежедневно в 3:00 утра. Вы настраиваете повторение с использованием кронообразного выражения или простого интервала (минуты, часы, дни, недели, месяцы). Триггеры расписания также поддерживают расширенные параметры, такие как время начала, время окончания и часовой пояс. Они идеально подходят для повторяющихся рабочих мест ETL, таких как ночная нагрузка на хранилище данных или почасовое обновление отчетности.

Триггер событий

Триггеры событий реагируют на внешние события, чаще всего события из Azure Blob Storage или Azure Data Lake Storage Gen2. Например, можно создать триггер, который запускается при поступлении нового файла в конкретный контейнер, или при обновлении файла. ADF поддерживает две категории триггеров событий:

  • Перехватывающие триггеры событий хранения — активируются событиями хранения в виде блоба (т.е. BlobCreated, BlobDeleted). Можно фильтровать события по префиксу имени блоба, суффиксу и пути. Это широко используется для шаблонов приема в режиме реального времени, таких как обработка входящих файлов CSV из системы продаж.
  • Таможенные триггеры событий — Основанные на пользовательских темах Azure Event Grid. Это позволяет запускать трубопроводы в ответ на любое событие, связанное с конкретной областью, такое как завершенное обучение модели машинного обучения, действие пользователя или изменение в сторонней системе. Пользовательские триггеры делают ADF гибким оркестратором в архитектурах, управляемых событиями.

Триггеры событий не работают по фиксированному графику — они работают только тогда, когда происходит определенное событие, что делает их экономически эффективными и своевременными.

Триггеры с грохочущими окнами

Этот тип триггера находится между триггерами расписания и событий. Триггер с падением окна работает на фиксированной частоте, но также обеспечивает управление состоянием — он запоминает, какие окна уже были обработаны. Например, вы можете установить триггер с падением окна для запуска каждый час, и он будет запускать точно в начале каждого окна (например, 00:00–01:00, 01:00–02:00). Каждое окно является независимым, и триггер обеспечивает точно один раз семантику обработки. Падение окна особенно полезно для дополнительных нагрузок данных, где вам нужно обрабатывать данные для определенного временного диапазона без перекрытия или пропуска каких-либо интервалов.

Создание и управление триггерами

Триггеры могут создаваться и управляться с помощью нескольких интерфейсов:

  • Лазурный портал (UI): Самый простой способ для разовых настроек. Можно определить триггер, протестировать его и связать с одним или несколькими трубопроводами. Портал предоставляет визуальный интерфейс для настройки повторения, фильтров событий и параметров.
  • Azure CLI или PowerShell: Подходит для сценариев и интеграции DevOps. Например, можно использовать cmdlet для создания триггера программно.
  • ARM Шаблоны (Azure Resource Manager): Рекомендуемый подход для инфраструктуры в виде кода (IaC). Вы можете определить триггеры как ресурсы JSON внутри шаблона ARM и развернуть их через Azure DevOps или GitHub Actions.
  • REST API: Для продвинутой автоматизации или при интеграции с внешними системами оркестровки можно напрямую позвонить в ADF REST API.

Одна критическая точка: триггер должен быть явно связан с трубопроводом, прежде чем он сможет начать работу. Вы можете связать один триггер с несколькими трубопроводами или один трубопровод с несколькими триггерами, в зависимости от вашего рабочего процесса.

Продвинутый триггер и интеграция трубопроводов

Триггерные зависимости и цепь

В сложных ландшафтах данных вам может понадобиться один трубопровод для запуска за другим или триггер для ожидания определенного события перед началом. Azure Data Factory поддерживает цепные трубопроводы с использованием Execute Pipeline Activity — контрольная активность внутри родительского трубопровода, который работает синхронно или асинхронно. Для внешних зависимостей между триггерами вы можете комбинировать триггеры с шаблонами на основе событий. Например, у вас может быть триггер расписания, который запускает трубопровод А (который загружает необработанные данные), а затем триггер события, который запускается, когда трубопровод А записывает маркер завершения в хранилище, тем самым запуская трубопровод B (который преобразует данные). Этот подход отделяет трубопроводы и делает их отказоустойчивыми.

Интеграция с Azure Monitoring и Alerts

Azure Data Factory глубоко интегрируется с Azure Monitor и Log Analytics. Каждый запуск трубопровода, запуск активности и событие триггера регистрируются в диагностических журналах ADF. Вы можете передавать эти журналы в Log Analytics и создавать панели инструментов, пользовательские запросы и правила оповещения. Например, вы можете настроить оповещение, которое запускает электронную почту или веб-хук, если запуск трубопровода выходит из строя более трех раз в 15-минутном окне. Кроме того, вы можете использовать встроенный в ADF Alerts и Metrics клинок в портале для быстрой настройки уведомлений о сбоях трубопровода или запуска пропущенных окон.

Преимущества автоматизации рабочих процессов данных с помощью ADF

Использование триггеров и конвейеров Azure Data Factory для автоматизации рабочих процессов данных дает измеримые преимущества:

  • Операционная эффективность — ручная передача файлов и запланированные сценарии заменяются бессерверными управляемыми конвейерами. Это освобождает инженеров данных, чтобы сосредоточиться на логике, а не на инфраструктуре.
  • Надежность и согласованность — ADF автоматически перезаписывает неудачные действия, соблюдает тайм-ауты и регистрирует каждый шаг. Как только трубопровод спроектирован и протестирован, он работает последовательно без дрейфа.
  • Масштабируемость — ADF может обрабатывать петабайт данных и тысячи прогонов трубопровода в день. Базовая вычислительная (Azure Integration Runtime) шкала упругая, поэтому вам не нужно предоставлять серверы.
  • Контроль затрат — Вы платите только за вычисления, потребляемые действиями. Триггеры событий и триггеры окна с обрушением уменьшают отходы, работая только при необходимости. Вы также можете установить пороги для остановки дорогостоящих трубопроводов, если они превышают бюджет.
  • Конечная наблюдаемость — С помощью диагностических журналов, мониторинга и оповещения вы можете обнаружить и устранить сбои, прежде чем они повлияют на потребителей, находящихся ниже по течению.

Лучшие практики для триггеров и трубопроводов

Чтобы получить максимальную отдачу от триггеров и трубопроводов ADF в производственной среде, следуйте этим лучшим практикам:

  • Проектирование для модульности и повторного использования. Разбивайте большие трубопроводы на более мелкие, сфокусированные трубопроводы (например, один для проглатывания, один для очистки, один для загрузки). Используйте параметры и продавайте их между трубопроводами с использованием эксплуатационной активности трубопровода. Это облегчает тестирование, отладку и техническое обслуживание.
  • Используйте зависимости триггеров осторожно. Для рабочих нагрузок, требующих строгого последовательного выполнения, предпочтите цепь через активность Execute Pipeline, а не полагаться на внешние маркеры событий. Для слабо связанных этапов спусковые механизмы событий идеальны.
  • Внедрить надежную обработку ошибок. Внутри каждого трубопровода добавьте действия If Condition для проверки на успех или неудачу. При отказе введите ошибку в журнал и необязательно отправьте оповещение. Используйте зависимость «On Failure» для запуска конвейера восстановления (например, отправьте файл повторно, уведомите команду).
  • Параметризуйте все. Используйте параметры трубопровода для путей файлов, строк соединений или интервалов расписания. Избегайте значений жесткого кодирования. Это позволяет продвигать один и тот же артефакт в средах разработки, тестирования и производства.
  • Версия контролирует ваши трубопроводы. Экспортируйте свои трубопроводы и триггеры в качестве шаблонов ARM и храните их в репозитории Git (Azure Repos или GitHub). Используйте встроенную интеграцию Git ADF для связи репозитория с вашей фабрикой. Это позволяет сотрудничать, просматривать код и откатывать.
  • Мониторинг затрат и производительности. Включите диагностические журналы и отправьте их в Log Analytics.Запрос на дорогостоящие или длительные действия. Настройте настройки Azure Integration Runtime DIU (Data Integration Unit) для копирования действий для оптимизации пропускной способности.
  • Сначала тестируйте триггеры в непроизводственной среде. Всегда подтверждайте, что триггер загорается в нужное время или в нужном событии, прежде чем включить его в производство. Распространенной ошибкой является оставить активным триггер графика при разработке трубопровода, вызывая неожиданные прогоны.
  • Используйте фильтрацию событий для снижения шума. При создании триггеров событий укажите префиксы имени файла, суффиксы и пути, чтобы избежать запуска нерелевантных событий с блобами. Это экономит вычислительную стоимость и предотвращает потерянные запуски.

Случаи общего использования

Инкрементные загрузки данных

Один из наиболее распространенных шаблонов заключается в загрузке только новых или измененных данных из исходной системы (например, транзакционной базы данных) в хранилище данных. Триггер окна, работающий каждые 15 минут, может выполнять конвейер, который копирует строки, где «последняя измененная» временная метка попадает в это окно. Затем трубопровод может вставлять данные в Azure Synapse Analytics или Azure SQL Database.

Потребление файлов в реальном времени

Когда партнер загружает файл CSV в контролируемый контейнер Azure Blob Storage, триггер событий запускает конвейер, который проверяет схему, перемещает файл в папку «обработка», запускает поток данных для преобразования данных и, наконец, загружает его в базу данных SQL. Этот шаблон распространен в розничных и логистических системах.

Ночной пакетный процессинг

Триггер расписания, установленный до 2:00 утра UTC, выполняет ряд конвейеров: во-первых, копирует данные о дополнительных продажах с локального SQL Server на Azure Blob; во-вторых, запускает работу HDInsight Hive для агрегирования данных; в-третьих, выполняет сохраненную процедуру в базе данных Azure SQL для обновления таблиц отчетов.

Гибридная оркестровка данных

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

Мониторинг и устранение неполадок

Использование Azure Monitor и Log Analytics

Чтобы получить глубокое представление о выполнении триггеров и трубопроводов, настройте диагностические настройки на вашей фабрике данных для отправки журналов в рабочее пространство Log Analytics. После этого вы можете запускать запросы Kusto, такие как:

ADFActivityRun
| where ActivityName == 'Copy data1' and Status == 'Failed'
| project TimeGenerated, PipelineName, ActivityName, ErrorMessage

Настройте правила оповещения, чтобы уведомить вас о том, что трубопровод не работает или триггер не запускается в ожидаемом окне. Вы также можете визуализировать историю запуска с помощью Azure Workbooks или Power BI.

Общие вопросы и решения

  • Трейггер не зажигает: Проверить состояние триггера (старт/стоп). Проверить, что соответствующий трубопровод опубликован и находится в Активном состоянии. Для триггеров событий подтвердить, что учетная запись хранения или тема сетки событий правильно настроена и что подписка на событие не отфильтрована.
  • Пипелайн висит или раз отключается: Деятельность имеет по умолчанию тайм-аут 7 дней. Установите явные тайм-ауты для трубопроводов, которые должны быстро выйти из строя. Используйте активность «Валидация» для проверки существования файлов перед продолжением.
  • Параметрическое несоответствие: Если триггер проходит параметры трубопровода, которые не соответствуют определению трубопровода, прогон будет неуспешным. Убедитесь, что имена и типы параметров согласованы. Используйте значения по умолчанию в трубопроводе, чтобы обеспечить гибкость.
  • Проблемы с параллелями: По умолчанию трубопровод может работать до 100 одновременных экземпляров.Если у вас есть триггер с перекрывающимися окнами, установите максимальную параллель на триггере до 1 для обеспечения последовательной обработки.

Заключение

Триггеры и конвейеры Azure Data Factory обеспечивают мощную, гибкую платформу для автоматизации рабочих процессов данных в любом масштабе. Понимая различия между триггерами расписания, события и падения окон, а также применяя модульную конструкцию трубопровода и надежные методы мониторинга, инженеры данных могут создавать надежные, экономически эффективные и поддерживающие решения для интеграции данных. Независимо от того, управляете ли вы дополнительными нагрузками, проглатыванием событий в реальном времени или сложной партией ETL, ADF дает вам инструменты для автоматизации с уверенностью. Для дальнейшего чтения изучите официальную документацию по трубопроводу , обзор типов триггеров и руководство по мониторингу .