Справочник по API Feature Views

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: предложение SQL WHERE , применяемое перед агрегированием. Пример: "status = 'completed'".
  • transformation_sql: выражение SQL SELECT , применяемое к исходной таблице. Используйте это для переименования столбцов, типов приведения или вычислений производных столбцов перед агрегированием. Если опущено, все столбцы выбраны (*). Пример: "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 (необязательно): предложение SQL WHERE , применяемое к потоковым данным перед агрегированием, с помощью ссылок на столбцы с префиксом точек (например, "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)