Vistas de funciones

Importante

Esta característica está en versión preliminar pública. Los administradores del área de trabajo pueden controlar el acceso a esta característica desde la página Vistas previas . Consulte Administrar versiones preliminares de Azure Databricks.

Las vistas de características le permiten definir y calcular características de orígenes de datos. Las características se pueden definir mediante una variedad de orígenes (tabla Delta, flujo de Kafka y datos en tiempo de solicitud) y cálculos (agregaciones en período de tiempo, selecciones de columna simples, etc.). En esta guía se describen los siguientes flujos de trabajo:

  • Flujo de trabajo de desarrollo de funcionalidades
    • Use create_feature para definir objetos de características del catálogo de Unity que se pueden usar en el entrenamiento del modelo y servir flujos de trabajo.
    • Como alternativa, construya objetos Feature localmente y utilice register_feature para persistirlos en el Unity Catalog más adelante. Las características construidas localmente se pueden usar con create_training_set antes del registro.
  • Flujo de trabajo de entrenamiento del modelo
    • Use create_training_set para calcular características agregadas a un momento dado para el aprendizaje automático. Para obtener documentación detallada sobre el entrenamiento con vistas de características, consulte Entrenamiento de modelos con vistas de características.
  • Flujo de trabajo de materialización y provisión de características
    • Después de definir una característica con create_feature o recuperarla mediante get_feature, puede usar materialize_features para materializar la característica o el conjunto de características en un almacén sin conexión para una reutilización eficaz o para una tienda en línea para el servicio en línea.
    • Use create_training_set con la vista materializada para preparar un conjunto de datos de entrenamiento por lotes sin conexión.

Para obtener información detallada sobre la API, consulte Referencia de la API de vistas de características.

Requirements

  • Proceso sin servidor o un clúster de proceso clásico que ejecuta Databricks Runtime 17.0 ML o superior.

  • Debe instalar el paquete de Python personalizado. Ejecute las siguientes líneas de código cada vez que ejecute un cuaderno:

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

Ejemplo de inicio rápido

Para obtener un cuaderno de inicio rápido ejecutable, consulte Cuaderno de ejemplo.

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 support CronSchedule or TableTrigger, 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(),
)

Cuaderno de ejemplo de

Cuaderno de inicio rápido de vistas de atributos

Obtener el portátil

Características de streaming

Utiliza funciones de streaming cuando los valores de las funciones deban actualizarse continuamente en lugar de en un calendario por lotes. Las funciones de streaming y batch utilizan los mismos Feature constructores, funciones de agregación y flujos de trabajo de entrenamiento y servicio.

Las funciones de streaming no se materializan en una tienda offline. Para el entrenamiento y la inferencia por lotes, Databricks calcula valores de características a partir de la fuente.

Las funciones de streaming tienen los siguientes requisitos:

  • Debe proporcionar un online_config. Las funciones de streaming no soportan offline_config.
  • No puedes combinar funciones de streaming y lotes en una sola materialize_features llamada. Haz una llamada separada para cada tipo de disparador.
  • transformation_sql no está compatible con funciones de streaming.
  • La materialización en streaming solo procesa los registros que llegan después de que el oleoducto inicie y no rellena los registros históricos. Los agregados de ventana móvil solo devuelven resultados completos después de que llega la primera ventana completa de datos.

Elegir una fuente de función de streaming

Elige una fuente basada en tus necesidades de frescura y en la configuración de ingesta existente:

  • Usa un StreamSource cuando la actualización en menos de un segundo sea prioritaria. StreamSource Las características proporcionan una latencia extremo a extremo P99 de 200 milisegundos. Primero configura un flujo, luego haz referencia a él usando un StreamSource. Las fuentes de flujo admiten Kafka como entrada y mantienen automáticamente una tabla de ingestión (Delta) como copia histórica de los datos para el entrenamiento.
  • Use DeltaTableSource cuando ya tenga una ruta de ingestión de baja latencia hacia una tabla Delta. Espera una frescura del orden de decenas de segundos.
  • Usa Zerobus para llenar un DeltaTableSource cuando no tengas ya una ruta de ingestión de baja latencia. La ingestión de Zerobus tarda del orden de decenas de segundos, por lo que cabe esperar una actualización de las variables de menos de un minuto.

Define una función de streaming usando un StreamSource

Un StreamSource hace referencia a un Stream mediante su nombre de tres partes (catalog.schema.stream_name). Un objeto Stream no es un objeto protegible del catálogo de Unity, pero tiene como ámbito un esquema de catálogo de Unity y el acceso se rige por la tabla de ingesta de Stream. Las referencias de columna en las definiciones de entidad, series temporales y función deben ir precedidas de value. o key. para indicar qué parte del mensaje de Kafka se debe leer. Los campos anidados se admiten mediante la notación de puntos (por ejemplo, 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)),
    ),
)

Define una función de streaming usando una DeltaTableSource

Para materializar una característica definida en a DeltaTableSource como una característica de streaming, pase StreamingMode como el disparador a materialize_features. La definición de la funcionalidad utiliza las mismas API que una funcionalidad por lotes respaldada por un archivo DeltaTableSource. Las fuentes de la tabla Delta soportan agregación y funciones de selección de columnas.

La tabla Delta de origen debe tener activado el feed de datos de cambio (CDF) configurando delta.enableChangeDataFeed=true.

El siguiente ejemplo define y materializa una característica de agregación con una fuente de tabla Delta.

from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
    DeltaTableSource,
    AggregationFunction,
    OnlineStoreConfig,
    Sum,
    RollingWindow,
    StreamingMode,
)
from datetime import timedelta

client = FeatureEngineeringClient()

source = DeltaTableSource(
    catalog_name="my_catalog",
    schema_name="my_schema",
    table_name="transactions",
)

feature = client.create_feature(
    catalog_name="my_catalog",
    schema_name="my_schema",
    name="user_purchase_sum",
    source=source,
    entity=["user_id"],
    timeseries_column="event_time",
    function=AggregationFunction(
        operator=Sum(input="amount"),
        time_window=RollingWindow(window_duration=timedelta(hours=1)),
    ),
)

online_config = OnlineStoreConfig(
    catalog_name="my_catalog",
    schema_name="my_schema",
    table_name_prefix="streaming_features",
    online_store_name="my_online_store",
)

client.materialize_features(
    features=[feature],
    online_config=online_config,
    trigger=StreamingMode(),
)

Utiliza una tabla Delta poblada por Zerobus

Una tabla Delta poblada por Zerobus puede servir como fuente de funciones de streaming. Zerobus no configura delta.enableChangeDataFeed=true automáticamente. Debes establecer esta propiedad manualmente en la tabla Delta de destino antes de usarla como fuente de función de streaming.

Condiciones del filtro en fuentes de streaming

Usa filter_condition para filtrar filas antes de agregarlas en StreamSource o DeltaTableSource.

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

Selección de columnas a partir de fuentes de streaming

Los atributos ColumnSelection transmiten el valor más reciente de la fuente para cada clave de entidad sin agregación. Durante el entrenamiento, los valores de las variables mantienen la precisión correspondiente a ese momento temporal.

Las funciones de selección de columnas no tienen TTL. Para eliminar un valor seleccionado de la tienda online, la fuente debe emitir un valor nulo para la columna seleccionada.

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"),
)

Accede a campos anidados desde un StreamSource

Para un StreamSource, puedes acceder a campos JSON anidados usando notación de puntos (por ejemplo, value.nested_field.amount). En tiempo de servicio, la carga de solicitud y la respuesta usan nombres de nodo hoja (por ejemplo, amount en lugar de value.amount). Los nombres de los nodos hoja deben ser únicos en todas las columnas de salida de entidad, series temporales y características dentro de un modelo o una especificación de característica, ya que el punto de conexión de servicio usa los nombres de hoja para enrutar los valores.

Ventanas de tiempo para las características de streaming

Las funciones de streaming solo admiten RollingWindow para las agregaciones. Las ventanas móviles recalculan continuamente sobre los datos más recientes, lo que encaja con la naturaleza en tiempo real de las fuentes de datos en streaming. TumblingWindow y SlidingWindow están diseñados para el cálculo por lotes en intervalos históricos fijos.

Cuaderno de ejemplo de funcionalidades de streaming

Cuaderno de inicio rápido para vistas de atributos de streaming

Obtener el portátil

Entrenamiento e inferencia de modelos

Para entrenar modelos y ejecutar la inferencia por lotes con vistas de características, incluidos log_model(), score_batch()y create_training_set(), consulte Entrenamiento de modelos con vistas de características.

Materialización de características

Después de definir las características, puede materializarlas en almacenes offline u online para una reutilización eficaz en los flujos de trabajo de entrenamiento y provisión. Una vez materializadas las características, puede ejecutar los modelos mediante la ejecución de modelos en la CPU. Para obtener más información, consulte Materializar vistas de atributos.

procedimientos recomendados

Nomenclatura de características

  • Use nombres descriptivos para las características críticas para la empresa.
  • Siga las convenciones de nomenclatura coherentes entre los equipos.
  • Use nombres generados automáticamente a medida que empiece a desarrollar características.

Ventanas de plazo

  • Alinee los límites de ventana con ciclos de negocio (diario, semanal).
  • Las ventanas más cortas capturan tendencias recientes, pero pueden ser ruidosas. Las ventanas más largas producen distribuciones de características más estables, pero podrían perderse cambios de comportamiento recientes. Elija en función de la rapidez con la que cambia la señal subyacente para su caso de uso. Por ejemplo, una ventana de 7 días suaviza las fluctuaciones diarias y genera entradas de modelo coherentes, mientras que una ventana de 1 hora reacciona rápidamente a los cambios de comportamiento, pero podría introducir varianza que degrada el rendimiento del modelo. Si la precisión del modelo se degrada cuando cambia la distribución, use una ventana más larga para estabilizar las entradas.
  • Las ventanas de ráfagas y deslizantes son más escalables que las ventanas continuas. Comience con ventanas deslizantes para la mayoría de los casos de uso.

Performance

  • Materialice las características del mismo origen de datos en una sola materialize_features llamada para minimizar los exámenes de datos.
  • Use la misma granularidad (por ejemplo, todas las duraciones de diapositivas de 1 hora o de 1 día) para las características del mismo origen de datos para permitir una mejor agrupación durante la materialización.

Columnas de entidad frente a condiciones de filtro

Use esta guía de decisión al trabajar con características de la misma tabla de origen:

Use entity (en create_feature) cuando necesite diferentes niveles de agregación:

  • Características de nivel de cliente (una fila por cliente): entity=["customer_id"]
  • Características cliente-comerciante (varias filas por cliente): entity=["customer_id", "merchant_id"]
  • Los distintos niveles de agregación pueden compartir lo mismo DeltaTableSource: especificar valores diferentes entity en cada definición de característica

Use filter_condition (en DeltaTableSource) cuando necesite filtrar filas en el mismo nivel de agregación:

  • Solo transacciones de alto valor: filter_condition="amount > 100" (todavía agregadas por cliente)
  • Solo pedidos completados: filter_condition="status = 'completed'" (aún agregados por cliente)

Regla general: Si el cambio daría lugar a un número diferente de filas por valor de entidad, use valores diferentes entity en las definiciones de características. Si solo estás filtrando qué filas contribuyen a la misma agregación, usa filter_condition en el origen de los datos.

Patrones comunes

Análisis de clientes

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)))),
]

Análisis de tendencias

# 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))),
)

Patrones estacionales

# 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))),
)

Limitaciones

  • Los nombres de las columnas entity y timeseries deben coincidir entre el conjunto de datos de entrenamiento (etiquetado) y las definiciones de características cuando se usan en la create_training_set API.
  • El nombre de columna usado como columna label del conjunto de datos de entrenamiento no debe existir en las tablas de origen usadas para definir Features.
  • En la create_feature API se admite una lista limitada de funciones (UDAFs). Consulte Funciones admitidas.
  • Las columnas de entidad no pueden ser de tipo DATE o TIMESTAMP.
  • RequestSource solo admite tipos de datos escalares definidos en ScalarDataType (INTEGER, FLOAT, BOOLEAN, STRING, DOUBLE, LONG, TIMESTAMP, DATE, SHORT). No se admiten tipos complejos como matrices, asignaciones y estructuras.
  • RequestSource no admite funciones de agregación ni ventanas de tiempo. Solo se pueden usar ColumnSelection funciones.
  • El conjunto de nombres de columnas de entidad, nombres de columnas timeseries y nombres de columna de características de solicitud debe ser único a nivel global en todos los orígenes de un conjunto de entrenamiento o un punto de conexión de servicio.
  • score_batch podría no funcionar en entornos informáticos sin servidor. Para solucionarlo, use un clúster de cómputo clásico que ejecute Databricks Runtime 17.0 ML o posterior.

Para conocer las limitaciones específicas de la materialización, consulte Limitaciones.