Pular para o conteúdo principal

Referência da API de Views de recurso

info

Visualização

Este recurso está em Prévia Pública. Administradores do Workspace podem controlar o acesso a este recurso na página **Pré-visualizações**. Consulte Gerenciar prévias do Databricks.

Controle de acesso​

Recursos são objetos governáveis do Unity Catalog. O acesso a um recurso é controlado pelos privilégios CREATE FEATURE, READ FEATURE e MANAGE do Unity Catalog. Para descrições completas, consulte referência de privilégios do Unity Catalog.

  • CREATE FEATURE : necessário para criar um recurso em um esquema. create_feature e register_feature exigem CREATE FEATURE no esquema pai. Seguindo o princípio do privilégio mínimo, conceda CREATE FEATURE no nível do esquema; você também pode concedê-lo em um catálogo para permitir a criação de recursos em qualquer esquema desse catálogo.
  • READ FEATURE : necessário para ler metadados de recursos. get_feature, create_training_set e list_materialized_features exigem READ FEATURE no recurso. Esse privilégio não concede acesso aos dados de recursos na origem ou em tabelas de saída materializadas. Para ler esses dados para treinamento ou serving, você também deve ter SELECT nas tabelas aplicáveis. READ FEATURE concedido em um esquema ou catálogo se aplica a todos os recursos atuais e futuros que ele contém.
  • MANAGE : necessário para gerenciar o ciclo de vida e as concessões de um recurso. A exclusão de um recurso com delete_feature e a materialização de um recurso com materialize_features exigem MANAGE no recurso. A exclusão de um recurso materializado com delete_materialized_feature não é governada por MANAGE: apenas o criador do recurso materializado pode excluí-lo.

Todas as operações de recurso também exigem USE CATALOG no catálogo pai e USE SCHEMA no esquema pai. Para saber como MANAGE e READ FEATURE se aplicam à materialização, consulte Permissões.

API de View de recurso​

ConstrutorFeature e register_feature()​

A abordagem recomendada é construir um objeto Feature localmente e usar register_feature para persistir no Unity Catalog. Este fluxo de trabalho de duas etapas permite experimentar recursos (incluindo create_training_set) antes de registrá-los.

Python
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 um Feature construído localmente no Unity Catalog.

Python
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
Python
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, constrói e registra imediatamente um recurso no Unity Catalog em um único passo. Use isto quando você não precisar experimentar o recurso localmente primeiro.

Python
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: A fonte de dados usada na computação de recurso (DeltaTableSource, StreamSource, RequestSource ou FeatureViewSource).
  • function: um AggregationFunction que agrupa um operador e uma janela de tempo, ColumnSelection("column_name") para recursos de passagem ou CustomUDF para transformações linha a linha. Consulte Funções compatíveis para ver os tipos de origem compatíveis.
  • catalog_name: O nome do catálogo do Unity Catalog para o recurso.
  • schema_name: O nome do esquema do Unity Catalog para o recurso.
  • entity: Lista de nomes de colunas que definem a agregação ou as chaves de busca (chaves primárias). Obrigatório para DeltaTableSource e StreamSource. Por exemplo, ["user_id"] realiza agregações ou buscas por usuário. Omitir para RequestSource e FeatureViewSource.
  • timeseries_column: A coluna timestamp usada para agregação de janela de tempo ou seleção de valor mais recente. Necessário para DeltaTableSource e StreamSource. Omitir para RequestSource e FeatureViewSource.
  • name: Nome de recurso opcional. Se omitido, gerado automaticamente a partir da coluna de entrada, função e janela (por exemplo, amount_avg_rolling_7d).
  • description: Descrição opcional do recurso.

Retorna: Uma instância de recurso validada

Gera: ValueError se alguma validação falhar

delete_feature()​

Exclui um recurso do Unity Catalog pelo seu nome totalmente qualificado.

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

Antes de excluir um recurso, remova ou atualize todos os modelos ou especificações de recurso que façam referência a ele. Não é possível excluir um recurso enquanto ele ainda tiver recursos materializados. Exclua os recursos materializados primeiro e, em seguida, exclua o recurso. Consulte Como excluir um recurso materializado.

Nomes gerados automaticamente​

Quando name é omitido, um nome é gerado automaticamente. Os nomes gerados seguem o padrão: {column}_{function}_{window}. Por exemplo:

  • price_avg_rolling_1h (preço médio de 1 hora)
  • transaction_count_rolling_30d_1d (Contagem de 30 dias de transação com atraso de 1 dia do timestamp do evento)

Funções suportadas​

Funções de agregação​

nota

As funções de agregação são encapsuladas em um AggregationFunction junto com uma janela de tempo, conforme descrito em janelas de tempo. Cada função recebe um parâmetro input que especifica a coluna de origem a ser agregada.

Função

Descrição

Exemplo de caso de uso

Sum(input="column")

Total de valores

Uso diário do aplicativo por usuário em minutos

Avg(input="column")

Média de valores

Valor médio da transação

Count(input="column")

Número de registros

Número de logins por usuário

Min(input="column")

Valor mínimo

Menor frequência cardíaca registrada por um dispositivo wearable

Max(input="column")

Valor máximo

Maior valor de transação por sessão.

StddevPop(input="column")

Desvio padrão populacional

Variabilidade diária do valor da transação entre todos os clientes

StddevSamp(input="column")

Desvio padrão de amostra

Variabilidade das taxas de cliques de campanhas de anúncios

VarPop(input="column")

Variância populacional

Dispersão das leituras do sensor para dispositivos IoT em uma fábrica.

VarSamp(input="column")

Variância de amostra

Distribuição de avaliações de filmes em um grupo amostrado

ApproxCountDistinct(input="column", relativeSD=0.05)

Contagem única aproximada

Contagem distinta de itens comprados

ApproxPercentile(input="column", percentile=0.95, accuracy=100)

Percentil aproximado

Latência de resposta p95

First(input="column")

Primeiro valor

Primeiro login Timestamp

Last(input="column")

Último valor

Valor de compra mais recente

FirstN(input="column", n=3)

Primeiros n valores como uma matriz

Três primeiros produtos visualizados em uma sessão

LastN(input="column", n=3)

Últimos n valores como uma matriz

Três status de caso de suporte mais recentes

FirstDistinct(input="column", n=3)

Primeiros n valores distintos como uma matriz

Três primeiras categorias de produto distintas visualizadas

LastDistinct(input="column", n=3)

Últimos n valores distintos como uma matriz

Três categorias de comerciante distintas mais recentes

Função

Descrição

Exemplo de caso de uso

Sum(input="column")

Total de valores

Uso diário do aplicativo por usuário em minutos

Avg(input="column")

Média de valores

Valor médio da transação

Count(input="column")

Número de registros

Número de logins por usuário

Min(input="column")

Valor mínimo

Menor frequência cardíaca registrada por um dispositivo wearable

Max(input="column")

Valor máximo

Maior valor de transação por sessão.

StddevPop(input="column")

Desvio padrão populacional

Variabilidade diária do valor da transação entre todos os clientes

StddevSamp(input="column")

Desvio padrão de amostra

Variabilidade das taxas de cliques de campanhas de anúncios

VarPop(input="column")

Variância populacional

Dispersão das leituras do sensor para dispositivos IoT em uma fábrica.

VarSamp(input="column")

Variância de amostra

Distribuição de avaliações de filmes em um grupo amostrado

ApproxCountDistinct(input="column", relativeSD=0.05)

Contagem única aproximada

Contagem distinta de itens comprados

ApproxPercentile(input="column", percentile=0.95, accuracy=100)

Percentil aproximado

Latência de resposta p95

First(input="column")

Primeiro valor

Primeiro login Timestamp

Last(input="column")

Último valor

Valor de compra mais recente

FirstN(input="column", n=3)

Primeiros n valores como uma matriz

Três primeiros produtos visualizados em uma sessão

LastN(input="column", n=3)

Últimos n valores como uma matriz

Três status de caso de suporte mais recentes

FirstDistinct(input="column", n=3)

Primeiros n valores distintos como uma matriz

Três primeiras categorias de produto distintas visualizadas

LastDistinct(input="column", n=3)

Últimos n valores distintos como uma matriz

Três categorias de comerciante distintas mais recentes

nota

First, Last, FirstN, LastN, FirstDistinct e LastDistinct incluem valores nulos por default. Para ignorar nulos, adicione um filter_condition que exclua explicitamente colunas de entrada que sejam nulas.

FirstN, LastN, FirstDistinct e LastDistinct usam o timeseries_column do recurso para ordenar linhas de entrada e retornar uma matriz contendo até n valores. O parâmetro n deve ser um número inteiro positivo. FirstN e FirstDistinct selecionam valores do mais antigo para o mais recente. LastN e LastDistinct selecionam valores do mais recente para o mais antigo e, em seguida, retornam os valores selecionados na ordem de timestamp. FirstDistinct e LastDistinct removem valores duplicados ao selecionar valores nessa direção.

Por exemplo, se as linhas de origem de uma entidade forem ordenadas por event_time como ["A", "A", "B", "C", "B", "B"], as seguintes funções retornam:

Função

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"]

Função

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 e LastDistinct requerem a versão 0.17.0 ou posterior do databricks-feature-engineering.

CustomUDF​

CustomUDF aplica uma função definida pelo usuário (UDF) Python registrada no Unity Catalog a cada linha. Use-o para transformar entradas de solicitação ou combinar valores de recursos. Ele não agrega linhas ou define um intervalo temporal.

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

input_bindings maps each UDF parameter name to an input. For RequestSource, the input is a source column name. For FeatureViewSource, it is an upstream recurso reference. Bind every UDF parameter, including parameters with defaults. Input types must match the UDF parameter types exactly, without implicit numeric casts. Use scalar input and return types.

Origem

Comportamento

RequestSource

Transforma colunas do DataFrame de treinamento ou da requisição de inferência.

FeatureViewSource

Combina valores de recurso upstream. Veja FeatureViewSource.

Origem

Comportamento

RequestSource

Transforma colunas do DataFrame de treinamento ou da requisição de inferência.

FeatureViewSource

Combina valores de recurso upstream. Veja FeatureViewSource.

Os recursos CustomUDF com suporte do Delta não podem ser materializados ou disponibilizados online. Para transformar os valores de recurso baseados em tabela para treinamento e disponibilização, defina uma agregação baseada em Delta ou um recurso de seleção de coluna e faça referência a ele por meio de FeatureViewSource.

CustomUDF não é compatível com StreamSource. Para transformar a saída de um recurso de transmissão, faça referência a esse recurso por meio de FeatureViewSource.

CustomUDF with RequestSource requires databricks-feature-engineering version 0.17.0 or later.

Para usar um CustomUDF, você precisa do privilégio EXECUTE no UDF, do privilégio USE CATALOG no seu catálogo pai e do privilégio USE SCHEMA no seu esquema pai.

The following example uses NumPy to compute log(1 + amount), reducing the escala of large transaction amounts. Run it on serverless compute with custom UDF dependencies enabled. The main.ecommerce schema must exist.

Python
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 um recurso que vincula a coluna de solicitação transaction_amount ao parâmetro UDF amount:

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

O ENVIRONMENT do UDF configura dependências para computação offline. Para o serving online, declare também os pacotes em create_feature_spec(extra_pip_requirements=...) ou log_model(extra_pip_requirements=...). Eles não são copiados do UDF automaticamente. Consulte dependências do Feature Serving e dependências do modelo.

CustomUDF recursos não podem ser materializados. UDFs baseadas em solicitação e em recursos são executadas sob demanda durante o treinamento e a disponibilização. Cada UDF em uma cadeia de dependência adiciona computação, portanto, mantenha funções e cadeias pequenas. As UDFs devem lidar com entradas ausentes, que podem ser None offline ou NaN online.

For guidance on handling missing values, see How to handle missing recurso values.

ColumnSelection (pass-through)​

ColumnSelection seleciona uma única coluna de uma fonte sem aplicar nenhuma agregação. É encapsulado diretamente no parâmetro function (não dentro de AggregationFunction). O tipo de retorno é inferido do esquema de origem.

Função

Descrição

Exemplo de caso de uso

ColumnSelection("col")

Último valor de uma coluna (sem agregação)

Categoria de fornecedor mais recente, pass-through de um campo de solicitação

Função

Descrição

Exemplo de caso de uso

ColumnSelection("col")

Último valor de uma coluna (sem agregação)

Categoria de fornecedor mais recente, pass-through de um campo de solicitação

ColumnSelection supports the following fontes de dados:

  • DeltaTableSource : Retorna o valor mais recente por chave de entidade via um join pontual (sem agregação de janela de retrospectiva).
  • StreamSource : Retorna o valor mais recente por chave de entidade da transmissão (sem agregação de janela de retrospectiva).
  • RequestSource : Passa o valor fornecido no momento da inferência (ou extraído do DataFrame rotulado no momento do treinamento).

Para DeltaTableSource, os recursos ColumnSelection oferecem suporte a filter_condition e transformation_sql, aplicados antes da seleção do valor mais recente, da mesma forma que os recursos de agregação.

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

Exemplo: recursos de agregação e seleção de colunas​

O exemplo a seguir mostra recursos definidos sobre a mesma fonte de dados.

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

Recursos com condições de filtro​

O parâmetro filter_condition permite filtrar linhas da tabela de origem antes de calcular agregações ou selecionar o valor da coluna mais recente. Isso funciona como uma cláusula SQL WHERE aplicada antes de agrupar e agregar os dados.

nota

Para recursos de agregação, filter_condition filtra as linhas antes da agregação, como uma cláusula WHERE do SQL aplicada antes de GROUP BY. Isso não altera a granularidade, que é sempre definida por entity na definição do recurso.

Filtros são úteis ao trabalhar com grandes tabelas de origem que incluem um superconjunto de dados necessários para o compute de recursos, e minimizam a necessidade de criar views separadas sobre essas tabelas.

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

fonte de dados​

DeltaTableSource​

DeltaTableSource é um objeto Python efêmero usado para definir como os recursos são computados a partir de uma tabela de origem. Não cria uma nova tabela. Ele especifica a configuração para ler dados e agregar recursos.

Python
DeltaTableSource(
catalog_name: str, # Required: Catalog name
schema_name: str, # Required: Schema name
table_name: str, # Required: Table name
filter_condition: Optional[str] = None, # Optional: SQL WHERE clause to filter source data
transformation_sql: Optional[str] = None, # Optional: SQL SELECT expression for column transformations
dataframe_schema: Optional[str] = None, # Required if transformation_sql is set: schema of the resulting DataFrame
lateness: Optional[SourceLateness] = None, # Optional: Expected source settling time
)

Parâmetros:

  • catalog_name, schema_name, table_name: Identifique a tabela Delta de origem no Unity Catalog.
  • filter_condition: Uma cláusula SQL WHERE aplicada antes da agregação ou da seleção de colunas. Exemplo: "status = 'completed'".
  • transformation_sql: Uma expressão SQL SELECT aplicada à tabela de origem. Use isto para renomear colunas, converter tipos ou compute colunas derivadas antes da agregação ou seleção de colunas. Se omitido, todas as colunas serão selecionadas (*). Exemplo: "user_id, CAST(amount AS DOUBLE) AS amount, event_time".
  • dataframe_schema: o esquema do DataFrame resultante após as transformações, no formato JSON Spark StructType (de df.schema.json()). Obrigatório se transformation_sql for fornecido. Isso informa ao sistema os nomes e tipos de coluna que resultam da sua transformação.
  • lateness: um objeto SourceLateness que descreve quanto tempo a origem normalmente leva para ficar completa no tempo do evento. Se omitida, a origem será considerada imediatamente concluída.

Quando filter_condition e transformation_sql estão definidos, a query resultante é: SELECT {transformation_sql} FROM {table} WHERE {filter_condition}.

SourceLateness.settling_delay é a maneira recomendada de simular durante o treinamento um atraso de ETL consistente que afeta a materialização online. O Databricks retrocede o tempo de avaliação de treinamento elegível por esta duração para que um exemplo de treinamento não use dados que ainda estariam em trânsito online. Durante a materialização, o Databricks aguarda pela mesma duração antes de publicar uma janela concluída e serve a última janela concluída durante o período intermediário.

For example, suppose a daily ETL job completes 8 hours after midnight in a local time zone where midnight corresponds to 07:00 UTC. Use um atraso de estabilização de 8 horas e um deslocamento de janela de 7 horas:

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

O timeseries_column deve ser do tipo TimestampType ou TimestampNTZType. DateType não é compatível com séries temporais; converta a coluna para TimestampType primeiro (por exemplo, com transformation_sql).

Exemplo: Usando transformation_sql para transformações de coluna

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

Exemplo: Derivando transformation_sql e dataframe_schema de um PySpark DataFrame

É possível escrever a transformação como uma query PySpark e, então, extrair o esquema do DataFrame resultante:

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

Expressões transformation_sql suportadas​

As mesmas regras se aplicam a transformation_sql em DeltaTableSource e StreamSource.

transformation_sql suporta quaisquer expressões linha a linha; operações avaliadas independentemente para cada linha. Elas não alteram o número de linhas ou a correspondência um-para-um com a origem. As expressões linha a linha incluem renomeações de colunas, conversões, operações aritméticas e muito mais.

Operações que alteram a forma ou a contagem de linhas não são compatíveis, como agregações como SUM() ou COUNT(). Use AggregationFunction na definição do recurso em vez disso.

DeltaTableSource.from_sql()​

Por conveniência, você pode criar um DeltaTableSource a partir de uma query SQL. O método analisa a query para extrair automaticamente o nome da tabela, transformation_sql e filter_condition.

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

Apenas SELECT ... FROM ... [WHERE ...] queries simples são compatíveis. SQL complexo (JOINs, subqueries, CTEs, UNIONs) é rejeitado. Para queries complexas, construa DeltaTableSource diretamente com transformation_sql e filter_condition.

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

Iterar com to_dataframe()​

Use source.to_dataframe() para visualizar os dados que serão usados para o compute de recursos. Isso é útil para iterar em filter_condition e transformation_sql até que produzam os resultados esperados.

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

Compreendendo as entidades​

As colunas de entidade definem o nível de agregação para seus recursos. Elas são especificadas na definição de Feature, e não em DeltaTableSource. As entidades determinam:

  • Como os dados são agrupados : Os recursos são agregados por combinação única de valores de entidade (semelhante a GROUP BY em SQL)
  • A estrutura da key primária : Cada combinação de entidade exclusiva resulta em uma linha de recursos de compute

Exemplo: Recursos em nível de cliente.

O código a seguir agrega recursos no nível do cliente (uma linha por cliente):

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

Exemplo: Recursos no nível do cliente-loja

Para agregar recursos em um nível mais detalhado (uma linha por combinação cliente-loja), use múltiplas colunas de entidade:

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

Quando precisar de recursos em diferentes níveis de agregação (por exemplo, nível de cliente e nível de cliente-loja), use valores de entity diferentes em suas definições de recursos. O mesmo DeltaTableSource pode ser compartilhado entre recursos com diferentes configurações de entidade.

StreamSource​

StreamSource referencia uma Transmissão. A Transmissão contém conexão, autenticação, esquema e configuração de ingestão para a fonte de transmissão. Para Kafka, as referências de coluna nas definições de recurso devem ser prefixadas com value. ou key. para indicar qual parte da mensagem ler.

Python
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: O nome completo de três partes de uma transmissão (por exemplo, "my_catalog.my_schema.my_stream").
  • filter_condition (opcional): Uma cláusula SQL WHERE aplicada a dados de transmissão antes da agregação, usando referências de coluna com prefixo de ponto (por exemplo, "value.event_type = 'purchase'").
  • transformation_sql (opcional): uma expressão SQL SELECT aplicada antes da agregação ou seleção de coluna, usando referências prefixadas por ponto para as structs key e value. Suporta as mesmas expressões linha a linha que DeltaTableSource. Se omitido, a origem usa todas as colunas (*).
  • dataframe_schema: O esquema JSON Spark StructType da saída projetada. Obrigatório se você definir transformation_sql.
  • lateness: A SourceLateness object that describes how long the transmissão normally takes to become complete in event time. See SourceLateness.settling_delay.
Python
from databricks.feature_engineering.entities import StreamSource

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

Derive dataframe_schema executando a projeção na tabela de ingestão da transmissão, que expõe as estruturas key e value.

Python
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 um esquema para dados que são fornecidos no momento da inferência na carga útil da solicitação, em vez de pesquisados em uma tabela pré-materializada. Durante o treinamento, estas colunas são extraídas do DataFrame rotulado passado para create_training_set. Durante o servindo modelo, o chamador deve incluí-los na carga útil da solicitação HTTP.

RequestSource pode ser usado com funções de Feature View do CustomUDF ou ColumnSelection. Ele não oferece suporte a funções de agregação ou janelas de tempo.

Definindo o esquema​

Defina o esquema como uma lista de objetos FieldDefinition, cada um especificando um nome de coluna e um ScalarDataType:

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

Tipos de dados compatíveis​

RequestSource suporta os tipos escalares definidos em ScalarDataType: INTEGER, FLOAT, BOOLEAN, STRING, DOUBLE, LONG, TIMESTAMP, DATE, SHORT. Tipos complexos como matrizes, mapas e structs não são compatíveis.

Como os dados da solicitação são hidratados​

Contexto

Comportamento

**Treinamento** create_training_set()

Colunas são extraídas do DataFrame rotulado. Os tipos são validados em relação ao esquema declarado. Incompatibilidades geram um erro (sem conversão implícita).

**Serviço** (Endpoint de modelo)

As colunas são extraídas de dataframe_records ou dataframe_split na solicitação HTTP. Os valores JSON são convertidos para os tipos declarados (por exemplo, número JSON → DOUBLE).

Contexto

Comportamento

**Treinamento** create_training_set()

Colunas são extraídas do DataFrame rotulado. Os tipos são validados em relação ao esquema declarado. Incompatibilidades geram um erro (sem conversão implícita).

**Serviço** (Endpoint de modelo)

As colunas são extraídas de dataframe_records ou dataframe_split na solicitação HTTP. Os valores JSON são convertidos para os tipos declarados (por exemplo, número JSON → DOUBLE).

Assinatura do modelo​

Quando um modelo é registrado usando log_model com um conjunto de treinamento que inclui RequestSource recursos, as colunas RequestSource são adicionadas à assinatura do modelo MLflow como entradas obrigatórias. Isso significa que o esquema da API do Endpoint de serviço reflete quais campos os chamadores devem fornecer no momento da inferência.

FeatureViewSource​

FeatureViewSource uses the outputs of other Feature Views as inputs to a CustomUDF. O encadeamento de recursos cria um gráfico acíclico direcionado (DAG). Por exemplo, um recurso de margem pode combinar agregados de receita e custo, e outro recurso pode transformar a margem.

Use a databricks-feature-engineering versão 0.18.0 ou posterior para FeatureViewSource.

Passe uma lista de objetos Feature para features, não strings de nome de recurso. Recupere recursos registrados com get_feature. Em input_bindings, use o full_name de cada recurso registrado. Para um recurso local não registrado, use seu name em vez disso.

O exemplo a seguir pressupõe dois recursos registrados, revenue_sum_7d e cost_sum_7d, que retornam valores DOUBLE por customer_id e usam event_time para computação pontual:

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

Aplicam-se as seguintes restrições:

  • Apenas CustomUDF é suportado como a função. Os recursos upstream podem ser agregações, seleções de colunas ou outros recursos CustomUDF.
  • Omita entity e timeseries_column no recurso derivado. Cada recurso upstream retém sua própria entidade, Timestamp e definição de janela.
  • Um recurso tem uma fonte. Para combinar um valor de solicitação com um recurso apoiado por tabela, defina um recurso RequestSource e faça referência a ambos por meio de FeatureViewSource.
  • Todo recurso upstream declarado deve ser usado em input_bindings. Ciclos não são permitidos.
  • Registre recursos upstream antes de registrar o recurso derivado. Gráficos locais e não registrados podem ser usados com create_training_set para experimentação.
  • Para treinamento ou disponibilização, você precisa do privilégio READ FEATURE ou MANAGE no recurso derivado e seus recursos upstream transitivos. Use nomes de recursos distintos em todo o gráfico para registro e disponibilização, mesmo em catálogos ou esquemas.
  • Um recurso pode fazer referência a até 20 recursos upstream diretos. Gráficos registrados oferecem suporte a uma profundidade máxima de cinco recursos ao longo de um caminho de dependência, incluindo o recurso base.
  • FeatureViewSource os recursos não podem ser materializados ou avaliados com compute_features. Use create_training_set para avaliá-los offline. Para a disponibilização online, materialise os recursos upstream baseados em tabela suportados.

For dependency evaluation and output selection, see Ensinar com FeatureViewSource features. For deployment, see Serve derived recursos.

API de treinamento e inferência​

create_training_set e score_batch compute valores de recurso corretos em determinado momento sob demanda a partir dos dados de origem. Para recursos que suportam a materialização offline, como agregações de janela deslizante em fontes de tabela delta, materializar os recursos primeiro em um armazenamento offline melhora o desempenho de ambas as operações. Quando recursos offline materializados estão disponíveis, as operações leem dados offline pré-computados em vez de recalcular os valores de recurso da fonte. Consulte Materializar Views de Recurso para materializar recursos em um armazenamento offline.

create_training_set()​

Cria um dataset de treinamento com cálculo de recurso correto pontual. Para obter detalhes, consulte Ensinar modelos com Recurso Views.

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

Logs um modelo com metadados de recurso para acompanhamento de linhagem e consulta automática de recursos durante a inferência. Para obter detalhes, consulte Ensinar modelos com Recurso Views.

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

Executa inferência em lote offline com pesquisa automática de recursos. Utiliza os metadados de recursos armazenados com o modelo para calcular recursos corretos para um ponto no tempo, garantindo consistência com o treinamento.

Python
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

O DataFrame de entrada deve conter as colunas de entidade e séries temporais usadas durante o treinamento. Recursos são automaticamente computados a partir dos dados de origem.

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

Janelas de tempo​

As Feature Views dão suporte a quatro tipos de janela para controlar o comportamento de retrospectiva para agregações baseadas em janela de tempo. Os tipos de janela disponíveis dependem da origem do recurso: recursos de origem de transmissão podem usar janelas móveis (rolling) e em dente de serra (sawtooth), e recursos de origem em lotes podem usar janelas móveis (rolling), em salto (tumbling) e deslizantes (sliding).

  • Janelas contínuas retroagem a partir do tempo do evento. Duração e atraso são explicitamente definidos.
  • Janelas em cascata são janelas de tempo fixas e não sobrepostas. Cada ponto de dados pertence a exatamente uma janela.
  • As janelas deslizantes são janelas de tempo sobrepostas e contínuas com um intervalo de deslizamento configurável.
  • As janelas sawtooth mantêm uma longa janela de retrospectiva atualizada sobre uma fonte de transmissão usando um caminho híbrido de lotes e transmissão. Consulte Janela sawtooth.

A ilustração a seguir mostra os tipos de janela tumbling, sliding, rolling e sawtooth.

Janelas de retrospectiva do tipo tumbling, sliding, rolling e sawtooth.

Temporização da janela de tempo​

Use delay para avaliar uma janela em um ponto analítico anterior no tempo. Por exemplo, uma janela de 30 dias com um atraso de 7 dias calcula um valor de 30 dias a partir de uma semana antes do horário de avaliação. delay é independente do horário de chegada da fonte. Para modelar o tempo que os dados de origem levam para chegar, configure SourceLateness.settling_delay em vez disso.

When both settings are present, they compose. O Databricks trata a janela como concluída após o atraso de estabilização da origem e a avalia usando o atraso analítico.

Use offset to change the alignment of fixed window boundaries. By default, tumbling windows and sliding windows are aligned to midnight UTC. For example, an offset of 22 hours aligns a daily boundary to 22:00 UTC. To approximate boundaries in a local time zone, configure a static offset relative to UTC. The offset does not adjust for daylight saving time, shift the evaluation time, or model late-arriving data.

A tabela a seguir resume o suporte a esses campos:

campo

Supported windows

Restrição

delay

Rolling, tumbling, and sliding

Deve ser não negativo datetime.timedelta

offset

Fixas e deslizantes

Deve ser não negativo e menor que o período*

SourceLateness.settling_delay

Recursos (features) deslizantes, fixos e cumulativos

Deve ser não negativo datetime.timedelta

start_time

Rolling, tumbling, and sliding

Deve ser um datetime.datetime

campo

Supported windows

Restrição

delay

Rolling, tumbling, and sliding

Deve ser não negativo datetime.timedelta

offset

Fixas e deslizantes

Deve ser não negativo e menor que o período*

SourceLateness.settling_delay

Recursos (features) deslizantes, fixos e cumulativos

Deve ser não negativo datetime.timedelta

start_time

Rolling, tumbling, and sliding

Deve ser um datetime.datetime

*Período: Para uma janela fixa (tumbling window), o deslocamento deve ser menor que window_duration. Para uma janela deslizante (sliding window), ela deve ser menor que slide_duration.

Hora de início​

Use start_time para definir o limite de horário do evento mais antigo em UTC no qual um recurso pode emitir uma saída. O limite é inclusivo. start_time controla as saídas. Isso não restringe as linhas de origem históricas que uma janela pode ler e não altera o alinhamento da janela. Se start_time estiver entre dois limites alinhados, a primeira saída de janela fixa elegível será o próximo limite.

Com start_time, janelas de duração fixa podem ser emitidas antes que a duração total de uma janela tenha decorrido na origem. Estes primeiros resultados usam qualquer história de origem disponível. Por exemplo, considere uma janela deslizante com um window_duration de um ano e um slide_duration de um dia, sobre uma origem cujos dados começam em 1º de janeiro de 2024:

  • Sem start_time, o recurso emite pela primeira vez em 1º de janeiro de 2025, assim que uma janela completa de um ano puder ser formada.
  • Com start_time definido como 21 de agosto de 2024, o recurso é emitido pela primeira vez em 21 de agosto de 2024. Essa saída cobre apenas a história de origem disponível até o momento, a partir de 1º de janeiro de 2024. A janela atinge sua extensão completa de um ano em 1º de janeiro de 2025 e produz saídas completas a partir de então.

Como start_time não altera o alinhamento da janela, um valor entre dois limites alinhados não cria um novo limite. Para uma janela deslizante com limites diários à meia-noite UTC, um start_time de 06:00 UTC é emitido primeiro no próximo limite de meia-noite. Um start_time que cai exatamente em um limite é emitido nesse limite, porque o limite é inclusivo.

Se start_time não estiver definido, as janelas deslizantes (tumbling windows) e as janelas deslizantes de duração fixa emitem primeiro em um limite alinhado após a formação de uma janela completa. As janelas deslizantes de tempo de vida (lifetime sliding windows) e as janelas correntes (rolling windows) são emitidas assim que houver dados de origem elegíveis.

nota

start_time é compatível com recursos de lote que usam DeltaTableSource com janelas rolantes, basculantes ou deslizantes. Não há compatibilidade com StreamSource ou SawtoothWindow.

Por exemplo:

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

Janela deslizante​

nota

RollingWindow era anteriormente denominado ContinuousWindow. Se estiver migrando de uma versão anterior do SDK, atualize suas importações de acordo.

Janelas deslizantes são agregados atualizados e em tempo real, tipicamente usados em dados de transmissão. Em pipelines de transmissão, a janela deslizante emite uma nova linha apenas quando o conteúdo da janela de comprimento fixo muda, como quando um evento entra ou sai. Quando um recurso de janela deslizante é usado em pipelines de treinamento, um cálculo preciso de recurso pontual é realizado nos dados de origem usando a duração de janela de comprimento fixo imediatamente anterior ao Timestamp de um evento específico. Isso ajuda a evitar distorções online-offline ou vazamento de dados. Recursos no tempo T agregam eventos de [T − duração, T).

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

A tabela a seguir lista os parâmetros para uma janela contínua. Os horários de início e término da janela são baseados nestes parâmetros da seguinte forma:

  • Hora de início: evaluation_time - window_duration - delay (inclusive)
  • Hora de término: evaluation_time - delay (exclusiva)

Parâmetro

Restrições

delay (opcional)

Deve ser ≥ 0. Desloca a janela analítica para trás em relação ao Timestamp de avaliação. Use SourceLateness.settling_delay para modelar uma linha de base consistente para o atraso de chegada da origem na sua transmissão.

window_duration

Deve ser > 0

start_time (opcional)

Limite de tempo do evento mais antigo em que o recurso pode emitir uma saída.

Parâmetro

Restrições

delay (opcional)

Deve ser ≥ 0. Desloca a janela analítica para trás em relação ao Timestamp de avaliação. Use SourceLateness.settling_delay para modelar uma linha de base consistente para o atraso de chegada da origem na sua transmissão.

window_duration

Deve ser > 0

start_time (opcional)

Limite de tempo do evento mais antigo em que o recurso pode emitir uma saída.

Python
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 uma janela contínua com atraso usando o código abaixo.

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

Exemplos de janela deslizante​

  • window_duration=timedelta(days=7): Isso cria uma janela de retrospectiva de 7 dias terminando no tempo de avaliação atual. Para um evento às 14:00 no Dia 7, isto inclui todos os eventos das 14:00 do Dia 0 até (mas sem incluir) as 14:00 do Dia 7.

  • window_duration=timedelta(hours=1), delay=timedelta(minutes=30): isso cria uma janela de retrospectiva de 1 hora que termina 30 minutos antes do tempo de avaliação. Para um evento às 15:00, isso inclui todos os eventos das 13:30 até (mas sem incluir) as 14:30.

Use Last para limitar a atualização de um valor mais recente​

Combine Last com RollingWindow quando um valor mais recente for válido apenas por um tempo limitado. Em um momento de avaliação, o recurso retorna o valor da linha com o timestamp mais recente neste intervalo:

[evaluation_time - delay - window_duration, evaluation_time - delay)

Se a última linha no intervalo contiver um valor nulo, o recurso retornará null. Se você quiser excluir valores de entrada nulos, defina um filter_condition na origem.

Esta combinação difere de ColumnSelection. O ColumnSelection retorna o último valor observado não nulo sem expirá-lo com base na idade.

Para recursos de lote, esta combinação tem um modo especial de materialização exclusivo para online. Ele oferece suporte apenas a DeltaTableSource, Last, RollingWindow e TableTrigger. Consulte Materializar valores mais recentes delimitados por frescor.

Janela em cascata​

Para recursos definidos usando janelas deslizantes, as agregações são calculadas em uma janela de comprimento fixo pré-determinada que avança por um intervalo de slide, produzindo janelas não sobrepostas que particionam totalmente o tempo. Como resultado, cada evento na origem contribui para exatamente uma janela. Recursos no tempo t agregam dados de janelas terminando em ou antes de t (exclusivo). Windows começam na época Unix.

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

A tabela a seguir lista os parâmetros para uma janela em cascata.

Parâmetro

Restrições

window_duration

Deve ser > 0

delay (opcional)

Must be >= 0. Shifts the analytic window backward from the evaluation timestamp.

offset (opcional)

Deve ser ≥ 0 e inferior a window_duration. Desloca os limites da janela a partir da meia-noite UTC.

start_time (opcional)

Limite de tempo do evento mais antigo em que o recurso pode emitir uma saída.

Parâmetro

Restrições

window_duration

Deve ser > 0

delay (opcional)

Must be >= 0. Shifts the analytic window backward from the evaluation timestamp.

offset (opcional)

Deve ser ≥ 0 e inferior a window_duration. Desloca os limites da janela a partir da meia-noite UTC.

start_time (opcional)

Limite de tempo do evento mais antigo em que o recurso pode emitir uma saída.

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

Exemplo de janela de rolagem​

  • window_duration=timedelta(days=5): Isso cria janelas de comprimento fixo predeterminadas de 5 dias cada. Exemplo: a Janela #1 abrange do Dia 0 ao Dia 4, a Janela #2 abrange do Dia 5 ao Dia 9, a Janela #3 abrange do Dia 10 ao Dia 14 e assim por diante. Especificamente, a Janela #1 inclui todos os eventos com timestamps que começam em 00:00:00.00 no Dia 0 até (mas não incluindo) quaisquer eventos com timestamp 00:00:00.00 no Dia 5. Cada evento pertence a exatamente uma janela.

Janela deslizante​

Para recursos definidos usando janelas deslizantes, as agregações são calculadas em uma janela que avança por um intervalo de deslizamento. Uma janela deslizante pode ter uma duração fixa ou uma duração de tempo de vida. As janelas de duração fixa se sobrepõem, portanto, cada evento de origem pode contribuir para a agregação de recursos para várias janelas. Uma janela de tempo de vida inclui todos os eventos de origem antes do final da janela. Os recursos no momento t agregam dados de janelas que terminam em ou antes de t (exclusivo). Windows são alinhadas à era Unix.

Python
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

A tabela a seguir lista os parâmetros para uma janela deslizante.

Parâmetro

Restrições

window_duration

Deve ser positivo para uma janela de duração fixa. Defina como None para uma janela de tempo de vida.

slide_duration

Deve ser positivo. Para uma janela de duração fixa, ela também deve ser inferior a window_duration.

delay (opcional)

Must be >= 0. Shifts the analytic window backward from the evaluation timestamp.

offset (opcional)

Deve ser ≥ 0 e inferior a slide_duration. Desloca os limites da janela a partir da meia-noite UTC.

start_time (opcional)

Limite de tempo do evento mais antigo em que o recurso pode emitir uma saída.

Parâmetro

Restrições

window_duration

Deve ser positivo para uma janela de duração fixa. Defina como None para uma janela de tempo de vida.

slide_duration

Deve ser positivo. Para uma janela de duração fixa, ela também deve ser inferior a window_duration.

delay (opcional)

Must be >= 0. Shifts the analytic window backward from the evaluation timestamp.

offset (opcional)

Deve ser ≥ 0 e inferior a slide_duration. Desloca os limites da janela a partir da meia-noite UTC.

start_time (opcional)

Limite de tempo do evento mais antigo em que o recurso pode emitir uma saída.

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

Exemplo de janela deslizante​

  • window_duration=timedelta(days=5), slide_duration=timedelta(days=1): Isso cria janelas sobrepostas de 5 dias que avançam em 1 dia a cada vez. Exemplo: A janela nº 1 abrange do Dia 0 ao Dia 4, a janela nº 2 abrange do Dia 1 ao Dia 5, a janela nº 3 abrange do Dia 2 ao Dia 6, e assim por diante. Cada janela inclui eventos de 00:00:00.00 no dia de início até (mas sem incluir) 00:00:00.00 no dia de término. Como as janelas se sobrepõem, um único evento pode pertencer a várias janelas (neste exemplo, cada evento pertence a até 5 janelas diferentes).

Janela de tempo de vida​

Defina window_duration=None para criar uma janela de tempo de vida. Em cada limite de deslizamento, o recurso agrega todos os eventos de origem para a entidade com Timestamp anteriores a esse limite. Por exemplo, um deslizamento de um dia produz um valor cumulativo uma vez por dia.

Janelas de tempo de vida são suportadas apenas por SlidingWindow. RollingWindow e TumblingWindow requerem um window_duration finito.

nota

As janelas de tempo de vida exigem uma versão do cliente databricks-feature-engineering que seja compatível com window_duration=None e com a ativação do workspace. Versões anteriores do cliente não são compatíveis com esta sintaxe.

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

Janela em dente de serra​

info

Beta

SawtoothWindow está em Beta.

Uma janela sawtooth é uma agregação que oferece suporte a atualizações altamente recentes para eventos recentes, juntamente com a compactação diária de dados históricos. Sua borda final (mais antiga) avança em etapas diárias fixas, enquanto a borda inicial (recente) permanece atualizada com os eventos mais recentes, de modo que o comprimento efetivo da janela "serra" ao longo de cada dia. A maior parte da janela é atendida a partir de dados na tabela de ingestão da transmissão, e apenas os dois dias mais recentes vêm da transmissão ao vivo. Este é um compromisso que realiza o compute de forma eficiente de janelas de longa duração (escalando para anos) enquanto permanece responsivo a atualizações recentes.

Uma janela sawtooth: a borda inicial rastreia os eventos mais recentes, enquanto a borda final avança em passos diários, de modo que a janela coberta &quot;serra&quot; ao longo de cada dia.

As janelas sawtooth são materializadas em um caminho híbrido de lotes e transmissão. Um pipeline de lotes mantém a maior parte da janela, enquanto um pipeline de transmissão mantém os dados mais recentes atualizados em tempo real. Os dois são mesclados (merge) na leitura, portanto, para o modelo ou consumidor de serviço, trata-se de um único recurso.

Como a porção histórica da janela é computada pelo pipeline de lotes, um recurso em dente de serra está pronto para servir pouco depois que a materialização começa, mesmo quando a janela abrange meses ou anos. Uma janela deslizante só é concluída após o decurso de sua duração total de janela. O window_duration mínimo deve ser superior a dois dias (o limite inferior aplicado). A Databricks recomenda uma janela em dente de serra para durações superiores a 7 dias. Para janelas superiores a dois dias e de até sete dias, escolha entre a precisão de comprimento fixo de uma janela deslizante e a prontidão de produção mais rápida de uma janela em dente de serra.

nota

Um recurso de dente de serra depende do histórico já presente. A tabela de ingestão da transmissão deve conter dados que cubram pelo menos a duração total da janela, ou a janela calculada estará incompleta. Antes de se passarem 2 dias completos, o recurso reflete apenas os dados materializados até o momento. Não é recomendado disponibilizar o recurso em produção até que se passem 2 dias completos. Uma agregação em uma janela vazia retorna 0 para Sum e Count, e nulo para Avg, Min, Max, First, Last, VarPop, VarSamp, StddevPop e StddevSamp.

Para saber se um recurso dente de serra está pronto, abra a recurso view no Explorador de Catálogos. Na seção de recursos materializados, o preenchimento retroativo em lote é concluído assim que o tempo da última materialização do recurso avança e seu status indica sucesso. A parte de transmissão é materializada por um pipeline declarativo do Lakeflow. Após a Feature View passar pela validação, o recurso materializado é vinculado a esse pipeline, onde você pode monitorar o status da sua execução.

As janelas sawtooth exigem um(a) StreamSource e são materializadas com StreamingMode.

Python
class SawtoothWindow(TimeWindow):
window_duration: datetime.timedelta

As bordas de uma janela sawtooth movem-se de forma diferente das de uma janela deslizante: a borda inicial rastreia o evento mais recente, enquanto a borda final avança uma vez por dia, em vez de continuamente. Todos os dias, em um corte fixo às 18:00 UTC, a borda final avança para o limite da meia-noite UTC daquele dia. Como resultado, a janela efetiva é ligeiramente maior que window_duration e cresce ao longo do dia antes de retornar um dia no próximo corte. O treinamento e a disponibilização usam o mesmo corte de 18:00 UTC, portanto, o treinamento offline e a disponibilização online permanecem consistentes.

Parâmetro

Restrições

window_duration

Deve ser maior que dois dias. Uma duração que não seja um número inteiro de dias (por exemplo, timedelta(days=3, minutes=15)) é permitida, mas a janela ainda é atualizada com granularidade diária.

Parâmetro

Restrições

window_duration

Deve ser maior que dois dias. Uma duração que não seja um número inteiro de dias (por exemplo, timedelta(days=3, minutes=15)) é permitida, mas a janela ainda é atualizada com granularidade diária.

As janelas em dente de serra oferecem suporte às funções de agregação Sum, Avg, Count, Min, Max, First, Last, VarPop, VarSamp, StddevPop e StddevSamp.

Exemplo de janela dente de serra​

O exemplo a seguir mostra uma contagem de 7 dias das transações de um usuário. A borda inicial rastreia o evento atual, enquanto a borda final avança um dia de cada vez. Para eventos em 10 de março, a janela retrocede até cerca de 3 de março. À medida que 10 de março avança, a borda inicial continua avançando enquanto a borda final permanece, de modo que o intervalo coberto aumenta. Então, no início de 11 de março, a borda final avança para cerca de 4 de março. A janela efetiva é sempre um pouco maior do que sete dias. Os dois dias mais recentes são servidos a partir da transmissão ao vivo, e os dias anteriores são servidos a partir da tabela de ingestão da transmissão.

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

Limitações da janela dente de serra​

  • O parâmetro delay não é compatível.
  • SourceLateness.settling_delay não é suportado.
  • Funções de agregação diferentes de Sum, Avg, Count, Min, Max, First, Last, VarPop, VarSamp, StddevPop e StddevSamp não são suportadas (por exemplo, ApproxCountDistinct, ApproxPercentile, FirstN, LastN, FirstDistinct e LastDistinct).
  • Janelas dente de serra requerem um StreamSource. Um DeltaTableSource não é compatível.

Trigger de materialização​

Triggers controlam quando um pipeline de materialização é executado. O tipo de trigger depende do tipo de recurso.

CronSchedule​

Use CronSchedule para recursos de agregação em lotes. By default, Databricks derives a schedule from the aggregation window. Um programar derivado contabiliza o período da janela, a janela delay e offset, e a origem settling_delay para que uma execução não publique uma janela antes de se esperar que seus dados de origem estejam completos. Programações derivadas dão suporte a janelas fixas e deslizantes.

To request a derived schedule, omit the cron expression. CronSchedule() and the explicit CronSchedule(mode=CronScheduleMode.DERIVED) form are equivalent:

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

trigger = CronSchedule(mode=CronScheduleMode.DERIVED)

Do not set quartz_cron_expression with CronScheduleMode.DERIVED. When you retrieve the materialized recurso, the returned schedule can contain the cron expression that Databricks computed.

Para controlar o programar diretamente, forneça uma expressão cron do Quartz. CronScheduleMode.MANUAL é inferido quando você fornece uma expressão:

Python
from databricks.feature_engineering.entities import CronSchedule

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

TableTrigger​

Use TableTrigger para recursos ColumnSelection ou recursos de agregação (AggregationFunction) suportados por um DeltaTableSource. O pipeline é executado sempre que a tabela Delta upstream recebe um novo commit.

Para recursos de agregação, o pipeline é limitado para que não seja executado a cada commit. O pipeline é executado no máximo uma vez por metade do comprimento da janela do recurso, mas nunca com frequência superior a cada 5 minutos. Por exemplo, um recurso com uma janela de tombamento de 1 hora é executado no máximo uma vez a cada 30 minutos, ou um recurso com uma janela de 8 horas é executado no máximo uma vez a cada 4 horas. O limite mínimo de 5 minutos se aplica quando metade da janela é menor que isso, portanto, janelas de 10 minutos ou menos são executadas no máximo uma vez a cada 5 minutos. Recursos de agregação cuja janela é inferior a 5 minutos não podem usar TableTrigger; use um trigger de transmissão em vez disso.

Python
from databricks.feature_engineering.entities import TableTrigger

trigger = TableTrigger()

StreamingMode​

Use StreamingMode para recursos apoiados por um StreamSource. O pipeline é executado como um pipeline de transmissão contínua.

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

Escolhendo um Trigger​

Cada recurso usa um trigger; as opções por tipo de recurso são:

Tipo de recurso

Trigger

Quando é executado

Agregação (AggregationFunction) de DeltaTableSource

CronSchedule

Em um programar derivado ou manual

Agregação (AggregationFunction) de DeltaTableSource

TableTrigger

Em cada commit da tabela de origem

ColumnSelection (de DeltaTableSource)

TableTrigger

Em cada commit da tabela de origem

Recursos de StreamSource

StreamingMode

Transmissão contínua

Tipo de recurso

Trigger

Quando é executado

Agregação (AggregationFunction) de DeltaTableSource

CronSchedule

Em um programar derivado ou manual

Agregação (AggregationFunction) de DeltaTableSource

TableTrigger

Em cada commit da tabela de origem

ColumnSelection (de DeltaTableSource)

TableTrigger

Em cada commit da tabela de origem

Recursos de StreamSource

StreamingMode

Transmissão contínua

Não é possível materializar recursos que exigem tipos de trigger diferentes em uma única chamada materialize_features. Faça chamadas separadas em vez disso.