Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Используйте ожидания для применения ограничений качества, которые проверяют данные по мере потоков через конвейеры ETL. Ожидания обеспечивают более подробные сведения о метриках качества данных и позволяют завершать обновления или удалять записи при обнаружении недопустимых записей.
Дополнительные варианты использования и рекомендуемые рекомендации см. в рекомендациях по ожиданиям и расширенных шаблонах.
Что такое ожидания?
Ожидания — это необязательные предложения в материализованном представлении конвейера, потоковой таблице или инструкциях создания представлений, которые применяют проверки качества данных к каждой записи, передаваемой через запрос. Ожидания используют стандартные логические инструкции SQL для указания ограничений. Вы можете объединить несколько ожиданий для одного набора данных и задать ожидания во всех объявлениях набора данных в конвейере.
Замечание
Вы также можете определить ожидания для потоковых таблиц и материализованных представлений, поддерживаемых автономным конвейером, созданным в Databricks SQL. Используйте условие CONSTRAINT expectation_name EXPECT (expectation_expr) в CREATE STREAMING TABLE и CREATE MATERIALIZED VIEW.
В следующих разделах представлены три компонента ожидания и приведены примеры синтаксиса.
Название ожидания
Каждое ожидание должно иметь имя, которое используется в качестве идентификатора для отслеживания и мониторинга ожидания. Выберите имя, которое передает проверяемые метрики. В следующем примере определяется ожидание valid_customer_age чтобы подтвердить, что возраст находится в диапазоне от 0 до 120 лет.
Это важно
Имя ожидания должно быть уникальным для заданного набора данных. Ожидания можно использовать повторно в нескольких наборах данных в производственной линии. См. переносимые и многократно используемые ожидания.
Питон
@dp.table
@dp.expect("valid_customer_age", "age BETWEEN 0 AND 120")
def customers():
return spark.readStream.table("datasets.samples.raw_customers")
SQL
CREATE OR REFRESH STREAMING TABLE customers(
CONSTRAINT valid_customer_age EXPECT (age BETWEEN 0 AND 120)
) AS SELECT * FROM STREAM(datasets.samples.raw_customers);
Ограничение для оценки
Предложение ограничения — это условный оператор SQL, который должен иметь значение true или false для каждой записи. Ограничение содержит фактическую логику для проверяемого объекта. Если запись не удовлетворяет этому условию, срабатывает ожидание.
Ограничения должны использовать допустимый синтаксис SQL и не могут содержать следующее:
- Пользовательские функции Python
- Вызовы внешних служб
- Вложенные запросы, ссылающиеся на другие таблицы
Ниже приведены примеры ограничений, которые можно добавить в инструкции создания набора данных:
Питон
Синтаксис ограничения в Python:
@dp.expect(<constraint-name>, <constraint-clause>)
Можно указать несколько ограничений:
@dp.expect(<constraint-name>, <constraint-clause>)
@dp.expect(<constraint2-name>, <constraint2-clause>)
Примеры.
# Simple constraint
@dp.expect("non_negative_price", "price >= 0")
# SQL functions
@dp.expect("valid_date", "year(transaction_date) >= 2020")
# CASE statements
@dp.expect("valid_order_status", """
CASE
WHEN type = 'ORDER' THEN status IN ('PENDING', 'COMPLETED', 'CANCELLED')
WHEN type = 'REFUND' THEN status IN ('PENDING', 'APPROVED', 'REJECTED')
ELSE false
END
""")
# Multiple constraints
@dp.expect("non_negative_price", "price >= 0")
@dp.expect("valid_purchase_date", "date <= current_date()")
# Complex business logic
@dp.expect(
"valid_subscription_dates",
"""start_date <= end_date
AND end_date <= current_date()
AND start_date >= '2020-01-01'"""
)
# Complex boolean logic
@dp.expect("valid_order_state", """
(status = 'ACTIVE' AND balance > 0)
OR (status = 'PENDING' AND created_date > current_date() - INTERVAL 7 DAYS)
""")
SQL
Синтаксис ограничения в SQL:
CONSTRAINT <constraint-name> EXPECT ( <constraint-clause> )
Несколько ограничений должны быть разделены запятой:
CONSTRAINT <constraint-name> EXPECT ( <constraint-clause> ),
CONSTRAINT <constraint2-name> EXPECT ( <constraint2-clause> )
Примеры.
-- Simple constraint
CONSTRAINT non_negative_price EXPECT (price >= 0)
-- SQL functions
CONSTRAINT valid_date EXPECT (year(transaction_date) >= 2020)
-- CASE statements
CONSTRAINT valid_order_status EXPECT (
CASE
WHEN type = 'ORDER' THEN status IN ('PENDING', 'COMPLETED', 'CANCELLED')
WHEN type = 'REFUND' THEN status IN ('PENDING', 'APPROVED', 'REJECTED')
ELSE false
END
)
-- Multiple constraints
CONSTRAINT non_negative_price EXPECT (price >= 0),
CONSTRAINT valid_purchase_date EXPECT (date <= current_date())
-- Complex business logic
CONSTRAINT valid_subscription_dates EXPECT (
start_date <= end_date
AND end_date <= current_date()
AND start_date >= '2020-01-01'
)
-- Complex boolean logic
CONSTRAINT valid_order_state EXPECT (
(status = 'ACTIVE' AND balance > 0)
OR (status = 'PENDING' AND created_date > current_date() - INTERVAL 7 DAYS)
)
Действие на недействительную запись
Необходимо указать действие, чтобы определить, что происходит при сбое проверки записи. В следующей таблице описываются доступные действия.
| Действие | Синтаксис SQL | Синтаксис Python | Result |
|---|---|---|---|
| предупреждение (по умолчанию) | EXPECT |
dp.expect |
Недопустимые записи записываются в целевой объект. |
| выпадение | EXPECT ... ON VIOLATION DROP ROW |
dp.expect_or_drop |
Недопустимые записи удаляются перед записью данных в целевой объект. Количество удаленных записей регистрируется вместе с другими метриками набора данных. |
| провал | EXPECT ... ON VIOLATION FAIL UPDATE |
dp.expect_or_fail |
Недопустимые записи препятствуют успешному обновлению. Перед повторной обработкой требуется вмешательство вручную. |
Вы также можете реализовать расширенную логику для карантина недопустимых записей без сбоя или удаления данных. См. Карантин недопустимых записей.
Метрики отслеживания ожиданий
Вы можете просмотреть метрики отслеживания для действий warn или drop из пользовательского интерфейса конвейера. Поскольку fail прерывает обновление при обнаружении недопустимой записи, метрики не записываются.
Замечание
Для потоковых таблиц и материализованных представлений, поддерживаемых автономным конвейером, созданным в Databricks SQL, вкладка "Качество данных " в пользовательском интерфейсе конвейера недоступна. Запросите журнал событий, чтобы просмотреть метрики ожидания. См. Запрос по качеству данных или метрики ожиданий.
Чтобы просмотреть метрики ожидания, выполните следующие действия.
- На боковой панели рабочей области Azure Databricks щелкните "Задания и конвейеры".
- Щелкните название вашего конвейера.
- Щелкните набор данных с определенным ожиданием.
- Перейдите на вкладку "Качество данных " в правой боковой панели.
Вы можете просматривать метрики качества данных, запрашивая журнал событий конвейера Lakeflow. См. Запрос по качеству данных или метрики ожиданий.
Сохранение недопустимых записей
Сохранение недопустимых записей — это поведение по умолчанию для ожиданий. Используйте оператор expect, если вы хотите сохранить записи, которые нарушают ожидание, но при этом собирать метрики о количестве записей, удовлетворяющих или не соответствующих ограничениям. Записи, которые нарушают ожидание, добавляются в целевой набор данных вместе с допустимыми записями:
Питон
@dp.expect("valid timestamp", "timestamp > '2012-01-01'")
SQL
CONSTRAINT valid_timestamp EXPECT (timestamp > '2012-01-01')
Удаление недопустимых записей
expect_or_drop Используйте оператор, чтобы предотвратить дальнейшую обработку недопустимых записей. Записи, которые нарушают ожидание, удаляются из целевого набора данных:
Питон
@dp.expect_or_drop("valid_current_page", "current_page_id IS NOT NULL AND current_page_title IS NOT NULL")
SQL
CONSTRAINT valid_current_page EXPECT (current_page_id IS NOT NULL and current_page_title IS NOT NULL) ON VIOLATION DROP ROW
Сбой при обнаружении недопустимых записей
Если недопустимые записи неприемлемы, используйте expect_or_fail оператор, чтобы остановить выполнение сразу после сбоя проверки записи. Если операция является обновлением таблицы, система атомарно откатывает транзакцию:
Питон
@dp.expect_or_fail("valid_count", "count > 0")
SQL
CONSTRAINT valid_count EXPECT (count > 0) ON VIOLATION FAIL UPDATE
Это важно
В активированном конвейере сбой одного потока не приводит к сбою других параллельных потоков. В непрерывном конвейере сбой ожидания останавливает поток, все зависимые потоки и конвейер выдает сообщение, объясняющее, почему он остановлен.
Для более точного управления оркестрацией рабочих процессов, если проверка завершается с ошибкой, разделите проверку и последующую обработку на отдельные пайплайны и координируйте их с помощью потока управления между задачами пайплайнов. См. таблицы проверки и поток управления конвейером.
Устранение неполадок с обновлениями, которые не оправдали ожидания
Если конвейер завершается сбоем из-за нарушения ожидания, необходимо исправить код конвейера, чтобы правильно обрабатывать недопустимые данные перед повторной запуском конвейера.
Ожидания, настроенные для обработки сбоев в конвейерах, изменяют Spark план запросов ваших преобразований для отслеживания сведений, необходимых для обнаружения и составления отчетов о нарушениях. Эти сведения можно использовать для определения входной записи, которая привела к нарушению для многих запросов. Конвейеры Lakeflow предоставляют выделенное сообщение об ошибке для сообщения о таких нарушениях. Ниже приведен пример сообщения об ошибке нарушения ожидания:
[EXPECTATION_VIOLATION.VERBOSITY_ALL] Flow 'sensor-pipeline' failed to meet the expectation. Violated expectations: 'temperature_in_valid_range'. Input data: '{"id":"TEMP_001","temperature":-500,"timestamp_ms":"1710498600"}'. Output record: '{"sensor_id":"TEMP_001","temperature":-500,"change_time":"2024-03-15 10:30:00"}'. Missing input data: false
Управление несколькими ожиданиями
Замечание
Хотя SQL и Python поддерживают несколько ожиданий в одном наборе данных, только Python позволяет группировать несколько ожиданий и указывать коллективные действия.
Можно объединить несколько ожиданий и указать коллективные действия с помощью функций expect_all, expect_all_or_dropа также expect_all_or_fail.
Эти декораторы принимают словарь Python в качестве аргумента, где ключ является именем ожидания и значением является ограничение ожидания. Вы можете повторно использовать один набор ожиданий в нескольких наборах данных в конвейере. Ниже приведены примеры каждого из expect_all операторов Python:
valid_pages = {"valid_count": "count > 0", "valid_current_page": "current_page_id IS NOT NULL AND current_page_title IS NOT NULL"}
@dp.table
@dp.expect_all(valid_pages)
def raw_data():
# Create a raw dataset
@dp.table
@dp.expect_all_or_drop(valid_pages)
def prepared_data():
# Create a cleaned and prepared dataset
@dp.table
@dp.expect_all_or_fail(valid_pages)
def customer_facing_data():
# Create cleaned and prepared to share the dataset
Ограничения
- Поскольку для этих типов объектов поддерживаются только потоковые таблицы, материализованные представления и временные представления, метрики качества данных поддерживаются только для этих типов объектов.
- Метрики качества данных недоступны, если:
- В запросе не определены ожидания.
- Поток использует оператор, который не поддерживает необходимые функции.
- Тип потока, например приемники, не поддерживает ожидания.
- Для заданного запуска потока отсутствуют обновления связанной таблицы потоковой передачи или материализованного представления.
- Конфигурация конвейера не включает необходимые параметры для записи метрик, таких как
pipelines.metrics.flowTimeReporter.enabled.
- В некоторых случаях
COMPLETEDпоток может не содержать метрик. Вместо этого метрики сообщаются в каждом микропакетеflow_progressв событии с статусомRUNNING. - Так как представления вычисляются только при запросе, метрики качества данных могут быть недоступны для определенного представления. Кроме того, представление, запрашиваемое в нескольких подчиненных наборах данных, может иметь несколько наборов метрик качества данных.
- Ожидания не поддерживаются с
AUTO CDC FROM SNAPSHOT.