Представления функций

Important

Эта функция доступна в общедоступной предварительной версии. Администраторы рабочей области могут управлять доступом к этой функции на странице "Предварительные версии ". См. статью "Управление предварительными версиями Azure Databricks".

Представления функций позволяют определять и вычислять функции из источников данных. Признаки можно определять с использованием различных источников (таблицы Delta, потока Kafka и данных, доступных во время запроса) и вычислений (агрегирования по временным окнам, простых выборок столбцов и многого другого). В этом руководстве рассматриваются следующие рабочие процессы:

  • Рабочий процесс разработки компонентов
    • Используется create_feature для определения объектов функций каталога Unity, которые можно использовать в обучении модели и обслуживании рабочих процессов.
    • Кроме того, создайте объекты Feature локально и позже используйте register_feature для их сохранения в каталоге Unity. Локально созданные особенности можно использовать с create_training_set до регистрации.
  • Рабочий процесс обучения модели
  • Материализация компонентов и рабочий процесс обслуживания
    • После определения функции с помощью create_feature или получения ее с помощью get_feature, можно использовать materialize_features для материализации функции или набора функций в автономное хранилище для эффективного повторного использования или в онлайн-хранилище для обслуживания в режиме онлайн.
    • Используйте create_training_set с материализованным представлением для подготовки офлайн набора данных для пакетного обучения.

Подробные сведения об API см. в справочнике по API Feature Views.

Требования

  • Бессерверные вычислительные ресурсы или кластер классических вычислительных ресурсов, использующий Databricks Runtime 17.0 ML или выше.

  • Необходимо установить пользовательский пакет Python. Выполните следующие строки кода при каждом запуске записной книжки:

    %pip install databricks-feature-engineering>=0.16.0
    dbutils.library.restartPython()
    

Пример быстрого начала

Для записной книжки быстрого запуска см. пример записной книжки.

from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
    CronSchedule, DeltaTableSource, Feature, AggregationFunction,
    Sum, Avg, ColumnSelection, TableTrigger,
    TumblingWindow, SlidingWindow,
    OfflineStoreConfig, OnlineStoreConfig,
)
from datetime import timedelta

CATALOG_NAME = "main"
SCHEMA_NAME = "feature_store"
TABLE_NAME = "transactions"

# 1. Create data source
source = DeltaTableSource(
    catalog_name=CATALOG_NAME,
    schema_name=SCHEMA_NAME,
    table_name=TABLE_NAME,
)

# 2. Define features locally (no catalog/schema needed yet)
avg_feature = Feature(
    source=source,
    entity=["user_id"],
    timeseries_column="transaction_time",
    function=AggregationFunction(Avg(input="amount"), TumblingWindow(window_duration=timedelta(days=30))),
    name="avg_transaction_30d",
)

sum_feature = Feature(
    source=source,
    entity=["user_id"],
    timeseries_column="transaction_time",
    function=AggregationFunction(Sum(input="amount"), SlidingWindow(window_duration=timedelta(days=7), slide_duration=timedelta(days=1))),
    # name auto-generated: "amount_sum_sliding_7d_1d"
)

fe = FeatureEngineeringClient()

# 3. Explore features with compute_features
feature_df = fe.compute_features(features=[avg_feature, sum_feature])
feature_df.display()

# 4. Create training set using local features
# `labeled_df` should have columns "user_id", "transaction_time", and "target".
training_set = fe.create_training_set(
    df=labeled_df,
    features=[avg_feature, sum_feature],
    label="target",
)
training_set.load_df().display()

# 5. Register features in Unity Catalog
avg_feature = fe.register_feature(
    feature=avg_feature,
    catalog_name=CATALOG_NAME,
    schema_name=SCHEMA_NAME,
)
sum_feature = fe.register_feature(
    feature=sum_feature,
    catalog_name=CATALOG_NAME,
    schema_name=SCHEMA_NAME,
)

# 6. Or use create_feature for a one-step define-and-register workflow
latest_amount = fe.create_feature(
    source=source,
    function=ColumnSelection("amount"),
    entity=["user_id"],
    timeseries_column="transaction_time",
    catalog_name=CATALOG_NAME,
    schema_name=SCHEMA_NAME,
    name="latest_amount",
)

# 7. Train model
with mlflow.start_run():
    training_df = training_set.load_df()

    # training code

    fe.log_model(
        model=model,
        artifact_path="recommendation_model",
        flavor=mlflow.sklearn,
        training_set=training_set,
        registered_model_name=f"{CATALOG_NAME}.{SCHEMA_NAME}.recommendation_model",
    )

# 8. (Optional) Materialize features for serving
# Features must be registered in UC before calling materialize_features
online_config = OnlineStoreConfig(
    catalog_name=CATALOG_NAME,
    schema_name=SCHEMA_NAME,
    table_name_prefix="customer_features_serving",
    online_store_name="customer_features_store",
)

# Aggregation features use CronSchedule and support both offline and online configs
fe.materialize_features(
    features=[avg_feature, sum_feature],
    offline_config=OfflineStoreConfig(
        catalog_name=CATALOG_NAME,
        schema_name=SCHEMA_NAME,
        table_name_prefix="customer_features",
    ),
    online_config=online_config,
    trigger=CronSchedule(
        quartz_cron_expression="0 0 * * * ?",  # Hourly
        timezone_id="UTC",
    ),
)

# ColumnSelection features use TableTrigger and only support online config
fe.materialize_features(
    features=[latest_amount],
    online_config=online_config,
    trigger=TableTrigger(),
)

пример записной книжки

Ноутбук для быстрого запуска представлений функций

Получите ноутбук

Функции потоковой передачи

Помимо пакетных функций из таблиц Delta можно определить функции из источников потоковой передачи для вариантов использования в режиме реального времени. Функции потоковой передачи используют тот же класс компонентов, что и пакетные функции— Feature одни и те же конструкторы, те же функции агрегирования, те же рабочие процессы обучения и обслуживания, поэтому для обновления пакета до реального времени требуются минимальные изменения кода. После материализации потоковые признаки обеспечивают сквозную актуальность менее чем за секунду (задержка p99 — 200 мс) непосредственно в конечных точках обслуживания моделей.

Чтобы использовать функции потоковой передачи, сначала настройте поток, а затем сошлитесь на него с помощью StreamSource. Источники потоковых данных поддерживают Kafka в качестве источника входных данных и автоматически поддерживают таблицу загрузки (Delta) в качестве исторической копии данных для обучения.

Определение функции потоковой передачи

StreamSource ссылается на Stream по его трёхкомпонентному имени (catalog.schema.stream_name). Поток не является защищаемым объектом Unity Catalog, но относится к схеме Unity Catalog, а доступ к нему регулируется таблицей приема данных потока. Ссылки на столбцы в определениях сущностей, временных рядов и функций должны иметь префикс value. или key., чтобы указать, какую часть сообщения Kafka следует читать. Вложенные поля поддерживаются с помощью нотации точек (например, value.user.address.city).

from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
    StreamSource,
    Feature,
    AggregationFunction,
    Sum,
    RollingWindow,
)
from datetime import timedelta

client = FeatureEngineeringClient()

stream_source = StreamSource(
    full_name="my_catalog.my_schema.my_stream",
)

feature = Feature(
    name="user_purchase_sum",
    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)),
    ),
)

Фильтрация условий в StreamSource

Используется filter_condition для фильтрации строк из потока перед агрегированием, как и в DeltaTableSource.

stream_source = StreamSource(
    full_name="my_catalog.my_schema.my_stream",
    filter_condition="value.event_type = 'purchase'",
)

Выбор столбцов из потоков данных

ColumnSelection функции работают с источниками потоковой передачи. Выбранный столбец содержит последнее значение из потока для каждой сущности с соблюдением точности на определённый момент времени.

from databricks.feature_engineering.entities import ColumnSelection

passenger_count = Feature(
    name="passenger_count",
    source=stream_source,
    entity=["value.user_id"],
    timeseries_column="value.event_time",
    function=ColumnSelection(column="value.passenger_count"),
)

Доступ к вложенным полям

Доступ к вложенным полям JSON можно получить с помощью нотации точек (например, value.nested_field.amount). Во время обслуживания полезные данные запроса и ответ используют имена конечных узлов (например, amount вместо value.amount). Имена листовых узлов должны быть уникальными среди всех столбцов вывода сущностей, временных рядов и признаков в рамках модели или Feature Spec, поскольку конечная точка сервинга использует имена листовых узлов для маршрутизации значений.

Временные окна для функций потоковой передачи

Возможности потоковой передачи поддерживают только RollingWindow при агрегировании. Скользящие окна постоянно пересчитываются на основе самых последних данных, что соответствует природе потоковых источников, работающих в реальном времени. TumblingWindow и SlidingWindow предназначены для пакетного вычисления через фиксированные исторические интервалы.

Пример записной книжки для функций потоковой передачи

Записная книжка быстрого запуска представлений функций потоковой передачи

Получите ноутбук

Обучение модели и вывод

Обучение моделей и выполнение пакетного вывода с помощью представлений компонентов, включая log_model(), score_batch()и create_training_set(), см. в разделе "Обучение моделей с помощью представлений функций".

Материализация характеристик

После определения функций их можно материализовать в автономные или интернет-магазины для эффективного использования в рабочих процессах обучения и обслуживания. После материализации функций можно обслуживать модели с помощью обслуживания модели ЦП. Дополнительные сведения см. в разделе "Материализация представлений функций".

Лучшие практики

Именование компонентов

  • Используйте описательные имена для критически важных для бизнеса функций.
  • Следуйте согласованным соглашениям об именовании между командами.
  • При разработке функций используйте автоматически созданные имена .

Окна времени

  • Согласуйте границы окон с бизнес-циклами (ежедневными, еженедельными).
  • Короткие окна захватывают последние тенденции, но могут быть шумными. Более длинные окна создают более стабильные дистрибутивы функций, но могут пропустить последние сдвиги поведения. Выберите, исходя из того, как быстро изменяется базовый сигнал в вашем случае использования. Например, 7-дневное окно сглаживает ежедневные колебания и создает согласованные входные данные модели, в то время как 1-часовое окно реагирует быстро на изменения поведения, но может привести к дисперсии, которая снижает производительность модели. Если точность модели снижается при смене распределения, используйте более длинное окно для стабилизации входных данных.
  • Переворачивающиеся и скользящие окна являются более масштабируемыми, чем скользящие (непрерывные) окна. Начните с скользящих окон для большинства вариантов использования.

Performance

  • Материализуйте функции из одного источника данных в одном materialize_features вызове, чтобы свести к минимуму сканирование данных.
  • Используйте тот же уровень детализации (например, все продолжительности слайдов по 1 часу или 1 дню) для функций одного источника данных, чтобы повысить их группировку во время материализации.

Столбцы сущностей vs условия фильтра

Используйте это руководство по принятию решений при работе с функциями из той же исходной таблицы:

Используйте entity (on create_feature) при необходимости различных уровней агрегирования:

  • Функции уровня клиента (одна строка для каждого клиента): entity=["customer_id"]
  • Функции клиента-торговца (несколько строк для каждого клиента): entity=["customer_id", "merchant_id"]
  • Разные уровни агрегирования могут совместно использовать одно и то же DeltaTableSource: укажите разные entity значения для каждого определения компонента.

Используйте filter_condition (on DeltaTableSource) при необходимости фильтрации строк на одном уровне агрегирования:

  • Только высокозначные транзакции: filter_condition="amount > 100" (по-прежнему агрегированные на каждого клиента)
  • Только завершенные заказы: filter_condition="status = 'completed'" (по-прежнему агрегированные на каждого клиента)

Правило большого пальца: Если изменение приведет к другому количеству строк на значение сущности, используйте разные entity значения в определениях признаков. Если вы просто фильтруете строки, которые вносят вклад в ту же агрегацию, используйте filter_condition в исходных данных.

Распространенные шаблоны

Аналитика клиентов

from databricks.feature_engineering.entities import AggregationFunction, Sum, Count, RollingWindow

fe = FeatureEngineeringClient()
features = [
    # Recency: Number of transactions in the last day
    fe.create_feature(catalog_name="main", schema_name="ecommerce", source=transactions,
            entity=["user_id"], timeseries_column="transaction_time",
            function=AggregationFunction(Count(input="transaction_id"), RollingWindow(window_duration=timedelta(days=1)))),

    # Frequency: transaction count over the last 90 days
    fe.create_feature(catalog_name="main", schema_name="ecommerce", source=transactions,
            entity=["user_id"], timeseries_column="transaction_time",
            function=AggregationFunction(Count(input="transaction_id"), RollingWindow(window_duration=timedelta(days=90)))),

    # Monetary: total spend in the last month
    fe.create_feature(catalog_name="main", schema_name="ecommerce", source=transactions,
            entity=["user_id"], timeseries_column="transaction_time",
            function=AggregationFunction(Sum(input="amount"), RollingWindow(window_duration=timedelta(days=30)))),
]

Анализ трендов

# Compare recent vs. historical behavior
fe = FeatureEngineeringClient()
recent_avg = fe.create_feature(
    catalog_name="main", schema_name="ecommerce",
    source=transactions, entity=["user_id"], timeseries_column="transaction_time",
    function=AggregationFunction(Avg(input="amount"), RollingWindow(window_duration=timedelta(days=7))),
)

historical_avg = fe.create_feature(
    catalog_name="main", schema_name="ecommerce",
    source=transactions, entity=["user_id"], timeseries_column="transaction_time",
    function=AggregationFunction(Avg(input="amount"), RollingWindow(window_duration=timedelta(days=7), delay=timedelta(days=7))),
)

Сезонные шаблоны

# Same day of week, 4 weeks ago
fe = FeatureEngineeringClient()
weekly_pattern = fe.create_feature(
    catalog_name="main", schema_name="ecommerce",
    source=transactions, entity=["user_id"], timeseries_column="transaction_time",
    function=AggregationFunction(Avg(input="amount"), RollingWindow(window_duration=timedelta(days=1), delay=timedelta(weeks=4))),
)

Ограничения

  • Имена столбцов объектов и временных рядов должны совпадать между обучающим набором данных (с метками) и определениями характеристик при использовании в create_training_set API.
  • Имя столбца, используемое в качестве label столбца в наборе данных обучения, не должно существовать в исходных таблицах, используемых для определения Feature.
  • Ограниченный список функций (UDAFs) поддерживается в create_feature API. См. раздел "Поддерживаемые функции".
  • Столбцы сущностей не могут быть типом DATE или TIMESTAMP.
  • RequestSource поддерживает только скалярные типы данных, определенные в ScalarDataType (INTEGER, FLOAT, BOOLEAN, STRING, DOUBLE, LONG, TIMESTAMP, DATE, SHORT). Сложные типы, такие как массивы, карты и структуры, не поддерживаются.
  • RequestSource не поддерживает функции агрегирования или временные окна. Можно использовать только ColumnSelection функции.
  • Набор имен столбцов сущностей, имен столбцов временных рядов и имен столбцов признаков запросов должен быть глобально уникальным для всех источников в наборе данных для обучения или конечной точке предоставления.
  • score_batch может не работать в бессерверной вычислительной среде. Это можно обойти с помощью классического вычислительного кластера под управлением Databricks Runtime 17.0 ML или более поздней версии.

Сведения об ограничениях, относящихся к материализации, см. в разделе "Ограничения".