Referencia de la API de vistas de características

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.

Control de acceso

Las características son objetos de catálogo de Unity que se pueden controlar. El acceso a una característica se controla mediante los CREATE FEATUREprivilegios de catálogo de , READ FEATUREy MANAGE Unity. Para obtener descripciones completas, consulte Referencia de privilegios del catálogo de Unity.

  • CREATE FEATURE: Necesario para crear una característica en un esquema. create_feature y register_feature requieren CREATE FEATURE en el esquema primario. Siguiendo el principio de privilegios mínimos, conceda CREATE FEATURE en el nivel de esquema; también puede concederlo en un catálogo para permitir la creación de características en cualquier esquema de ese catálogo.
  • READ FEATURE: Necesario para leer metadatos de características. get_feature, create_training_set, y list_materialized_features requieren READ FEATURE en la característica. Este privilegio no concede acceso a datos de características en tablas de salida fuente o materializadas. Para leer esos datos para entrenamiento o servicio, también debes estar SELECT en las tablas correspondientes. READ FEATURE concedido en un esquema o catálogo se aplica a todas las características actuales y futuras que contiene.
  • MANAGE: Necesario para gestionar el ciclo de vida y las subvenciones de una funcionalidad. Eliminar una característica con delete_feature, y materializar una característica con materialize_features, requiere MANAGE sobre la característica. Eliminar una característica materializada con delete_materialized_feature no se rige por MANAGE: solo el creador de la característica materializada puede eliminarla.

Todas las operaciones de características también requieren USE CATALOG en el catálogo primario y USE SCHEMA en el esquema primario. Para obtener información sobre cómo MANAGE y READ FEATURE aplicar a la materialización, consulte Permisos.

API de vista de características

Feature constructor y register_feature()

El enfoque recomendado es construir un Feature objeto localmente y usarlo register_feature para conservarlo en el catálogo de Unity. Este flujo de trabajo de dos pasos permite experimentar con características (incluidas create_training_set) antes de registrarlas.

Feature(
    source: DataSource,                                    # Required: DeltaTableSource, StreamSource, RequestSource, or FeatureViewSource
    function: Union[AggregationFunction, ColumnSelection, CustomUDF], # Required for all features
    entity: Optional[List[str]] = None,                    # Required for DeltaTableSource and StreamSource
    timeseries_column: Optional[str] = None,               # Required for DeltaTableSource and StreamSource
    name: Optional[str] = None,                            # Optional: Feature name (auto-generated if omitted)
    description: Optional[str] = None,                     # Optional: Feature description
)

FeatureEngineeringClient.register_feature() registra un objeto construido Feature localmente en el catálogo de 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() valida, construye e registra inmediatamente una característica en el catálogo de Unity en un solo paso. Úselo cuando no necesite experimentar con la característica localmente primero.

FeatureEngineeringClient.create_feature(
    source: DataSource,                                    # Required: DeltaTableSource, StreamSource, RequestSource, or FeatureViewSource
    function: Union[AggregationFunction, ColumnSelection, CustomUDF], # Required for all features
    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 DeltaTableSource and StreamSource
    timeseries_column: Optional[str] = None,               # Required for DeltaTableSource and StreamSource
    name: Optional[str] = None,                            # Optional: Feature name (auto-generated if omitted)
    description: Optional[str] = None,                     # Optional: Feature description
) -> Feature

Parámetros:

  • source: La fuente de datos utilizada en el cálculo de características (DeltaTableSource, StreamSource, RequestSource, o FeatureViewSource).
  • function: Y AggregationFunction que agrupa un operador y una ventana temporal, ColumnSelection("column_name") para características de paso o CustomUDF para transformaciones hilera. Consulte Funciones soportadas para tipos de fuente compatibles.
  • catalog_name: el nombre del catálogo del catálogo de Unity para la característica.
  • schema_name: el nombre del esquema del catálogo de Unity para la característica.
  • entity: lista de nombres de columna que definen las claves de agregación o búsqueda (claves principales). Requerida para DeltaTableSource y StreamSource. Por ejemplo, ["user_id"] agrega o busca por usuario. Omite y RequestSourceFeatureViewSource.
  • timeseries_column: columna de marca de tiempo usada para la agregación de período de tiempo o selección de valor más reciente. Requerida para DeltaTableSource y StreamSource. Omite y RequestSourceFeatureViewSource.
  • name: nombre de característica opcional. Si se omite, se genera automáticamente a partir de la columna de entrada, la función y la ventana (por ejemplo, amount_avg_rolling_7d).
  • description: descripción opcional de la característica.

Devuelve: Una instancia de función validada

Genera: ValueError si se produce un error en la validación

delete_feature()

Elimina una característica del catálogo de Unity por su nombre completo.

FeatureEngineeringClient.delete_feature(
    full_name: str,  # Required: '<catalog>.<schema>.<feature_name>'
) -> None
fe.delete_feature(full_name="main.store.amount_sum_rolling_7d")

Antes de eliminar una característica, quite o actualice los modelos o especificaciones de características que hacen referencia a ella. Una función no puede eliminarse mientras aún tenga características materializadas. Elimina primero las características materializadas y luego borra la funcionalidad. Consulte Eliminación de una característica materializada.

Nombres generados automáticamente

Cuando name se omite, se genera automáticamente un nombre. Los nombres generados siguen el patrón: {column}_{function}_{window}. Por ejemplo:

  • price_avg_rolling_1h (Precio medio de 1 hora)
  • transaction_count_rolling_30d_1d (Recuento de 30 días de la transacción con un retraso de 1d a partir de la marca de tiempo del evento)

Funciones compatibles

Funciones de agregación

Note

Las funciones de agregación se encapsulan en un AggregationFunction junto con un período de tiempo, como se describe en las ventanas de tiempo. Cada función toma un input parámetro que especifica la columna de origen que se va a agregar.

Function Description Ejemplo de caso de uso
Sum(input="column") Total de valores Uso diario de aplicaciones por usuario en minutos
Avg(input="column") Promedio de valores Cantidad media de transacción
Count(input="column") Número de registros Número de inicios de sesión por usuario
Min(input="column") Valor mínimo Frecuencia cardíaca más baja registrada por un dispositivo portátil
Max(input="column") Valor máximo Cantidad de transacción más alta por sesión
StddevPop(input="column") Desviación estándar de la población Variabilidad diaria de la cantidad de transacciones en todos los clientes
StddevSamp(input="column") Desviación estándar de ejemplo Variabilidad de las tasas de clics en campañas publicitarias
VarPop(input="column") Varianza de población Propagación de lecturas de sensores para dispositivos IoT en una fábrica
VarSamp(input="column") Varianza muestral Propagación de clasificaciones de películas en un grupo muestreado
ApproxCountDistinct(input="column", relativeSD=0.05) Recuento único aproximado Recuento distinto de artículos comprados
ApproxPercentile(input="column", percentile=0.95, accuracy=100) Percentil aproximado Latencia de respuesta p95
First(input="column") Primer valor Primera marca de tiempo de inicio de sesión
Last(input="column") Último valor Importe de compra más reciente
FirstN(input="column", n=3) Primeros n valores como un array Los tres primeros productos vistos en una sesión
LastN(input="column", n=3) Últimos n valores como un array Tres estados más recientes de casos de apoyo
FirstDistinct(input="column", n=3) Primero n , valores distintos como array Primeras tres categorías de productos distintas vistas
LastDistinct(input="column", n=3) Últimos n valores distintos como un array Las tres categorías de comerciantes más recientes y distintas

Note

First, Last, FirstN, LastN, FirstDistinct, , y LastDistinct incluyen valores nulos por defecto. Para omitir valores NULL, agregue un filter_condition que excluya explícitamente las columnas de entrada que son NULL.

FirstN, LastN, , y LastDistinct usar las timeseries_column características para ordenar las filas de entrada y devolver un array que contenga hasta n valoresFirstDistinct. El n parámetro debe ser un entero positivo. FirstN y FirstDistinct selecciona valores desde el más antiguo hasta el más reciente. LastN y LastDistinct selecciona los valores de más reciente a más antiguo, y luego devuelve los valores seleccionados en orden de marca de tiempo. FirstDistinct y LastDistinct eliminar valores duplicados al seleccionar valores en esa dirección.

Por ejemplo, si las filas fuente de una entidad están ordenadas por event_time , ["A", "A", "B", "C", "B", "B"]las siguientes funciones devuelven:

Function Resultado
FirstN(input="event_type", n=3) ["A", "A", "B"]
LastN(input="event_type", n=3) ["C", "B", "B"]
FirstDistinct(input="event_type", n=3) ["A", "B", "C"]
LastDistinct(input="event_type", n=3) ["A", "C", "B"]

FirstN, LastN, FirstDistinct, y LastDistinct requieren databricks-feature-engineering la versión 0.17.0 o posterior.

CustomUDF

CustomUDFaplica una función Python (UDF) registrada en el Catálogo de Unity a cada fila. Úsalo para transformar entradas de peticiones o combinar valores de características. No agrega filas ni define una ventana temporal.

CustomUDF(
    function_name="main.ecommerce.log_amount_udf",
    input_bindings={"amount": "transaction_amount"},
)

input_bindings asigna el nombre de cada parámetro UDF a una entrada. Para RequestSource, la entrada es un nombre de columna fuente. Para FeatureViewSource, es una referencia de características aguas arriba. Vincula todos los parámetros de la UDF, incluidos los parámetros con valores predeterminados. Los tipos de entrada deben coincidir exactamente con los tipos de parámetros UDF, sin lanzamientos numéricos implícitos. Utiliza tipos de entrada y retorno escalar.

Source Behavior
RequestSource Transforma columnas desde el DataFrame de entrenamiento o la solicitud de inferencia.
FeatureViewSource Combina los valores de las características aguas arriba. Consulta FeatureViewSource.

Las funciones respaldadas CustomUDF por Delta no pueden materializarse ni servirse en línea. Para transformar los valores de características respaldadas por tabla para entrenamiento y servicio, define una agregación respaldada por Delta o una característica de selección de columnas y referenciala a través de FeatureViewSource.

CustomUDF no está soportado con StreamSource. Para transformar la salida de una función de streaming, haz referencia a esa característica a través FeatureViewSourcede .

CustomUDF con RequestSource requiere databricks-feature-engineering la versión 0.17.0 o posterior.

Para usar un CustomUDF, necesitas el EXECUTE privilegio en la UDF, el USE CATALOG privilegio en su catálogo padre y el USE SCHEMA privilegio en su esquema padre.

El siguiente ejemplo utiliza NumPy para calcular log(1 + amount), reduciendo la escala de grandes cantidades de transacciones. Ejecutarlo en computación serverless con dependencias personalizadas de UDF activadas. El main.ecommerce esquema debe existir.

spark.sql("""
CREATE OR REPLACE FUNCTION main.ecommerce.log_amount_udf(amount DOUBLE)
RETURNS DOUBLE
LANGUAGE PYTHON
ENVIRONMENT (
  dependencies = '["numpy==1.26.4"]',
  environment_version = '5'
)
AS $$
import numpy as np

if amount is None or not np.isfinite(amount) or amount < 0:
    return None
return float(np.log1p(amount))
$$
""")

Registrar una característica que vincule la columna transaction_amount de solicitud al parámetro amountUDF:

from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
    CustomUDF, FieldDefinition, RequestSource, ScalarDataType,
)

fe = FeatureEngineeringClient()

log_transaction_amount = fe.create_feature(
    source=RequestSource(
        schema=[
            FieldDefinition(name="transaction_amount", data_type=ScalarDataType.DOUBLE),
        ]
    ),
    function=CustomUDF(
        function_name="main.ecommerce.log_amount_udf",
        input_bindings={"amount": "transaction_amount"},
    ),
    catalog_name="main",
    schema_name="ecommerce",
    name="log_transaction_amount",
)

El UDF ENVIRONMENT configura dependencias para el cálculo offline. Para el servicio online, también declara los paquetes en create_feature_spec(extra_pip_requirements=...) o log_model(extra_pip_requirements=...). No se copian automáticamente de la UDF. Consulta Dependencias de servicio de características y dependencias de modelo.

CustomUDF las características no se pueden materializar. Los UDF respaldados por solicitudes y funciones funcionan bajo demanda durante la formación y el servicio. Cada UDF en una cadena de dependencias añade cálculo, así que mantén las funciones y cadenas pequeñas. Los UDF deben gestionar las entradas ausentes, que pueden estar None fuera de línea o NaN fuera de línea.

Para orientación sobre cómo manejar los valores faltantes, véase Cómo manejar los valores de características ausentes.

ColumnSelection (paso)

ColumnSelection selecciona una sola columna de un origen sin aplicar ninguna agregación. Se ajusta directamente en el function parámetro (no dentro AggregationFunctionde ). El tipo de valor devuelto se deduce del esquema de origen.

Function Description Ejemplo de caso de uso
ColumnSelection("col") Valor más reciente de una columna (sin agregación) Categoría de proveedor más reciente, paso a través de un campo de solicitud

ColumnSelection Soporta las siguientes fuentes de datos:

  • DeltaTableSource: devuelve el valor más reciente por clave de entidad a través de una combinación a un momento dado (sin agregación de ventana de búsqueda).
  • StreamSource: Devuelve el último valor por clave de entidad del Stream (sin agregación de ventana de retorno).
  • RequestSource: pasa por el valor proporcionado en tiempo de inferencia (o extraído del dataframe etiquetado en tiempo de entrenamiento).
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",
)

Ejemplo: características de selección de agregaciones y columnas

En el ejemplo siguiente se muestran las características definidas en el mismo origen de datos.

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

Características con condiciones de filtro

El filter_condition parámetro permite filtrar filas de la tabla de origen antes de calcular agregaciones. Esto funciona como una cláusula SQL WHERE que se aplica antes de agrupar y agregar datos.

Note

filter_condition filtra las filas antes de la agregación, como una cláusula SQL WHERE aplicada antes de GROUP BY. No cambia la granularidad, que siempre se define en entity la definición de características.

Los filtros son útiles al trabajar con tablas de origen de gran tamaño que incluyen un superconjunto de datos necesarios para el cálculo de características y minimizar la necesidad de crear vistas independientes sobre estas tablas.

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

Orígenes de datos

DeltaTableSource

DeltaTableSource es un objeto Python efímero que se usa para definir cómo se calculan las características desde una tabla de origen. No crea una nueva tabla. Especifica la configuración para leer datos y agregar características.

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
    lateness: Optional[SourceLateness] = None,       # Optional: Expected source settling time
)

Parámetros:

  • catalog_name, , schema_nametable_name: identifique la tabla Delta de origen en el catálogo de Unity.
  • filter_condition: una cláusula SQL WHERE aplicada antes de la agregación. Ejemplo: "status = 'completed'".
  • transformation_sql: expresión SQL SELECT aplicada a la tabla de origen. Úselo para cambiar el nombre de las columnas, los tipos de conversión o las columnas derivadas de proceso antes de la agregación. Si se omite, se seleccionan todas las columnas (*). Ejemplo: "user_id, CAST(amount AS DOUBLE) AS amount, event_time".
  • dataframe_schema: el esquema del dataframe resultante después de las transformaciones, en formato JSON StructType de Spark (de df.schema.json()). Obligatorio si transformation_sql se proporciona. Esto indica al sistema los nombres y tipos de columna que resultan de la transformación.
  • lateness: Un SourceLateness objeto que describe cuánto tarda normalmente la fuente en completarse en tiempo de evento. Si se omite, la fuente se considera completa de inmediato.

Cuando se establecen y filter_conditiontransformation_sql , la consulta resultante es: SELECT {transformation_sql} FROM {table} WHERE {filter_condition}.

SourceLateness.settling_delay es la forma recomendada de simular durante el entrenamiento un retraso ETL consistente que afecta a la materialización en línea. Azure Databricks retroactiva el tiempo de evaluación de entrenamiento elegible en este periodo de duración para que un ejemplo de entrenamiento no utilice datos que aún estarían en tránsito online. Durante la materialización, Azure Databricks espera el mismo tiempo antes de publicar una ventana completada y sirve a la última ventana completada durante el periodo intermedio.

Por ejemplo, supongamos que un trabajo ETL diario se completa 8 horas después de la medianoche en una zona horaria local donde la medianoche corresponde a las 07:00 UTC. Utiliza un retraso de asentamiento de 8 horas y un desplazamiento de ventana de 7 horas:

from datetime import timedelta
from databricks.feature_engineering.entities import (
    DeltaTableSource,
    SourceLateness,
    TumblingWindow,
)

source = DeltaTableSource(
    catalog_name="main",
    schema_name="analytics",
    table_name="events",
    lateness=SourceLateness(settling_delay=timedelta(hours=8)),
)

window = TumblingWindow(
    window_duration=timedelta(days=1),
    offset=timedelta(hours=7),
)

Note

El timeseries_column debe ser de tipo TimestampType o TimestampNTZType. DateType no es compatible con series temporales; cast la columna a TimestampType first (por ejemplo, con transformation_sql).

Ejemplo: Uso transformation_sql de para transformaciones de columna

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

Ejemplo: Derivar transformation_sql y dataframe_schema de un dataFrame de PySpark

Puede escribir la transformación como una consulta pySpark y, a continuación, extraer el esquema de la trama de datos resultante:

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

Expresiones soportadas transformation_sql

Las mismas reglas se aplican a transformation_sql sobre DeltaTableSource y StreamSource.

transformation_sql soporta cualquier expresión por fila; las operaciones se evalúan de forma independiente para cada fila. No cambian el número de filas ni la correspondencia uno a uno con la fuente. Las expresiones por filas incluyen renombramientos de columnas, casts, operaciones aritméticas y más.

No se soportan operaciones que cambian la forma o el recuento de filas, como agregaciones como SUM() o COUNT(). Use AggregationFunction en la definición de características en su lugar.

DeltaTableSource.from_sql()

Como comodidad, puede crear a DeltaTableSource partir de una consulta SQL. El método analiza la consulta para extraer automáticamente el nombre de la tabla, transformation_sqly filter_condition.

DeltaTableSource.from_sql(
    sql: str,                           # Required: SQL SELECT query
    spark: SparkSession,                # Required: active SparkSession (for schema inference)
) -> DeltaTableSource

Solo se admiten consultas simples SELECT ... FROM ... [WHERE ...] . Se rechaza SQL complejo (JOINs, subconsultas, CTE, UNION). Para consultas complejas, construya DeltaTableSource directamente con transformation_sql y 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",
)

Iteración con to_dataframe()

Use source.to_dataframe() para obtener una vista previa de los datos que se usarán para el cálculo de características. Esto es útil para iterar en filter_condition y transformation_sql hasta que generan los resultados esperados.

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()

Descripción de las entidades

Las columnas de entidad definen el nivel de agregación de las características. Se especifican en la Feature definición, no en DeltaTableSource. Las entidades determinan:

  • Cómo se agrupan los datos: las características se agregan por combinación única de valores de entidad (similares a GROUP BY en SQL)
  • La estructura de clave principal: cada combinación de entidad única da como resultado una fila de características calculadas.

Ejemplo: Características de nivel de cliente

El código siguiente agrega características en el nivel de cliente (una fila por cliente):

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

Ejemplo: Características de nivel de almacén de clientes

Para agregar características en un nivel más detallado (una fila por combinación de almacén de clientes), use varias columnas de entidad:

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

Cuando necesite características en distintos niveles de agregación (por ejemplo, de nivel de cliente y de almacén de clientes), use valores diferentes entity en las definiciones de características. Lo mismo DeltaTableSource se puede compartir entre características con distintas configuraciones de entidad.

StreamSource

StreamSource hace referencia a un objeto Stream. Stream contiene la configuración de conexión, autenticación, esquema e ingesta para el origen de streaming. Para Kafka, las referencias de columna en las definiciones de características deben tener value. el prefijo o key. indicar qué parte del mensaje se va a leer.

StreamSource(
    full_name: str,                       # Required: Three-part Stream name (catalog.schema.stream)
    filter_condition: Optional[str] = None,      # Optional: SQL WHERE clause applied before aggregation
    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
    lateness: Optional[SourceLateness] = None,    # Optional: Expected source settling time
)

Parámetros:

  • full_name: nombre completo de tres partes de un objeto Stream (por ejemplo, "my_catalog.my_schema.my_stream").
  • filter_condition (opcional): una cláusula SQL WHERE aplicada a los datos del flujo antes de la agregación, mediante referencias de columna con prefijo de punto (por ejemplo, "value.event_type = 'purchase'").
  • transformation_sql(opcional): Una expresión SQL SELECT aplicada antes de la agregación o selección de columnas, usando referencias con prefijo de puntos a las key estructuras de y.value Soporta las mismas expresiones por filas que DeltaTableSource. Si se omite, la fuente utiliza todas las columnas (*).
  • dataframe_schema: El esquema JSON de Spark StructType de la salida proyectada. Obligatorio si configuras transformation_sql.
  • lateness: Un SourceLateness objeto que describe cuánto tiempo tarda normalmente el flujo en completarse en tiempo de evento. Consulte SourceLateness.settling_delay.
from databricks.feature_engineering.entities import StreamSource

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

Se deriva dataframe_schema ejecutando la proyección contra la tabla de ingestión del Stream, que expone las key estructuras y value .

transformation_sql = (
    "value.amount * value.conversion_rate AS converted_amount, "
    "struct(value.user_id AS user_id, value.event_time AS time) AS event"
)

ingestion_table = "my_catalog.my_schema.events_ingestion"
dataframe_schema = spark.sql(
    f"SELECT {transformation_sql} FROM {ingestion_table} LIMIT 0"
).schema.json()

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

RequestSource

RequestSource define un esquema para los datos que se proporcionan en el momento de inferencia en la carga de la solicitud en lugar de buscar desde una tabla materializada previamente. Durante el entrenamiento, estas columnas se extraen de la etiqueta DataFrame pasada a create_training_set. Durante la servicio de modelos, el autor de la llamada debe incluirlos en la carga de la solicitud HTTP.

RequestSource puede usarse con funciones de Vista de Características CustomUDF o ColumnSelection . No admite funciones de agregación ni ventanas de tiempo.

Definición del esquema

Defina el esquema como una lista de FieldDefinition objetos, cada uno de los cuales especifica un nombre de columna y un 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),
    ]
)

Supported data types (Tipos de datos admitidos)

RequestSourceadmite los tipos escalares definidos en ScalarDataType: INTEGER, FLOATBOOLEANSTRINGDOUBLELONG, TIMESTAMP, , DATE, . SHORT No se admiten tipos complejos como matrices, mapas y estructuras.

Cómo se hidratan los datos de solicitud

Context Behavior
Entrenamiento (create_training_set) Las columnas se extraen de la etiqueta DataFrame. Los tipos se validan con el esquema declarado. Los errores de coincidencia generan un error (sin conversión implícita).
Servicio (punto de conexión del modelo) Las columnas se extraen de dataframe_records o dataframe_split en la solicitud HTTP. Los valores JSON se convierten en los tipos declarados (por ejemplo, el número JSON → DOUBLE).

Firma de modelo

Cuando se registra un modelo mediante log_model con un conjunto de entrenamiento que incluye RequestSource características, las RequestSource columnas se agregan a la firma del modelo de MLflow según las entradas necesarias. Esto significa que el esquema de API del punto de conexión de servicio refleja qué campos deben proporcionar los autores de llamada en el momento de la inferencia.

FeatureViewSource

FeatureViewSource utiliza las salidas de otras Vistas de Características como entradas para un CustomUDFarchivo . El encadenamiento de características crea un grafo dirigido acíclico (DAG). Por ejemplo, una característica de margen puede combinar agregados de ingresos y costes, y otra característica puede transformar el margen.

Usa databricks-feature-engineering la versión 0.18.0 o posterior para FeatureViewSource.

Pasa una lista de Feature objetos a features, no cadenas de nombres de características. Recuperar características registradas con get_feature. En input_bindings, utiliza la full_namecaracterística de cada registro . Para una función local no registrada, usa su name en su lugar.

El siguiente ejemplo asume dos características registradas, revenue_sum_7d y cost_sum_7d, que devuelven DOUBLE valores por customer_id y usan event_time para el cálculo en un momento en el tiempo:

from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import CustomUDF, FeatureViewSource

fe = FeatureEngineeringClient()
revenue = fe.get_feature(full_name="main.ecommerce.revenue_sum_7d")
cost = fe.get_feature(full_name="main.ecommerce.cost_sum_7d")

spark.sql("""
CREATE OR REPLACE FUNCTION main.ecommerce.margin_udf(revenue DOUBLE, cost DOUBLE)
RETURNS DOUBLE
LANGUAGE PYTHON
AS $$
import math

if revenue is None or cost is None:
    return None
if not math.isfinite(revenue) or not math.isfinite(cost) or revenue <= 0:
    return None
return (revenue - cost) / revenue
$$
""")

margin = fe.create_feature(
    source=FeatureViewSource(features=[revenue, cost]),
    function=CustomUDF(
        function_name="main.ecommerce.margin_udf",
        input_bindings={"revenue": revenue.full_name, "cost": cost.full_name},
    ),
    catalog_name="main",
    schema_name="ecommerce",
    name="margin",
)

Se aplican las restricciones que se indican a continuación:

  • Solo CustomUDF se soporta como función. Las características aguas arriba pueden ser agregaciones, selecciones de columnas u otras CustomUDF características.
  • Omite entity y timeseries_column en la característica derivada. Cada característica upstream conserva su propia entidad, marca de tiempo y definición de ventana.
  • Una función tiene una sola fuente. Para combinar un valor de solicitud con una característica respaldada por tabla, define una RequestSource característica y referencia ambas mediante FeatureViewSource.
  • Cada característica ascendente declarada debe usarse en input_bindings. No se permiten ciclos.
  • Registra las características aguas arriba antes de registrar la característica derivada. Se pueden usar grafos locales no registrados con create_training_set para experimentación.
  • Para el entrenamiento o servicio, necesitas el READ FEATURE privilegio o MANAGE sobre la característica derivada y sus características ascendentes transitivas. Utiliza nombres de características distintas a lo largo del gráfico para registro y servicio, incluso entre catálogos o esquemas.
  • Una característica puede hacer referencia hasta 20 funciones directas de arriba. Los grafos registrados soportan una profundidad máxima de cinco características a lo largo de un camino de dependencias, incluyendo la característica base.
  • FeatureViewSource las características no pueden materializarse ni evaluarse con compute_features. Úsalo create_training_set para evaluarlos fuera de línea. Para el servicio online, materializa las funciones de upstream respaldadas con tabla soportadas en su lugar.

Para la evaluación de dependencias y la selección de salida, véase Entrenar con características FeatureViewSource. Para el despliegue, véase características derivadas de Serve.

API de entrenamiento e inferencia

create_training_set y score_batch calculan los valores de características correctos a un momento dado a petición de los datos de origen. En el caso de las características que admiten la materialización sin conexión, como las agregaciones de ventana deslizantes en orígenes de tabla delta, la materialización de las características en primer lugar en un almacén sin conexión mejora el rendimiento de ambas operaciones. Cuando están disponibles las características sin conexión materializadas, las operaciones leen los datos sin conexión predefinidos en lugar de volver a calcular los valores de características del origen. Consulta Materializar vistas de características para materializar las características en un almacén sin conexión.

create_training_set()

Crea un conjunto de datos de entrenamiento con el cálculo de características correcto a un momento dado. Para obtener más información, consulte Entrenamiento de modelos con vistas de características.

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()

Registra un modelo con metadatos de características para el seguimiento de linaje y la búsqueda automática de características durante la inferencia. Para obtener más información, consulte Entrenamiento de modelos con vistas de características.

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()

Realiza la inferencia por lotes sin conexión con la búsqueda automática de características. Usa los metadatos de características almacenados con el modelo para calcular características correctas a un momento dado, lo que garantiza la coherencia con el entrenamiento.

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

La trama de datos de entrada debe contener las columnas de entidad y timeseries usadas durante el entrenamiento. Las características se calculan automáticamente a partir de los datos de origen.

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()

Ventanas de plazo

Las vistas de características soportan cuatro tipos de ventanas para controlar el comportamiento de retroalimentación en agregaciones basadas en ventanas temporales. Los tipos de ventanas disponibles dependen de la fuente de la característica: las características de fuente en streaming pueden usar ventanas enrollables y de dientes de sierra, y las de fuente por lotes pueden usar ventanas enrollables, de tumbling y deslizantes.

  • Las ventanas que se revierten desde la hora del evento. La duración y el retraso se definen explícitamente.
  • Las ventanas de deslizamiento son ventanas de tiempo fijas y no superpuestas. Cada punto de datos pertenece exactamente a una ventana.
  • Las ventanas deslizantes son ventanas de tiempo superpuestas y en movimiento, con un intervalo de desplazamiento configurable.
  • Las ventanas de diente de sierra mantienen fresca una ventana de retroceso largo sobre una fuente de streaming usando un recorrido híbrido por lotes y un recorrido de streaming. Ver ventana de dientes de sierra.

La siguiente ilustración muestra los tipos de ventanas de rodadura, deslizamiento, rodante y de diente de sierra.

Ventanas retrocedientes con rodas, deslizamientos, rodamientos y miradas en forma de dientes de sierra.

Temporización de la ventana temporal

Utilízalo delay para evaluar una ventana en un momento analítico anterior. Por ejemplo, una ventana de 30 días con un retraso de 7 días calcula un valor de 30 días a partir de una semana antes de la hora de la evaluación. delay es independiente del momento de llegada de la fuente. Para modelar el tiempo que tarda en llegar ese dato fuente, configura SourceLateness.settling_delay en su lugar.

Cuando ambos escenarios están presentes, componen. Azure Databricks trata la ventana como completa tras el retraso de asentamiento de la fuente y la evalúa usando el retardo analítico.

Úsalo offset para cambiar la alineación de los límites fijos de las ventanas. Por defecto, las ventanas deslizantes y las ventanas deslizantes están alineadas a medianoche UTC. Por ejemplo, un desplazamiento de 22 horas alinea un límite diario a las 22:00 UTC. Para aproximar los límites en una zona horaria local, configura un desplazamiento estático respecto a UTC. El desplazamiento no ajusta el horario de verano, no desplaza la hora de evaluación ni modela los datos que llegan tarde.

La siguiente tabla resume el soporte para estos campos:

Campo Ventanas compatibles Restricción
delay Rodar, girar y deslizarse Debe ser un no negativo datetime.timedelta
offset Tumbling y deslizamiento Debe ser no negativo y más corto que el periodo*
SourceLateness.settling_delay Características de rodamiento, volteo y deslizamiento Debe ser un no negativo datetime.timedelta
start_time Rodar, girar y deslizarse Debe ser un datetime.datetime

*Punto: Para una ventana de caída, el desplazamiento debe ser menor que window_duration. Para una ventana deslizante, debe ser más corta que slide_duration.

Hora de inicio

Úsase start_time para establecer el límite de tiempo de evento más temprano en UTC en el que una característica puede emitir una salida. El límite es inclusivo. start_time Salidas de gates. No restringe las filas de fuente históricas que puede leer una ventana, ni cambia la alineación de la ventana. Si start_time cae entre dos fronteras alineadas, la primera salida de ventana fija elegible es la siguiente frontera.

Con start_time, las ventanas de duración fija pueden emitirse antes de que haya transcurrido una ventana completa en la fuente. Estas primeras salidas utilizan cualquier historia de fuentes disponible. Por ejemplo, consideremos una ventana deslizante con un año window_duration y otro de un día slide_duration, sobre una fuente cuyos datos comienzan el 1 de enero de 2024:

  • Sin start_time, la función se emite por primera vez el 1 de enero de 2025, una vez que se pueda formar una ventana completa de un año.
  • Con start_time fecha prevista para el 21 de agosto de 2024, la función se emite por primera vez el 21 de agosto de 2024. Esa publicación cubre únicamente la historia de fuentes disponible hasta ahora, desde el 1 de enero de 2024. La ventana alcanza su periodo completo de un año el 1 de enero de 2025 y produce resultados completos a partir de entonces.

Dado que start_time no cambia la alineación de ventanas, un valor entre dos límites alineados no crea un nuevo límite. Para una ventana de caída con límites diarios a medianoche UTC, a start_time de las 06:00 UTC emite primero en el siguiente límite de medianoche. Un start_time que aterriza exactamente en un límite emite en ese límite, porque el límite es inclusivo.

Si start_time no está ajustado, las ventanas deslizantes y las ventanas deslizantes de duración fija emiten primero en un límite alineado después de que se pueda formar una ventana completa. Las ventanas deslizantes y las ventanas móviles de por vida emiten tan pronto como existen datos de origen elegibles.

Note

start_time está soportado para funciones por lotes que se usan DeltaTableSource con ventanas enrollables, de tumbling o correderas. No está soportado con StreamSource ni SawtoothWindowcon .

Por ejemplo:

from datetime import datetime, timedelta
from databricks.feature_engineering.entities import SlidingWindow

window = SlidingWindow(
    window_duration=timedelta(days=365),
    slide_duration=timedelta(days=1),
    start_time=datetime(2024, 8, 21),
)

Ventana enrollable

Note

RollingWindow anteriormente se denominaba ContinuousWindow. Si va a migrar desde una versión anterior del SDK, actualice las importaciones en consecuencia.

Las ventanas graduales se up-toagregados en tiempo real y de fecha, que normalmente se usan a través de datos de streaming. En las canalizaciones de streaming, la ventana gradual emite una nueva fila solo cuando cambia el contenido de la ventana de longitud fija, como cuando un evento entra o sale. Cuando se usa una característica de ventana gradual en canalizaciones de entrenamiento, se realiza un cálculo preciso de características a un momento dado en los datos de origen mediante la duración de la ventana de longitud fija inmediatamente anterior a la marca de tiempo de un evento específico. Esto ayuda a evitar la distorsión en línea y fuera de línea o la filtración de datos. Las funciones en el tiempo T agregan eventos de [T − duración, T).

class RollingWindow(TimeWindow):
    window_duration: datetime.timedelta
    delay: Optional[datetime.timedelta] = None
    start_time: Optional[datetime.datetime] = None

En la tabla siguiente se enumeran los parámetros de una ventana gradual. Las horas de inicio y finalización de la ventana se basan en estos parámetros de la siguiente manera:

  • Hora de inicio: evaluation_time - window_duration - delay (incluido)
  • Hora de finalización: evaluation_time - delay (exclusivo)
Parámetro Constraints
delay (opcional) Debe ser ≥ 0. Desplace la ventana analítica hacia atrás desde la marca temporal de la evaluación. Úsalo SourceLateness.settling_delay para modelar una línea base consistente para el retraso de llegada de la fuente en tu stream.
window_duration Debe ser > 0
start_time (opcional) Límite de tiempo de evento más temprano en el que la característica puede emitir una salida.
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))

Defina una ventana gradual con retraso mediante el código siguiente.

# Compute a 7-day value as of one day before the evaluation time
window = RollingWindow(
    window_duration=timedelta(days=7),
    delay=timedelta(days=1)
)

Ejemplos de ventanas graduales

  • window_duration=timedelta(days=7): crea una ventana retrospectiva de 7 días que termina en el tiempo de evaluación actual. Para un evento a las 2:00 p. m. del día 7, esto incluye todos los eventos de las 2:00 p. m. del día 0 hasta (pero sin incluir) 2:00 p. m. el día 7.

  • window_duration=timedelta(hours=1), delay=timedelta(minutes=30): esto crea una ventana de búsqueda retrospectiva de 1 hora que finaliza 30 minutos antes del tiempo de evaluación. Para un evento a las 3:00 p. m., esto incluye todos los eventos de 1:30 p. m. hasta (pero sin incluir) las 2:30 p. m.

Úsase Last para limitar la frescura de un valor reciente

Combina Last con RollingWindow cuando un valor más reciente es válido solo por un tiempo limitado. En un momento de evaluación, la característica devuelve el valor de la fila con la última marca de tiempo en este intervalo:

[evaluation_time - delay - window_duration, evaluation_time - delay)

Si la última fila del intervalo contiene un valor nulo, la característica devuelve nulo. Si quieres excluir valores nulos de entrada, pon a filter_condition en la fuente.

Esta combinación difiere de ColumnSelection. ColumnSelection devuelve el último valor no nulo observado sin caducarlo según la edad.

Para las características por lotes, esta combinación tiene un modo especial de materialización solo online. Solo soporta DeltaTableSource, Last, RollingWindow, y TableTrigger. Consulta Materialize valores más recientes con límites de frescura.

Ventana deslizante

En el caso de las características definidas mediante ventanas deslizantes de tamaño constante, las agregaciones se calculan en una ventana de longitud fija predefinida que avanza por un intervalo de deslizamiento, lo que produce ventanas no superpuestas que particionan completamente el tiempo. Como resultado, cada evento del origen contribuye exactamente a una ventana. Las características en el momento t agregan datos de ventanas que finalizan en o antes de t (exclusivo). Windows comienza en la época de Unix.

class TumblingWindow(TimeWindow):
    window_duration: datetime.timedelta
    delay: Optional[datetime.timedelta] = None
    offset: Optional[datetime.timedelta] = None
    start_time: Optional[datetime.datetime] = None

En la tabla siguiente se enumeran los parámetros de una ventana deslizante.

Parámetro Constraints
window_duration Debe ser > 0
delay (opcional) Debe ser ≥ 0. Desplace la ventana analítica hacia atrás desde la marca temporal de la evaluación.
offset (opcional) Debe ser ≥ 0 y más corto que window_duration. Desplaza los límites de la ventana desde la medianoche UTC.
start_time (opcional) Límite de tiempo de evento más temprano en el que la característica puede emitir una salida.
from databricks.feature_engineering.entities import TumblingWindow
from datetime import timedelta

window = TumblingWindow(
    window_duration=timedelta(days=1),
    delay=timedelta(hours=2),
    offset=timedelta(hours=22),
)

Ejemplo de ventana deslizante fija

  • window_duration=timedelta(days=5): crea ventanas de longitud fija determinadas previamente de 5 días cada una. Ejemplo: La ventana n.º 1 abarca el día 0 al día 4, la ventana #2 abarca el día 5 al día 9, la ventana #3 abarca el día 10 al día 14, etc. En concreto, Window #1 incluye todos los eventos con marcas de tiempo a partir del 00:00:00.00 día 0 hasta (pero sin incluir) ningún evento con marca 00:00:00.00 de tiempo el día 5. Cada evento pertenece exactamente a una ventana.

Ventana deslizante

Para características definidas mediante ventanas deslizantes, las agregaciones se calculan sobre una ventana que avanza en un intervalo de deslizamiento. Una ventana corredera puede tener una duración fija o una duración de toda la vida. Las ventanas de duración fija se solapan, por lo que cada evento fuente puede contribuir a la agregación de características en varias ventanas. Una ventana de por vida incluye todos los eventos fuente antes de que termine la ventana. Las características en el momento t agregan datos de ventanas que finalizan en o antes de t (exclusivo). Windows está alineado con la época Unix.

class SlidingWindow(TimeWindow):
    window_duration: Optional[datetime.timedelta]
    slide_duration: datetime.timedelta
    delay: Optional[datetime.timedelta] = None
    offset: Optional[datetime.timedelta] = None
    start_time: Optional[datetime.datetime] = None

En la tabla siguiente se enumeran los parámetros de una ventana deslizante.

Parámetro Constraints
window_duration Debe ser positivo para una ventana de duración fija. Configurado para None una ventana de por vida.
slide_duration Debe ser positivo. Para una ventana de duración fija, también debe ser más corta que window_duration.
delay (opcional) Debe ser ≥ 0. Desplace la ventana analítica hacia atrás desde la marca temporal de la evaluación.
offset (opcional) Debe ser ≥ 0 y más corto que slide_duration. Desplaza los límites de la ventana desde la medianoche UTC.
start_time (opcional) Límite de tiempo de evento más temprano en el que la característica puede emitir una salida.
from databricks.feature_engineering.entities import SlidingWindow
from datetime import timedelta

window = SlidingWindow(
    window_duration=timedelta(days=7),
    slide_duration=timedelta(days=1),
    delay=timedelta(hours=2),
    offset=timedelta(hours=22),
)

Ejemplo de ventana deslizante

  • window_duration=timedelta(days=5), slide_duration=timedelta(days=1): esto crea ventanas superpuestas de 5 días que avanzan en un día cada vez. Ejemplo: La ventana n.º 1 abarca el día 0 al día 4, la ventana #2 abarca el día 1 al día 5, la ventana #3 abarca el día 2 al día 6, etc. Cada ventana incluye eventos desde el día de 00:00:00.00 inicio hasta (pero no incluido) 00:00:00.00 en el día de finalización. Dado que las ventanas se superponen, un único evento puede pertenecer a varias ventanas (en este ejemplo, cada evento pertenece a hasta 5 ventanas diferentes).

Ventana de por vida

Configura window_duration=None para crear una ventana de por vida. En cada límite de deslizamiento, la característica agrega todos los eventos fuente de la entidad con marcas de tiempo anteriores a ese límite. Por ejemplo, una diapositiva de un día produce un valor acumulado una vez al día.

Las ventanas de por vida solo son soportadas por SlidingWindow. RollingWindow y TumblingWindow requieren un finito window_duration.

Note

Las ventanas de por vida requieren una databricks-feature-engineering versión de cliente que soporte window_duration=None la habilitación del espacio de trabajo. Las versiones anteriores del cliente no soportan esta sintaxis.

from datetime import timedelta
from databricks.feature_engineering.entities import (
    AggregationFunction,
    DeltaTableSource,
    Feature,
    SlidingWindow,
    Sum,
)

lifetime_spend = Feature(
    source=DeltaTableSource(
        catalog_name="main",
        schema_name="store",
        table_name="transactions",
    ),
    entity=["user_id"],
    timeseries_column="transaction_time",
    function=AggregationFunction(
        Sum(input="amount"),
        SlidingWindow(
            window_duration=None,
            slide_duration=timedelta(days=1),
        ),
    ),
    name="lifetime_spend",
)

Ventana de diente de sierra

Importante

SawtoothWindow está en Beta.

Una ventana de dientes de sierra es una agregación que permite actualizaciones muy recientes de eventos recientes junto con la compactación diaria de datos históricos. Su borde de salida (antiguo) avanza en pasos fijos y diarios, mientras que el borde delantero (reciente) se mantiene actualizado con los últimos acontecimientos, por lo que la longitud efectiva de la ventana "sierra" a lo largo de cada día. La mayor parte de la ventana se sirve a partir de los datos de la tabla de ingesta del Stream, y solo los dos días más recientes provienen del live stream. Este es un compromiso que calcula eficientemente ventanas de larga duración (escalando a años) mientras se mantiene sensible a actualizaciones recientes.

Ventana en forma de diente de sierra: el borde de ataque sigue los últimos acontecimientos mientras el borde de salida avanza en pasos diarios, por lo que la ventana cubierta

Las ventanas de diente de sierra se materializan en un recorrido híbrido por lotes y en corriente. Una pipeline por lotes mantiene la mayor parte de la ventana, mientras que una pipeline de streaming mantiene los datos más recientes frescos en tiempo real. Ambos se fusionan en lectura, así que para el modelo o consumidor que sirve es una sola característica.

Como la parte histórica de la ventana se calcula mediante la tubería por lotes, una estructura de diente de sierra está lista para servir poco después de que comience la materialización, incluso cuando la ventana abarca meses o años. Una ventana enrollable solo se completa después de que haya transcurrido toda su duración de ventana. El mínimo window_duration debe ser mayor que dos días (el límite inferior impuesto). Databricks recomienda una ventana de diente de sierra para duraciones superiores a 7 días. Para ventanas de más de dos días y hasta siete días, elige entre la precisión de longitud fija de una ventana enrollable y la mayor preparación para la producción de una ventana de diente de sierra.

Note

Una característica de diente de sierra se basa en la historia que ya está presente. La tabla de ingestión del Flujo debe contener datos que cubran al menos la duración completa de la ventana, o la ventana calculada está incompleta. Antes de que transcurran 2 días completos, la función refleja solo los datos materializados hasta el momento. No se recomienda ofrecer la producción de la película hasta que hayan transcurrido 2 días completos. Una agregación sobre una ventana vacía devuelve 0 para Sum y , y nulo para Avg, Min, Max, First, Last, VarPop, VarSamp, , StddevPop, y StddevSampCount.

Para saber si una característica de dientes de sierra está lista, abre la Vista de Características en el Explorador de Catálogos. En la sección de características materializadas, el relleno por lotes se completa una vez que avanza el último tiempo de materialización de la característica y su estado muestra éxito. La parte de flujo está materializada por un oleoducto declarativo Lakeflow. Tras la validación de la Vista de Funcionalidades, la característica materializada enlaza a esa pipeline, donde puedes monitorizar su estado de ejecución.

Las ventanas de diente de sierra requieren un StreamSource y se materializan con StreamingMode.

class SawtoothWindow(TimeWindow):
    window_duration: datetime.timedelta

Los bordes de una ventana en forma de sierra se mueven de forma diferente a los de una ventana enrollable: el borde de ataque sigue el último evento, mientras que el borde de salida avanza una vez al día en lugar de continuamente. Cada día, a un límite fijo de las 18:00 UTC, el borde de salida avanza hasta el límite UTC-medianoche de ese día. Como resultado, la ventana efectiva es ligeramente más larga y window_duration crece a lo largo del día antes de volver a aparecer un día en el siguiente límite. El entrenamiento y el servicio utilizan el mismo límite de las 18:00 UTC, por lo que la formación offline y el servicio online se mantienen constantes.

Parámetro Constraints
window_duration Debe de ser más de dos días. Se permite una duración que no sea un número entero de días (por ejemplo, timedelta(days=3, minutes=15)), pero la ventana se actualiza con detalle diario.

Las ventanas de diente de sierra soportan las Sumfunciones de agregación , Avg, Count, Min, VarPopFirstMaxVarSampStddevPopLastStddevSamp y agregación.

Ejemplo de ventana de diente de sierra

El siguiente ejemplo muestra un recuento de 7 días de las transacciones de un usuario. El borde de ataque sigue el evento actual mientras que el borde de salida avanza día a día. Para los eventos del 10 de marzo, la ventana se remonta aproximadamente al 3 de marzo. A medida que avanza el 10 de marzo, el borde de ataque sigue avanzando mientras el borde de salida aguanta, por lo que el tramo cubierto crece. Luego, a principios del 11 de marzo, el borde de salida avanza hasta aproximadamente el 4 de marzo. La ventana efectiva siempre es un poco mayor de siete días. Los dos días más recientes se sirven desde la retransmisión en directo, y los días anteriores desde la tabla de ingesta de la emisión.

from databricks.feature_engineering.entities import SawtoothWindow
from datetime import timedelta

# 7-day window kept continuously fresh with streaming data
window = SawtoothWindow(window_duration=timedelta(days=7))

Limitaciones de la ventana de diente de sierra

  • No se admite el delay parámetro .
  • No se admite SourceLateness.settling_delay.
  • No se soportan funciones de agregación distintas de Sum, Avg, Count, MinMax, VarPopVarSampFirstLast, StddevPop, y StddevSamp (por ejemplo, ApproxCountDistinct, ApproxPercentile, LastNFirstN, , FirstDistinct, y ).LastDistinct
  • Las ventanas de diente de sierra requieren un StreamSource. A DeltaTableSource no está soportado.

Desencadenadores de materialización

Desencadena el control cuando se ejecuta una canalización de materialización. El tipo de desencadenador depende del tipo de característica.

CronSchedule

Úsalo CronSchedule para funciones de agregación por lotes. Por defecto, Azure Databricks obtiene un horario a partir de la ventana de agregación. Un calendario derivado tiene en cuenta el periodo de la ventana, la ventana delay y offset, y la fuente settling_delay , de modo que una ejecución no publica una ventana antes de que se espere que sus datos fuente estén completos. Los horarios derivados permiten ventanas de tumbling y deslizamiento.

Para solicitar un horario derivado, omite la expresión cron. CronSchedule() y la forma explícita CronSchedule(mode=CronScheduleMode.DERIVED) son equivalentes:

from databricks.feature_engineering.entities import (
    CronSchedule,
    CronScheduleMode,
)

trigger = CronSchedule(mode=CronScheduleMode.DERIVED)

No fijes quartz_cron_expression con CronScheduleMode.DERIVED. Cuando recuperas la característica materializada, el schedule devuelto puede contener la expresión cron que Azure Databricks calculó.

Para controlar el calendario directamente, proporciona una expresión cron de Cuarzo. CronScheduleMode.MANUAL se infiere cuando se proporciona una expresión:

from databricks.feature_engineering.entities import CronSchedule

trigger = CronSchedule(
    quartz_cron_expression="0 0 * * * ?",  # Hourly
    timezone_id="UTC",
)

TableTrigger

Uso TableTrigger para ColumnSelection características o agregación (AggregationFunction) respaldadas por un DeltaTableSource. La canalización se ejecuta cada vez que la tabla Delta ascendente recibe una nueva confirmación.

Para las funciones de agregación, la tubería está limitada para que no se ejecute en cada commit. La tubería se ejecuta como mucho una vez por cada mitad de la ventana de la función, pero nunca más a menudo que cada 5 minutos. Por ejemplo, una característica con una ventana de tumbling de 1 hora se ejecuta como máximo una vez cada 30 minutos, o una función con una ventana de 8 horas se ejecuta como máximo una vez cada 4 horas. El suelo de 5 minutos se aplica cuando la mitad de la ventana es más pequeña, así que las ventanas de 10 minutos o menos funcionan como máximo una vez cada 5 minutos. Las funciones de agregación cuyo periodo inferior a 5 minutos no pueden usar TableTrigger, utilizan un disparador de streaming en su lugar.

from databricks.feature_engineering.entities import TableTrigger

trigger = TableTrigger()

StreamingMode

Use StreamingMode para las características respaldadas por .StreamSource La canalización se ejecuta como una canalización de streaming continua.

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

Elección de un desencadenador

Cada característica utiliza un disparador; Las opciones por tipo de característica son:

Tipo de característica Trigger Cuando se ejecuta
Agregación (AggregationFunction) de DeltaTableSource CronSchedule En un horario derivado o manual
Agregación (AggregationFunction) de DeltaTableSource TableTrigger En cada confirmación de tabla de origen
ColumnSelection (de DeltaTableSource) TableTrigger En cada confirmación de tabla de origen
Características de StreamSource StreamingMode Streaming continuo

No se pueden materializar características que requieran diferentes tipos de desencadenadores en una sola materialize_features llamada. Emita llamadas independientes en su lugar.