Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Important
Эта функция доступна в общедоступной предварительной версии. Администраторы рабочей области могут управлять доступом к этой функции на странице "Предварительные версии ". См. статью "Управление предварительными версиями Azure Databricks".
Управление доступом
Функции являются объектами каталога Unity с возможностью управления. Доступ к компоненту управляется CREATE FEATUREREAD FEATUREпривилегиями каталога Unity, а MANAGE также правами каталога Unity. Полные описания см. в справочнике по привилегиям каталога Unity.
-
CREATE FEATURE— требуется для создания компонента в схеме.create_featureиregister_featureтребуетсяCREATE FEATUREдля родительской схемы. Следуя принципу минимальных привилегий, предоставьтеCREATE FEATUREего на уровне схемы; вы также можете предоставить его в каталоге, чтобы разрешить создавать функции в любой схеме в этом каталоге. -
READ FEATURE— требуется для чтения функции и ее данных.get_feature,create_training_setи чтение материализованных данных функций для обучения или обслуживания, необходимыхREAD FEATUREдля функции.READ FEATUREпредоставлено схеме или каталогу, применяется ко всем текущим и будущим функциям, которые он содержит. -
MANAGE— требуется для управления жизненным циклом и грантами компонента. Удаление компонента сdelete_featureпомощью и материализация компонента сmaterialize_featuresпомощью илиdelete_materialized_feature, требуемаяMANAGEдля функции.
Все операции функций также требуются USE CATALOG в родительском каталоге и USE SCHEMA родительской схеме. Сведения о том, как MANAGE и READ FEATURE применяться к материализации, см. в разделе "Разрешения".
API представления функций
Feature конструктор и register_feature()
Рекомендуемый Feature подход — создать объект локально и использовать register_feature для сохранения его в каталоге Unity. Этот двухэтапный рабочий процесс позволяет экспериментировать с функциями (включая create_training_set) перед их регистрацией.
Feature(
source: DataSource, # Required: DeltaTableSource, StreamSource, or RequestSource
function: Union[AggregationFunction, ColumnSelection], # Required: Aggregation or column selection
entity: Optional[List[str]] = None, # Required for all sources except RequestSource: entity columns
timeseries_column: Optional[str] = None, # Required for all sources except RequestSource: timestamp column
name: Optional[str] = None, # Optional: Feature name (auto-generated if omitted)
description: Optional[str] = None, # Optional: Feature description
)
FeatureEngineeringClient.register_feature() регистрирует локально созданный Feature в каталоге Unity.
FeatureEngineeringClient.register_feature(
feature: Feature, # Required: A Feature instance (not already registered)
catalog_name: str, # Required: UC catalog name
schema_name: str, # Required: UC schema name
) -> Feature
from databricks.feature_engineering.entities import Feature, DeltaTableSource, AggregationFunction, Sum, RollingWindow
from datetime import timedelta
# Step 1: Construct the feature locally
feature = Feature(
source=DeltaTableSource(catalog_name="main", schema_name="store", table_name="transactions"),
entity=["user_id"],
timeseries_column="transaction_time",
function=AggregationFunction(Sum(input="amount"), RollingWindow(window_duration=timedelta(days=7))),
)
# Step 2: Register in Unity Catalog
fe = FeatureEngineeringClient()
registered_feature = fe.register_feature(
feature=feature,
catalog_name="main",
schema_name="store",
)
create_feature()
FeatureEngineeringClient.create_feature() проверяет, создает и немедленно регистрирует функцию в каталоге Unity на одном шаге. Используйте это, если вам не нужно экспериментировать с функцией локально.
FeatureEngineeringClient.create_feature(
source: DataSource, # Required: DeltaTableSource, StreamSource, or RequestSource
function: Union[AggregationFunction, ColumnSelection], # Required: Aggregation or column selection
catalog_name: str, # Required: The catalog name for the feature
schema_name: str, # Required: The schema name for the feature
entity: Optional[List[str]] = None, # Required for all sources except RequestSource: entity columns
timeseries_column: Optional[str] = None, # Required for all sources except RequestSource: timestamp column
name: Optional[str] = None, # Optional: Feature name (auto-generated if omitted)
description: Optional[str] = None, # Optional: Feature description
) -> Feature
Параметры:
-
source: источник данных, используемый в вычислениях функций (DeltaTableSourceилиStreamSourceRequestSource). -
function: этоAggregationFunctionобъединяет оператор (например,Sum(input="amount")входной столбец и период времени). ИлиColumnSelection("column_name")для сквозных функций. -
catalog_name: имя каталога каталога Unity для функции. -
schema_name: имя схемы каталога Unity для функции. -
entity: список имен столбцов, определяющих агрегирование или ключи подстановки (первичные ключи). Требуется для всех типов источников, кромеRequestSource. Например,["user_id"]агрегирует или ищет каждого пользователя. -
timeseries_column: столбец метки времени, используемый для агрегирования периода времени или выбора последнего значения. Требуется для всех типов источников, кромеRequestSource. -
name: необязательное имя функции. Если опущено, автоматически создано из входного столбца, функции и окна (например,amount_avg_rolling_7d). -
description: необязательное описание функции.
Возвращает: Проверенный экземпляр компонента
Поднимает: ValueError, если валидация не проходит
delete_feature()
Удаляет функцию из каталога Unity по полному имени.
FeatureEngineeringClient.delete_feature(
full_name: str, # Required: '<catalog>.<schema>.<feature_name>'
) -> None
fe.delete_feature(full_name="main.store.amount_sum_rolling_7d")
Перед удалением компонента удалите или обновите все модели или спецификации компонентов, ссылающиеся на нее. Если компонент был материализован, сначала удалите материализованную функцию. Узнайте , как удалить материализованную функцию.
Автоматически созданные имена
Если name опущено, имя создается автоматически. Созданные имена соответствуют шаблону: {column}_{function}_{window} Рассмотрим пример.
-
price_avg_rolling_1h(средняя цена за 1 час) -
transaction_count_rolling_30d_1d(30-дневное число транзакций с 1d задержкой от метки времени события)
Поддерживаемые функции
Функции агрегирования
Note
Функции агрегирования упаковываются вместе AggregationFunction с временным окном, как описано в периодах времени. Каждая функция принимает параметр, указывающий исходный input столбец для статистической обработки.
| Функция | Description | Пример варианта использования |
|---|---|---|
Sum(input="column") |
Общее количество значений | Ежедневное использование приложения на пользователя в минутах |
Avg(input="column") |
Среднее значение значений | Средняя сумма транзакции |
Count(input="column") |
Количество записей | Количество входов на одного пользователя |
Min(input="column") |
Минимальное значение | Самая низкая частота пульса, записанная носимым устройством |
Max(input="column") |
Максимальное значение | Максимальное количество транзакций на сеанс |
StddevPop(input="column") |
Стандартное отклонение популяции | Изменчивость ежедневной суммы транзакций для всех клиентов |
StddevSamp(input="column") |
Стандартное отклонение образца | Изменчивость показателей кликабельности рекламных кампаний |
VarPop(input="column") |
Дисперсии населения | Разброс показаний датчиков для устройств Интернета вещей на заводе |
VarSamp(input="column") |
Выборочная дисперсия | Распространение рейтингов фильмов по выборке группы |
ApproxCountDistinct(input="column", relativeSD=0.05) |
Приблизительное уникальное число | Различное количество приобретенных товаров |
ApproxPercentile(input="column", percentile=0.95, accuracy=100) |
Приблизительный процентиль | Задержка отклика p95 |
First(input="column") |
Первое значение | Первая метка времени входа |
Last(input="column") |
Последнее значение | Последняя сумма покупки |
Note
First и Last включите значения NULL по умолчанию. Чтобы пропустить значения NULL, добавьте filter_condition значение, которое явно исключает входные столбцы, которые имеют значение NULL.
ColumnSelection (сквозная передача)
ColumnSelection Выбирает один столбец из источника без применения агрегирования. Он упаковывается непосредственно в function параметр (не внутри AggregationFunction). Возвращаемый тип выводится из исходной схемы.
| Функция | Description | Пример варианта использования |
|---|---|---|
ColumnSelection("col") |
Последнее значение столбца (без агрегирования) | Последняя категория поставщика, сквозная передача поля запроса |
ColumnSelection можно использовать с любым источником данных:
-
DeltaTableSource: возвращает последнее значение для ключа сущности с помощью соединения между точками во времени (без агрегирования окна поиска). -
StreamSource: возвращает последнее значение для ключа сущности из потока (без агрегирования окна поиска). -
RequestSource: передает значение, предоставленное во время вывода (или извлекается из помеченного кадра данных во время обучения).
from databricks.feature_engineering.entities import (
ColumnSelection, DeltaTableSource, Feature, FieldDefinition,
RequestSource, ScalarDataType,
)
delta_source = DeltaTableSource(
catalog_name="main", schema_name="feature_store", table_name="transactions",
)
request_source = RequestSource(
schema=[
FieldDefinition(name="session_duration", data_type=ScalarDataType.DOUBLE),
]
)
# ColumnSelection from a Delta table
latest_amount = Feature(
source=delta_source,
function=ColumnSelection("amount"),
entity=["user_id"],
timeseries_column="transaction_time",
name="latest_transaction_amount",
)
# ColumnSelection from a RequestSource
session_feature = Feature(
source=request_source,
function=ColumnSelection("session_duration"),
name="session_duration",
)
Пример: агрегирование и функции выбора столбцов
В следующем примере показаны функции, определенные в одном источнике данных.
from databricks.feature_engineering.entities import (
AggregationFunction, Feature, Sum, Avg, ApproxCountDistinct,
ColumnSelection, RollingWindow,
)
from datetime import timedelta
window = RollingWindow(window_duration=timedelta(days=7))
sum_feature = Feature(
source=source,
entity=["user_id"],
timeseries_column="event_time",
function=AggregationFunction(Sum(input="amount"), window),
)
avg_feature = Feature(
source=source,
entity=["user_id"],
timeseries_column="event_time",
function=AggregationFunction(Avg(input="amount"), window),
)
distinct_count = Feature(
source=source,
entity=["user_id"],
timeseries_column="event_time",
function=AggregationFunction(ApproxCountDistinct(input="product_id", relativeSD=0.01), window),
)
# Column selection (no aggregation, no time window)
latest_amount = Feature(
source=source,
function=ColumnSelection("amount"),
entity=["user_id"],
timeseries_column="event_time",
name="latest_amount",
)
Функции с условиями фильтра
Параметр filter_condition позволяет фильтровать строки из исходной таблицы перед вычислением агрегатов. Это функция в качестве предложения SQL WHERE , применяемого до группировки и агрегирования данных.
Note
filter_condition фильтрует строки перед агрегированием, например предложение SQL WHERE , примененное ранее GROUP BY. Она не изменяет степень детализации, которая всегда определяется определением entity функции.
Фильтры полезны при работе с большими исходными таблицами, которые включают супермножество данных, необходимых для вычислений функций, и свести к минимуму потребность в создании отдельных представлений на вершине этих таблиц.
from databricks.feature_engineering.entities import AggregationFunction, Sum, Count, RollingWindow, DeltaTableSource
from datetime import timedelta
# Source with filter applied at the source level
high_value_transactions = DeltaTableSource(
catalog_name="main",
schema_name="ecommerce",
table_name="transactions",
filter_condition="amount > 100", # Only transactions over $100
)
high_value_sales = Feature(
source=high_value_transactions,
entity=["user_id"],
timeseries_column="transaction_time",
function=AggregationFunction(Sum(input="amount"), RollingWindow(window_duration=timedelta(days=30))),
)
# Multiple conditions
completed_orders_source = DeltaTableSource(
catalog_name="main",
schema_name="ecommerce",
table_name="orders",
filter_condition="status = 'completed' AND payment_method = 'credit_card'",
)
completed_orders = Feature(
source=completed_orders_source,
entity=["user_id"],
timeseries_column="order_time",
function=AggregationFunction(Count(input="order_id"), RollingWindow(window_duration=timedelta(days=7))),
)
# Filter on a StreamSource
from databricks.feature_engineering.entities import StreamSource
purchase_stream = StreamSource(
full_name="main.ecommerce.transactions_stream",
filter_condition="value.event_type = 'purchase'",
)
purchase_total = Feature(
source=purchase_stream,
entity=["value.user_id"],
timeseries_column="value.event_time",
function=AggregationFunction(Sum(input="value.amount"), RollingWindow(window_duration=timedelta(hours=1))),
)
Источники данных
DeltaTableSource
DeltaTableSource является временным объектом Python, используемым для определения того, как функции вычисляются из исходной таблицы. Она не создает новую таблицу. Он задает конфигурацию для чтения данных и агрегирования функций.
DeltaTableSource(
catalog_name: str, # Required: Catalog name
schema_name: str, # Required: Schema name
table_name: str, # Required: Table name
filter_condition: Optional[str] = None, # Optional: SQL WHERE clause to filter source data
transformation_sql: Optional[str] = None, # Optional: SQL SELECT expression for column transformations
dataframe_schema: Optional[str] = None, # Required if transformation_sql is set: schema of the resulting DataFrame
)
Параметры:
-
catalog_name, :schema_nametable_nameопределение исходной таблицы Delta в каталоге Unity. -
filter_condition: предложение SQLWHERE, применяемое перед агрегированием. Пример:"status = 'completed'". -
transformation_sql: выражение SQLSELECT, применяемое к исходной таблице. Используйте это для переименования столбцов, типов приведения или вычислений производных столбцов перед агрегированием. Если опущено, все столбцы выбраны (*). Пример:"user_id, CAST(amount AS DOUBLE) AS amount, event_time". -
dataframe_schema: схема результирующего кадра данных после преобразований в формате JSON Spark StructType (изdf.schema.json()). Обязательный, еслиtransformation_sqlуказан параметр . Это указывает системе имена столбцов и типы, которые приводят к преобразованию.
Если задано оба filter_condition и transformation_sql задано, результирующий запрос: SELECT {transformation_sql} FROM {table} WHERE {filter_condition}
Note
timeseries_column (указанный в определении компонента, а не в DeltaTableSource) должен иметь тип TimestampType или DateType. Целые типы могут работать, но привести к потере точности для агрегатов временных периодов.
Пример. Использование transformation_sql для преобразований столбцов
source = DeltaTableSource(
catalog_name="main",
schema_name="analytics",
table_name="raw_events",
transformation_sql="user_id, CAST(price_cents AS DOUBLE) / 100 AS price, event_time",
filter_condition="event_type = 'purchase'",
dataframe_schema=spark.sql(
"SELECT user_id, CAST(price_cents AS DOUBLE) / 100 AS price, event_time FROM main.analytics.raw_events LIMIT 0"
).schema.json(),
)
Пример: извлечение transformation_sql и dataframe_schema извлечение из кадра данных PySpark
Вы можете написать преобразование в виде запроса PySpark, а затем извлечь схему из результирующего кадра данных:
df = spark.sql(f"""
SELECT user_id, CAST(amount AS DOUBLE) / 100 AS amount_dollars, event_time
FROM main.analytics.events
WHERE event_date >= date_sub(current_date(), 7)
LIMIT 0
""")
# Use df.schema.json() as the dataframe_schema
source = DeltaTableSource(
catalog_name="main",
schema_name="analytics",
table_name="events",
transformation_sql="user_id, CAST(amount AS DOUBLE) / 100 AS amount_dollars, event_time",
filter_condition="event_date >= date_sub(current_date(), 7)",
dataframe_schema=df.schema.json(),
)
Note
transformation_sql поддерживает только выражения со строками (переименование столбцов, приведение, арифметика). Функции агрегирования, такие как COUNT(*) или SUM() не поддерживаются. Вместо этого используйте AggregationFunction определение функции.
DeltaTableSource.from_sql()
В качестве удобства можно создать из DeltaTableSource SQL-запроса. Метод анализирует запрос для автоматического извлечения имени таблицы и transformation_sqlfilter_condition.
DeltaTableSource.from_sql(
sql: str, # Required: SQL SELECT query
spark: SparkSession, # Required: active SparkSession (for schema inference)
) -> DeltaTableSource
Поддерживаются только простые SELECT ... FROM ... [WHERE ...] запросы. Сложный SQL (JOINs, вложенные запросы, CTEs, UNIONs) отклоняется. Для сложных запросов создайте DeltaTableSource напрямую с transformation_sql помощью и filter_condition.
from databricks.feature_engineering.entities import (
AggregationFunction,
DeltaTableSource,
Feature,
Sum,
TumblingWindow,
)
source = DeltaTableSource.from_sql(
spark=spark,
sql=f"SELECT customer_id, event_ts, amount * 2 AS doubled_amount, amount FROM {CATALOG}.{SCHEMA}.{TABLE}",
)
feature = Feature(
source=source,
function=AggregationFunction(Sum(input="doubled_amount"), time_window=TumblingWindow(window_duration=timedelta(days=7))),
entity=["customer_id"], timeseries_column="event_ts",
)
Итерации с to_dataframe()
Используется source.to_dataframe() для предварительного просмотра данных, которые будут использоваться для вычислений функций. Это полезно для итерации filter_condition и transformation_sql до тех пор, пока они не будут получать ожидаемые результаты.
source = DeltaTableSource(
catalog_name="main",
schema_name="analytics",
table_name="events",
filter_condition="event_type = 'purchase'",
)
# Preview the filtered source data
source.to_dataframe().display()
Общие сведения о сущностях
Столбцы сущностей определяют уровень агрегирования для ваших функций. Они указываются в определении Feature , а не в DeltaTableSource. Сущности определяют:
-
Как сгруппированы данные: функции агрегируются на уникальное сочетание значений сущностей (аналогично
GROUP BYв SQL) - Структура первичного ключа: каждая уникальная комбинация сущностей приводит к одной строке вычисляемых функций.
Пример: функции уровня клиента
Следующие функции кода агрегируются на уровне клиента (одна строка для каждого клиента):
from databricks.feature_engineering.entities import DeltaTableSource
source = DeltaTableSource(
catalog_name="main",
schema_name="analytics",
table_name="user_events",
)
Feature(
source=source,
entity=["user_id"], # Features aggregated per user
timeseries_column="event_time", # Timestamp for time windows
function=AggregationFunction(Sum(input="amount"), RollingWindow(window_duration=timedelta(days=7))),
)
Пример: функции уровня Customer-Store
Чтобы агрегировать функции на более подробном уровне (одна строка для сочетания хранилища клиентов), используйте несколько столбцов сущностей:
source = DeltaTableSource(
catalog_name="main",
schema_name="retail",
table_name="transactions",
)
Feature(
source=source,
entity=["user_id", "store_id"], # Features aggregated per user-store pair
timeseries_column="transaction_time",
function=AggregationFunction(Sum(input="amount"), RollingWindow(window_duration=timedelta(days=7))),
)
Если вам нужны функции на разных уровнях агрегирования (например, уровень клиента и уровень хранилища клиентов), используйте различные entity значения в определениях функций. Одно и то же DeltaTableSource можно совместно использовать между функциями с различными конфигурациями сущностей.
StreamSource
StreamSource ссылается на Stream. Поток содержит конфигурацию подключения, проверки подлинности, схемы и приема для источника потоковой передачи. Для Kafka ссылки на столбцы в определениях компонентов должны быть префиксированы или value.key. указывать, какая часть сообщения должна быть прочитана.
StreamSource(
full_name: str, # Required: Three-part Stream name (catalog.schema.stream)
filter_condition: Optional[str], # Optional: SQL WHERE clause applied before aggregation
)
Параметры:
-
full_name: полное трехкомпонентное имя потока (например,"my_catalog.my_schema.my_stream"). -
filter_condition(необязательно): предложение SQLWHERE, применяемое к потоковым данным перед агрегированием, с помощью ссылок на столбцы с префиксом точек (например,"value.event_type = 'purchase'").
from databricks.feature_engineering.entities import StreamSource
stream_source = StreamSource(
full_name="my_catalog.my_schema.my_stream",
filter_condition="value.event_type = 'purchase'",
)
RequestSource
RequestSource определяет схему для данных, предоставляемых во время вывода в полезных данных запроса, а не просматриваемых из предварительно материализованной таблицы. Во время обучения эти столбцы извлекаются из помеченного кадра данных, переданного create_training_setв . Во время обслуживания модели вызывающий объект должен включать их в полезные данные HTTP-запроса.
RequestSource используется с ColumnSelection (для передачи значения напрямую). Она не поддерживает функции агрегирования или временные окна.
Определение схемы
Определите схему как список объектов, каждый из которых задает имя столбца FieldDefinition и :ScalarDataType
from databricks.feature_engineering.entities import (
FieldDefinition, RequestSource, ScalarDataType,
)
request_source = RequestSource(
schema=[
FieldDefinition(name="transaction_amount", data_type=ScalarDataType.DOUBLE),
FieldDefinition(name="vendor_id", data_type=ScalarDataType.STRING),
FieldDefinition(name="transaction_id", data_type=ScalarDataType.STRING),
FieldDefinition(name="transaction_time", data_type=ScalarDataType.DATE),
]
)
Поддерживаемые типы данных
RequestSourceподдерживает скалярные типы, определенные в ScalarDataType: INTEGER, FLOATBOOLEANSTRINGDOUBLELONG, . TIMESTAMPDATESHORT Сложные типы, такие как массивы, карты и структуры, не поддерживаются.
Как данные запроса гидратируются
| Контекст | Behavior |
|---|---|
Обучение (create_training_set) |
Столбцы извлекаются из помеченного кадра данных. Типы проверяются на основе объявленной схемы. Несовпадения типов приводят к ошибке (неявное приведение типов не выполняется). |
| Обслуживание (конечная точка модели) | Столбцы извлекаются из dataframe_records HTTP-запроса или dataframe_split из них. Значения JSON приведение к объявленным типам (например, номер JSON → DOUBLE). |
Подпись модели
Если модель регистрируется с помощью log_model набора обучения, включающего RequestSource функции, RequestSource столбцы добавляются в подпись модели MLflow в качестве необходимых входных данных. Это означает, что схема API конечной точки обслуживания отражает, какие вызывающие поля должны предоставляться во время вывода.
API обучения и вывода
create_training_set и score_batch вычисляемые значения функций по запросу от исходных данных. Для функций, поддерживающих автономную материализацию, например агрегирование скользящих окон в источниках разностных таблиц, материализация функций сначала в автономном хранилище повышает производительность обеих операций. Если доступны материализованные автономные функции, операции считывают предварительно вычисляемые автономные данные вместо повторной компиляции значений признаков из источника. См. статью "Материализация представлений функций", чтобы материализовать функции в автономном хранилище.
create_training_set()
Создает обучающий набор данных с правильным вычислением функции на определенный момент времени. Дополнительные сведения см. в разделе "Обучение моделей с представлениями компонентов".
FeatureEngineeringClient.create_training_set(
df: DataFrame, # DataFrame with training data
features: Optional[List[Feature]], # List of Feature objects
label: Union[str, List[str], None], # Label column name(s)
exclude_columns: Optional[List[str]] = None, # Optional: columns to exclude
) -> TrainingSet
log_model()
Регистрирует модель с метаданными компонентов для отслеживания происхождения и автоматического поиска признаков во время вывода. Дополнительные сведения см. в разделе "Обучение моделей с представлениями компонентов".
FeatureEngineeringClient.log_model(
model, # Trained model object
artifact_path: str, # Path to store model artifact
flavor: ModuleType, # MLflow flavor module (e.g., mlflow.sklearn)
training_set: TrainingSet, # TrainingSet used for training
registered_model_name: Optional[str], # Optional: register model in Unity Catalog
)
score_batch()
Выполняет вывод автономной пакетной службы с автоматическим поиском компонентов. Использует метаданные компонентов, хранящиеся в модели, для вычисления правильных функций на определенный момент времени, обеспечивая согласованность с обучением.
FeatureEngineeringClient.score_batch(
model_uri: str, # URI of logged model (e.g., "models:/catalog.schema.model/1")
df: DataFrame, # DataFrame with entity keys and timestamps
) -> DataFrame
Входной кадр данных должен содержать столбцы сущности и времени, используемые во время обучения. Функции автоматически вычисляются из исходных данных.
fe = FeatureEngineeringClient()
# Batch scoring with automatic feature lookup
predictions = fe.score_batch(
model_uri="models:/main.ecommerce.fraud_model/1",
df=inference_df,
)
predictions.display()
Окна времени
Представления функций поддерживают три различных типа окон для управления поведением обратного просмотра для агрегатов на основе временных окон: скользящей, переворачивающейся и скользящей.
- Переворачивая окна от времени события. Длительность и задержка определяются явным образом.
- Скользящие окна являются фиксированными, не перекрывающимися временными интервалами. Каждая точка данных принадлежит ровно одному окну.
- Скользящие окна — это перекрывающиеся, временные окна с возможностью настройки интервала скольжения.
На следующем рисунке показано, как они работают.
Скользящей окна
Note
RollingWindow ранее был назван ContinuousWindow. При переходе с более ранней версии пакета SDK обновите импорт соответствующим образом.
Скользящие окна — это агрегаты up-toдаты и реального времени, которые обычно используются для потоковой передачи данных. В конвейерах потоковой передачи скользяшее окно выдает новую строку только при изменении содержимого окна фиксированной длины, например при входе или выходе события. При использовании функции скользящего окна в конвейерах обучения точный расчет функции на определенный момент времени выполняется для исходных данных с использованием длительности окна фиксированной длины, непосредственно предшествующей метке времени конкретного события. Это помогает предотвратить искажение данных между онлайн и оффлайн режимами или их утечку. Характеристики в момент T агрегируют события в интервале [T – продолжительность, T).
class RollingWindow(TimeWindow):
window_duration: datetime.timedelta
delay: Optional[datetime.timedelta] = None
В следующей таблице перечислены параметры для скользящего окна. Время начала и окончания окна основано на следующих параметрах:
- Время начала:
evaluation_time - window_duration - delay(включительно) - Время окончания:
evaluation_time - delay(эксклюзивное)
| Parameter | Constraints |
|---|---|
delay (необязательно) |
Должно быть ≥ 0 (сдвигает окно назад во времени с метки времени оценки). Используйте delay для учета любой задержки системы между созданием события и меткой времени события, чтобы предотвратить утечку будущих событий в обучающие наборы данных. Например, если в течение одной минуты между созданными событиями и этими событиями в конечном итоге приземляются в исходную таблицу, в которой они назначают метку времени, задержка будет timedelta(minutes=1). |
window_duration |
Должно быть > 0 |
from databricks.feature_engineering.entities import RollingWindow
from datetime import timedelta
# Look back 7 days from evaluation time
window = RollingWindow(window_duration=timedelta(days=7))
Определите скользякое окно с задержкой, используя приведенный ниже код.
# Look back 7 days, offset by 1 minute to account for data ingestion delay
window = RollingWindow(
window_duration=timedelta(days=7),
delay=timedelta(minutes=1)
)
Примеры скользящего окна
window_duration=timedelta(days=7): это создает 7-дневное окно обратного просмотра, заканчивающееся текущим временем оценки. Для мероприятия в 14:00 в День 7, это включает все события с 14:00 в День 0 до (но не включая) 14:00 в День 7.window_duration=timedelta(hours=1), delay=timedelta(minutes=30): создается 1-часовое окно обратного просмотра, заканчивающееся 30 минут до времени оценки. Для события в 3:00 это включает все события с 1:30 до (но не включая) 2:30 вечера. Это полезно для учета задержек приема данных.
"Переворачивающееся" окно
Для функций, определенных с использованием скользящих окон, агрегации вычисляются по предварительно определенному окну фиксированной длины, которое сдвигается на заданный интервал, создавая неперекрывающиеся окна, которые полностью распределяют время. В результате каждое событие в источнике способствует ровно одному окну. Характеристики во время t агрегируют данные из окон, которые заканчиваются до t, исключая t. Windows начинается с эпохи Unix.
class TumblingWindow(TimeWindow):
window_duration: datetime.timedelta
В следующей таблице перечислены параметры для переворачивающегося окна.
| Parameter | Constraints |
|---|---|
window_duration |
Должно быть > 0 |
from databricks.feature_engineering.entities import TumblingWindow
from datetime import timedelta
window = TumblingWindow(
window_duration=timedelta(days=7)
)
Пример "Переворачивающееся окно"
-
window_duration=timedelta(days=5): при этом создаются предопределенные окна фиксированной длины по 5 дней. Пример: окно #1 охватывает день 0 до дня 4, окно #2 охватывает день 5 до дня 9, окно #3 охватывает день 10 до дня 14 и т. д. В частности, окно #1 включает все события со метками времени, начиная с00:00:00.00дня 0 до (но не включая) какие-либо события с меткой времени00:00:00.00на день 5. Каждое событие принадлежит ровно одному окну.
"Скользящее" окно
Для функций, определенных с использованием скользящих окон, агрегации вычисляются в рамках предварительно заданного окна фиксированной длины, которое сдвигается с каждым интервалом, создавая перекрывающиеся окна. Каждое событие в источнике может способствовать агрегации компонентов для нескольких окон. Характеристики во время t агрегируют данные из окон, которые заканчиваются до t, исключая t. Windows начинается с эпохи Unix.
class SlidingWindow(TimeWindow):
window_duration: datetime.timedelta
slide_duration: datetime.timedelta
В следующей таблице перечислены параметры скользящего окна.
| Parameter | Constraints |
|---|---|
window_duration |
Должно быть > 0 |
slide_duration |
Должно быть > 0 и <window_duration |
from databricks.feature_engineering.entities import SlidingWindow
from datetime import timedelta
window = SlidingWindow(
window_duration=timedelta(days=7),
slide_duration=timedelta(days=1)
)
Пример скользящего окна
-
window_duration=timedelta(days=5), slide_duration=timedelta(days=1): это создает перекрывающиеся 5-дневные окна, которые перемещаются на 1 день каждый раз. Пример: окно #1 охватывает день 0 до дня 4, окно #2 охватывает день 1 до дня 5, окно #3 охватывает день 2 до дня 6 и т. д. Каждое окно включает события с00:00:00.00начального дня до (но не включая)00:00:00.00в конечный день. Так как окна перекрываются, одно событие может принадлежать нескольким окнам (в этом примере каждое событие принадлежит до 5 разных окон).
Триггеры материализации
Активирует управление при запуске конвейера материализации. Тип триггера зависит от типа компонента.
CronSchedule
Используется CronSchedule для функций агрегирования (AggregationFunction). Конвейер выполняется по фиксированному расписанию, определенному выражением cron Изумения.
from databricks.feature_engineering.entities import CronSchedule
trigger = CronSchedule(
quartz_cron_expression="0 0 * * * ?", # Hourly
timezone_id="UTC",
)
TableTrigger
Используется TableTrigger для ColumnSelection функций, поддерживаемых a DeltaTableSource. Конвейер выполняется всякий раз, когда вышестоящей таблице Delta получает новую фиксацию.
from databricks.feature_engineering.entities import TableTrigger
trigger = TableTrigger()
StreamingMode
Используется StreamingMode для функций, поддерживаемых a StreamSource. Конвейер выполняется как конвейер непрерывной потоковой передачи.
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
StreamSource, Feature, AggregationFunction, Sum,
RollingWindow, OnlineStoreConfig, StreamingMode,
)
from datetime import timedelta
fe = FeatureEngineeringClient()
stream_source = StreamSource(full_name="my_catalog.my_schema.my_stream")
streaming_feature = fe.create_feature(
source=stream_source,
entity=["value.user_id"],
timeseries_column="value.event_time",
function=AggregationFunction(
operator=Sum(input="value.amount"),
time_window=RollingWindow(window_duration=timedelta(hours=1)),
),
catalog_name="my_catalog",
schema_name="my_schema",
name="user_purchase_sum",
)
fe.materialize_features(
features=[streaming_feature],
online_config=OnlineStoreConfig(
catalog_name="my_catalog",
schema_name="my_schema",
table_name_prefix="streaming_features_serving",
online_store_name="feature_store_online",
),
trigger=StreamingMode(),
)
Выбор триггера
| Тип компонента | Trigger | При запуске |
|---|---|---|
Агрегирование (AggregationFunction) из DeltaTableSource |
CronSchedule |
По фиксированному расписанию cron |
ColumnSelection (из DeltaTableSource) |
TableTrigger |
Фиксация каждой исходной таблицы |
Функции из StreamSource |
StreamingMode |
Непрерывная потоковая передача |
Невозможно материализовать функции, требующие разных типов триггеров в одном materialize_features вызове. Вместо этого нужно выдавать отдельные вызовы.
Перенос бета-функций в общедоступную предварительную версию
Общедоступная предварительная версия представлений функций содержит сущности компонентов первого класса в каталоге Unity, управляемые CREATE FEATURE правами и READ FEATURE требуя databricks-feature-engineering версии 0.16.0 или более поздней. Функции, созданные во время бета-версии (с версией 0.15.0), хранятся как функции каталога Unity и не поддерживают все функции общедоступной предварительной версии. Чтобы получить поддержку долгосрочной общедоступной предварительной версии, создайте бета-версии с версией 0.16.0. Компоненты должны быть удалены и повторно созданы, а не только повторно материализованы.
Дополнительные сведения о функциях см. в разделе "Представления компонентов".
Что вам нужно сделать
- Обновление до версии 0.16.0. Это необходимая версия клиента для функций общедоступной предварительной версии (пакетная версия и потоковая передача).
- Повторно создайте функции. Представления компонентов бета-версии должны быть удалены и повторно созданы, а не повторно материализованы, так как они не поддерживают все функции общедоступной предварительной версии.
- Миграция до закрытия окна. Существующие бета-версии должны быть перенесены до 22 июля 2026 г.
Определение функций бета-версии и общедоступной предварительной версии
Функции общедоступной предварительной версии отображаются как объект компонента в каталоге Unity, например в обозревателе каталогов. Функции бета-версии отображаются как функция с определением YAML. Любая функция, представленная как функция, является бета-функцией, которую необходимо перенести.
Перенос бета-функций
Перенос бета-функции состоит из трех частей:
- Повторно создайте функцию как общедоступную предварительную версию.
- Повторно материализуйте эту функцию, поэтому ее автономные и онлайн-таблицы перестроены под новой функцией.
- После проверки перенесенных функций удалите бета-функции и их материализации.
Повторное создание функций
Используется list_beta_feature_views для поиска бета-функций, Feature.clone() создания незарегистрированной копии и register_feature повторной регистрации каждой копии в качестве общедоступной предварительной версии. Клонирование очищает регистрацию, каталог и схему, чтобы компонент можно было повторно зарегистрировать.
Чтобы избежать конфликтов имен, зарегистрируйте перенесенные функции с другим именем или в схеме, отличной от бета-версий. Следующий пример повторно регистрирует каждую функцию в исходной схеме с суффиксом _migrated имени.
# Update this to the catalog whose beta Feature Views you want to migrate.
CATALOG_TO_MIGRATE = "main"
from databricks.feature_engineering import FeatureEngineeringClient
fe = FeatureEngineeringClient()
# 1. Find every beta Feature View in the catalog. Returns Feature objects,
# scanned across all schemas in the catalog.
beta_features = fe.list_beta_feature_views(catalog_name=CATALOG_TO_MIGRATE)
# Keep each beta feature paired with its migrated counterpart for the next steps.
migrations = []
for beta_feature in beta_features:
catalog_name, schema_name, leaf_name = beta_feature.full_name.split(".")
# 2. Clone the feature as an unregistered copy, renamed with a "_migrated" suffix.
cloned = beta_feature.clone(new_name=f"{leaf_name}_migrated")
# 3. Re-register the clone as a Public Preview feature.
migrated = fe.register_feature(
feature=cloned,
catalog_name=catalog_name,
schema_name=schema_name,
)
migrations.append((beta_feature, migrated))
Повторно материализуйте перенесенные функции
Если бета-версия компонента была материализована, повторно материализуйте свой аналог Public Preview, чтобы его автономные и онлайн-таблицы были перестроены под новой функцией. Предоставьте конфигурации автономного и интернет-магазина для перенесенной функции и перестроите триггер из существующей материализации бета-функции.
from databricks.feature_engineering.entities import (
CronSchedule,
OfflineStoreConfig,
OnlineStoreConfig,
TableTrigger,
)
for beta_feature, migrated in migrations:
# Inspect the beta feature's existing materializations to see what to rebuild and
# to reconstruct the same trigger.
trigger = None
needs_offline = needs_online = False
for mf in fe.list_materialized_features(feature_name=beta_feature.full_name):
needs_online = needs_online or bool(mf.is_online)
needs_offline = needs_offline or not mf.is_online
# Rebuild the trigger from the materialized feature.
if mf.cron_schedule_trigger is not None:
trigger = CronSchedule(
quartz_cron_expression=mf.cron_schedule_trigger.cron_expression,
timezone_id="UTC", # Materialized schedules run in UTC.
)
elif mf.table_trigger is not None:
trigger = TableTrigger()
elif mf.streaming_mode is not None:
# Streaming features use StreamingMode, which can be reused as-is.
trigger = mf.streaming_mode
if not (needs_offline or needs_online):
continue # The beta feature was never materialized.
catalog_name, schema_name, _ = migrated.full_name.split(".")
fe.materialize_features(
features=[migrated],
offline_config=OfflineStoreConfig(
catalog_name=catalog_name,
schema_name=schema_name,
table_name_prefix="migrated_features",
)
if needs_offline
else None,
online_config=OnlineStoreConfig(
catalog_name=catalog_name,
schema_name=schema_name,
table_name_prefix="migrated_features",
online_store_name="my_online_store",
)
if needs_online
else None,
trigger=trigger,
)
Note
Материализация каждой функции в собственном materialize_features вызове создает отдельный конвейер. Чтобы сократить затраты на вычисление, группируйте функции, которые совместно используют место для автономного и интерактивного назначения и запускают один materialize_features вызов, передав их вместе features.
Удаление бета-функций
Предупреждение
Удалите бета-функции и их материализации только после проверки правильности перенесенных функций и их материализованных данных. Удаление необратимо.
После проверки перенесенных функций удалите материализации каждого бета-компонента, а затем бета-версию.
for beta_feature, _ in migrations:
# Delete the beta feature's materializations first.
mfs = list(fe.list_materialized_features(feature_name=beta_feature.full_name))
offline_mfs = [mf for mf in mfs if not mf.is_online]
if offline_mfs:
# Aggregation features pair an offline and online table; deleting the offline
# materialized feature removes its paired online table too.
for mf in offline_mfs:
fe.delete_materialized_feature(materialized_feature=mf)
else:
# Online-only features (ColumnSelection, streaming) have no offline pair; delete
# the online materialized feature directly.
for mf in mfs:
fe.delete_materialized_feature(materialized_feature=mf)
# Then delete the beta feature definition.
fe.delete_feature(full_name=beta_feature.full_name)