Данные Azure Data Factory Поток для сложных преобразований данных
Введение в Azure Data Factory Data Flows
Azure Data Factory (ADF) выступает в качестве полностью управляемой облачной службы интеграции данных, которая позволяет организациям организовывать и автоматизировать движение и преобразование данных. По своей сути ADF обеспечивает безкодовую визуальную среду для построения трубопроводов ETL и ELT. Среди его наиболее мощных возможностей - функция FLT:0]Data Flow, которая позволяет инженерам данных проектировать сложные преобразования данных с использованием графического холста, а не писать традиционный код. Эта статья углубляется в архитектуру, компоненты и расширенные варианты использования ADF Data Flows, предлагая всеобъемлющее руководство для овладения сложными преобразованиями данных в масштабе.
Потоки данных построены на кластерах Apache Spark, управляемых Azure, обеспечивая эластичное, высокопроизводительное выполнение. Они позволяют выполнять широкий спектр операций, включая фильтрацию, агрегирование, объединение, поворот и применение пользовательских выражений, без необходимости писать код Spark. Эта абстракция сокращает время разработки, снижает барьер для менее технических пользователей и гарантирует, что преобразования остаются поддерживаемыми и проверяемыми. Независимо от того, слияние ли вы гетерогенных источников данных, очистка потоковых данных или подготовка наборов данных для машинного обучения, ADF Data Flows обеспечивает надежное решение.
Понимание архитектуры потоков данных ADF
Чтобы эффективно использовать потоки данных, важно понять их базовую архитектуру. Каждый поток данных работает на временном кластере Spark, который развернется во время выполнения и завершится после завершения. Этот дизайн обеспечивает экономическую эффективность - вы платите только за вычислительные ресурсы, потребляемые во время преобразования. Размер кластера, количество ядер и память могут быть настроены в соответствии с объемом и сложностью данных.
Режимы исполнения
ADF Data Flows поддерживает два основных режима выполнения:
- Режим отладки — используется для интерактивного тестирования и разработки. Он работает на небольшом кластере Spark (8 ядер) и позволяет просматривать данные на каждом этапе трансформации. Режим отладки необходим для проверки логики перед развертыванием производства.
- Режим работы с конвейером — Используется для запланированных или запущенных производственных исполнений. Вы можете указать параметры кластера, такие как тип вычислений (общая цель, оптимизированная память), количество ядер и время в пути (TTL) для оптимизации стоимости и производительности.
Понимание этого различия имеет решающее значение для оценки затрат и производительности.В производстве всегда тестируйте преобразования в режиме Debug локально, прежде чем развертывать их в трубопроводы.
Data Flow vs. Copy Activity (англ.) (недоступная ссылка).
ADF Copy Activity предназначен для высокоскоростного, схема-агностичного перемещения данных. Потоки данных, наоборот, предназначены для преобразований, учитывающих схемы. В то время как Copy Activity может выполнять простые отображения и преобразования типов с помощью вкладки Mapping, Data Flows предлагает десятки типов преобразований и возможность обрабатывать сложную бизнес-логику. Для сценариев, требующих нескольких соединений, условных расколов или оконных функций, потоки данных являются подходящим выбором.
Ключевые компоненты потока данных
Каждый поток данных состоит из трех основных категорий компонентов: Источники, Трансформации и Зубы. Кроме того, вы можете использовать параметры и Переменные , чтобы сделать ваши потоки динамичными и многоразовыми.
1.Источник
Источник определяет, откуда берутся ваши данные. Azure Data Factory поддерживает широкий спектр типов источников, включая Azure Blob Storage Gen2, Azure Data Lake Storage Gen2, Azure SQL Database, Synapse Analytics, Amazon S3, Google Cloud Storage и локальные базы данных через самоорганизующиеся среды выполнения интеграции. Каждый источник может быть настроен с деталями подключения, форматом файла (Parquet, CSV, JSON, Avro, ORC) и определением схемы. Используя Schema Drift, Data Flows может автоматически адаптироваться к изменениям в схеме источника — критическая функция для обработки полуструктурированных или развивающихся данных.
Наилучшая практика заключается в использовании форматов Parquet или Delta Lake для источника и поглотителя из-за их колоннального хранения и эффективности сжатия. Эти форматы значительно ускоряют операции чтения/записи и снижают стоимость.
2.Трансформации
ADF Data Flows предлагает богатую библиотеку мероприятий по трансформации.
- Расширители: Фильтр, сортировка и переменная строка (для операций вставки/обновления/удаления).
- Колонные модификаторы: Выберите, производную колонку, агрегируйте, окно, поворот, разворот и ранжирование.
- Множественные входы/выходы: Присоединяйтесь, смотрите, существует, союз и условный раскол.
- Модификаторы схем: Новая ветвь, Ассерт (правила качества данных) и Суррогатный ключ.
Преобразование Производная колонка особенно мощно — вы можете создавать выражения с помощью встроенного конструктора выражений, который включает функции для манипулирования строками, арифметики даты / времени, математических операций и сопоставления шаблонов (по аналогии с SQL). Например, вы можете создать новую колонку «Полное имя» путем объединения «Первое имя» и «Последнее имя» с пространством.
3. Потонуть
Sink определяет, где находятся преобразованные данные. Как и источники, поглотители могут быть любым поддерживаемым хранилищем данных. Критические настройки включают формат файла, стратегию раздела (Hash, Dynamic, Round Robin или File Name) и режим вывода (Append vs. Overwrite). Для поглотителей Delta Lake вы можете включить Слияние , Обновление или Upsert поведение, позволяя Data Flows действовать как мини-погрузчик хранилища данных.
Реализация сложных преобразований: подробный сценарий
Давайте рассмотрим реальный пример: Обогащение клиента 360. Представьте, что у вас есть три источника исходных данных:
- Профили клиентов (CSV из хранилища Blob)
- История транзакций (паркет ADLS Gen2)
- Каталог продуктов (Azure SQL Database)
Цель состоит в том, чтобы создать единый обогащенный набор данных, который содержит для каждого клиента: их демографию, общие расходы, предпочтения в категориях продуктов и ярлык уровня лояльности. Эта трансформация будет включать в себя несколько шагов потока данных, выполняемых в одном конвейере.
Шаг 1: Загрузка и чистые источники
Добавить три узла источника. Для профилей клиентов использовать производную колонку для стандартизации формата «DateOfBirth» и удалить строки с нулевыми адресами электронной почты. Для транзакций отфильтровать возвращенные транзакции (где «Amount < 0»). Для каталога продукта присоедините имя категории с идентификатором категории.
Шаг 2: Присоединяйтесь к транзакциям с клиентами
Добавить преобразование Присоединиться , чтобы объединить очищенные Профили клиентов и историю транзакций на «CustomerID». Используйте внутреннее соединение, чтобы исключить клиентов без транзакций. Затем используйте преобразование Выбрать , чтобы отбросить дублирующие столбцы (например, переименовать «CustomerID» со второго входа).
Шаг 3: Совокупность на одного клиента
Подключите объединенный выход к преобразованию Aggregate.Группу по 'CustomerID' и 'CustomerName', и вычислите Sum(Amount) как TotalSpending, Count(] как TransactionCount, и Max(TransactionDate) как LastPurchaseDate.
Шаг 4: Обогащайтесь предпочтениями продукта
Используйте вторую Присоединиться, чтобы прикрепить Каталог продукта на 'ProductID' (который существует в источнике транзакции). Затем добавьте преобразование Pivot , чтобы преобразовать названия категорий в столбцы (например, электроника, одежда, дом) с подсчетом покупок в категории.
Шаг 5: Определите уровень лояльности
Добавить преобразование Производная колонка , в котором используется вложенная логика if-else для присвоения уровней лояльности: «если (TotalSpending > 10000, «Золото», if (TotalSpending > 5000, «Серебряный», «Бронзовый»)»
Шаг 6: Напишите обогащенные данные
Подключите конечный вывод к Sink, который нацелен на таблицу базы данных Azure SQL или папку Delta Lake в ADLS Gen2. Настройте раковину для использования Upsert поведения на «CustomerID», чтобы последующие запуски обновляли существующие записи вместо их дублирования.
Весь этот процесс разработан визуально, каждый шаг тестируется в режиме отладки. Полученный трубопровод поддерживается, самодокументируется и может быть запланирован почасово или ежедневно.
Лучшие практики для высокопроизводительных потоков данных
Оптимизация производительности потока данных имеет важное значение при работе с терабайтами данных. Следуйте этим проверенным практикам:
- Используйте соответствующую размерность кластера: Для больших наборов данных выберите по меньшей мере 16-32 ядра. Для операций с интенсивной памятью (например, присоединений или агрегации) выберите Memory Optimized Compute.
- Разделите ваши данные: В настройках Источника включите обрезку разделов с помощью Параметров Раздела. Установите шаблон пути папки для чтения только соответствующих разделов.
- Минимизируйте перетасовку данных: Объединения и агрегации вызывают перетасовку операций по кластеру. Если вы можете, предварительно отфильтровать данные перед присоединением. Используйте Широковещательное присоединение для небольших таблиц поиска (например, таблица размеров 1 МБ).
- Оптимизируйте форматы файлов: Предпочитайте паркет или дельту по сравнению с CSV/JSON для источников и поглотителей. Эти столбцовые форматы уменьшают I/O и используют выталкивание предикатов.
- Уменьшить ветви преобразования: Каждый Новый Ветвь дублирует поток данных. Используйте условный сплит только тогда, когда это необходимо; в противном случае, слить условия в производных колонках.
- Использовать мониторинг потока данных: В мониторе ADF проверьте журналы выполнения потока данных на длительность этапов. Ищите долгосрочные преобразования и рассмотрите возможность их разбиения на более мелкие этапы.
Внешний ресурс: Официальное руководство Microsoft по производительности для потоков данных ADF
Мониторинг и отладка потоков данных
Эффективный мониторинг гарантирует надежную работу ваших конвейеров данных. ADF предоставляет встроенные возможности мониторинга потоков данных. Вы можете просматривать состояние исполнения, количество строк на каждом этапе и время, затрачиваемое на преобразование. Ключевые показатели для наблюдения включают:
- Время обработки — время выполнения кластера Total Spark.
- Data Skew — неравномерное распределение данных по разделам, видимое на выходе стадии.
- Счета строк — неожиданные падения строк могут указывать на проблемы с фильтром или присоединением.
Для отладки используйте режим Data Flow Debug . Он работает на небольшом кластере и позволяет вам интерактивно проверять выход каждого преобразования. Для дальнейшей диагностики сложных выражений вы можете использовать преобразование Assert для проверки правил качества данных (например, «isNotNull(CustomerID)») и фиксации сбоев.
Рассмотрение вопросов безопасности
Потоки данных часто обрабатывают конфиденциальную информацию. ADF интегрируется с Azure Key Vault для хранения строк соединений и учетных данных. Всегда используйте управляемую идентификацию или основную аутентификацию службы по ключам учетной записи. Для данных в пути Data Flows используют TLS; для данных в состоянии покоя убедитесь, что ваши места хранения зашифрованы (шифрование Azure Storage включено по умолчанию). Кроме того, вы можете применять преобразования уровня столбца, такие как маскировка или хеширование в выражениях Data Flow, используя такие функции, как sha2() или substring().
Интеграция потоков данных с другими сервисами Azure
Потоки данных ADF не работают изолированно. Они могут быть организованы с другими мероприятиями ADF для создания сквозных трубопроводов:
- Выполнить трубопроводную деятельность: Запустить другой трубопровод ADF после завершения Data Flow.
- Databricks Notebook: Для расширенной аналитики или вывода ML объедините поток данных с Databricks.
- Лазурные функции: Вызовите пользовательский бессерверный код для обогащения, который требует сторонних API.
- Power BI: Поглощают преобразованные данные непосредственно в наборы данных Power BI через разъем Power BI ADF.
Внешний ресурс: Обзор документации по потоку данных на фабрике данных Лазурного завода
Обычные подводные камни и как их избежать
- Сложный единичный поток данных: Разбейте монстра 50-трансформации на несколько потоков данных со таблицами постановки. Это улучшает управляемость и позволяет частичное повторение.
- Игнорирование дрейфа схем: Используйте опции Schema Drift в Source и Sink, чтобы изящно обрабатывать новые колонки без отказа трубопровода.
- Забывание времени на жизнь (TTL): Установите TTL 5-10 минут на вашем производственном кластере, чтобы сохранить теплые ресурсы для последующих потоков данных в том же конвейере.
- Не используя параметры: Названия таблиц жесткого кодирования или пути файлов делают трубопроводы жесткими. Используйте параметры трубопровода и передавайте их в параметры потока данных для максимальной повторной возможности использования.
Реальные примеры использования ADF Data Flows
Данные Lakehouse ELT
Многие организации используют Data Flows для преобразования сырых бронзовых/серебряных/золотых слоев в Data Lakehouse. Например, розничная компания глотает сырые данные о продажах в бронзовую зону, затем использует Data Flows для очистки, дублирования и объединения в серебро и, наконец, обогащает измерениями для создания золотого слоя для аналитики. Эта модель эффективно заменяет традиционные инструменты ETL, такие как SSIS.
Агрегация в реальном времени для панелей
Объедините потоки данных с триггерами на основе событий для обработки потоковых данных (например, показаний датчиков IoT) в режиме реального времени. Хотя потоки данных не передаются в потоковом режиме (они работают на микро-пакетах), они могут выполняться каждые 1-5 минут для создания агрегированных просмотров для Power BI.
Маскировка данных для соблюдения
Финансовые учреждения используют потоки данных для маскировки личной информации (PII) при перемещении данных из производственной среды в тестовую среду. Используя выражения Derived Column, они заменяют адреса электронной почты на «concat (слева (Email,1), ***@example.com)» и хеш-номера социального страхования.
Сравнение с Azure Databricks
Хотя ADF Data Flows и Azure Databricks могут выполнять сложные преобразования, они обслуживают разные персоны. Data Flows предлагают интерфейс без кода / низкого кода, подходящий для инженеров данных, которые предпочитают визуальный дизайн и управляемое управление. Databricks предоставляет интерфейс ноутбука для ученых-данных и инженеров, которым необходим полный контроль над кодом Spark, пользовательскими библиотеками и интеграцией машинного обучения. Часто лучший подход - гибрид: использование Data Flows для стандартной очистки и агрегации ETL и маршрутизация данных в Databricks для расширенной аналитики или обучения модели.
Внешний ресурс: Сравнение потоков данных ADF и Azure Databricks
Заключение
Azure Data Factory Data Flows предоставляет мощную, масштабируемую и визуальную платформу для решения сложных преобразований данных в облаке. Овладевая источниками, преобразованиями, поглотителями и их конфигурациями, инженеры по обработке данных могут создавать надежные трубопроводы ETL / ELT, которые уменьшают время до понимания при сохранении не требующей кода ремонтопригодности. С помощью лучших практик, мониторинга и интеграционных шаблонов, изложенных в этой статье, вы хорошо оснащены для реализации передовых решений для преобразования данных. Начните с малого с одного потока данных, тщательно проверьте в режиме Debug и постепенно расширяйте для организации потоков данных в масштабе предприятия.
Для дальнейшего чтения, изучите официальную документацию Microsoft по Data Flow Debug mode и expression functions reference.