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 menor privilégio, conceda CREATE FEATURE no nível do esquema; também é possível concedê-lo em um catálogo para permitir a criação de recursos em qualquer esquema nesse catálogo.
  • READ FEATURE — Necessário para ler um recurso e seus dados. get_feature, create_training_set, e a leitura de dados de recurso materializados para treinamento ou disponibilização exigem READ FEATURE no recurso. 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. Excluir um recurso com delete_feature, e materializar um recurso com materialize_features ou delete_materialized_feature, exige MANAGE no recurso.

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, or RequestSource
function: Union[AggregationFunction, ColumnSelection], # Required: Aggregation or column selection
entity: Optional[List[str]] = None, # Required for all sources except RequestSource: entity columns
timeseries_column: Optional[str] = None, # Required for all sources except RequestSource: timestamp column
name: Optional[str] = None, # Optional: Feature name (auto-generated if omitted)
description: Optional[str] = None, # Optional: Feature description
)

FeatureEngineeringClient.register_feature() 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, or RequestSource
function: Union[AggregationFunction, ColumnSelection], # Required: Aggregation or column selection
catalog_name: str, # Required: The catalog name for the feature
schema_name: str, # Required: The schema name for the feature
entity: Optional[List[str]] = None, # Required for all sources except RequestSource: entity columns
timeseries_column: Optional[str] = None, # Required for all sources except RequestSource: timestamp column
name: Optional[str] = None, # Optional: Feature name (auto-generated if omitted)
description: Optional[str] = None, # Optional: Feature description
) -> Feature

Parâmetros:

  • source: a fonte de dados usada no cálculo de recurso (DeltaTableSource, StreamSource ou RequestSource).
  • function: um AggregationFunction que agrupa o operador (por exemplo, Sum(input="amount")), a coluna de entrada e a janela de tempo. Ou ColumnSelection("column_name") para recursos pass-through.
  • 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 coluna que definem as chaves de agregação ou pesquisa (chaves primárias). Necessário para todos os tipos de origem, exceto RequestSource. Por exemplo, ["user_id"] agrega ou pesquisa por usuário.
  • timeseries_columnA coluna Timestamp usada para agregação por janela de tempo ou seleção do último valor. Obrigatório para todos os tipos de origem, exceto RequestSource.
  • 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 quaisquer modelos ou especificações de recurso que o referenciem. Se o recurso foi materializado, exclua o recurso materializado primeiro. 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.

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 pode ser usado com qualquer fonte 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).
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 filter_condition parâmetro permite filtrar linhas da tabela de origem **antes** de computar agregações. Isso funciona como uma cláusula SQL WHERE que é aplicada antes de agrupar e agregar dados.

nota

filter_condition filtra linhas antes da agregação, como uma cláusula WHERE 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
)

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. Exemplo: "status = 'completed'".
  • transformation_sql: Uma expressão SQL SELECT aplicada à tabela de origem. Use isso para renomear colunas, converter tipos ou calcular colunas derivadas antes da agregação. Se omitido, todas as colunas sã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.

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

nota

O timeseries_column (especificado na definição do recurso, não em DeltaTableSource) deve ser do tipo TimestampType ou DateType. Integer types podem funcionar, mas causam perda de precisão para agregações de janela de tempo.

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
)

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.
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 é usado com ColumnSelection (para passar um valor diretamente). 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.

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 Recurso Views suportam quatro tipos de janela para controlar o comportamento de retrospectiva para agregações baseadas em janelas de tempo. Os tipos de janela disponíveis dependem da origem do recurso:

  • Recursos de fonte de transmissão podem usar janelas deslizantes e dente de serra.

  • Recursos de fonte de lotes podem usar janelas rolling, tumbling e 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.

Janela contínua

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

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 para trás no tempo a partir do timestamp de avaliação). Use delay para contabilizar qualquer atraso do sistema entre o momento em que o evento é criado e o timestamp do evento, a fim de evitar vazamento futuro de eventos para datasets de treinamento. Por exemplo, se houver um atraso de um minuto entre o momento em que os eventos são criados e quando esses eventos são eventualmente inseridos em uma tabela de origem onde lhes é atribuído um timestamp, então o atraso seria timedelta(minutes=1).

window_duration

Deve ser > 0

Parâmetro

Restrições

delay (opcional)

Deve ser ≥ 0 (desloca a janela para trás no tempo a partir do timestamp de avaliação). Use delay para contabilizar qualquer atraso do sistema entre o momento em que o evento é criado e o timestamp do evento, a fim de evitar vazamento futuro de eventos para datasets de treinamento. Por exemplo, se houver um atraso de um minuto entre o momento em que os eventos são criados e quando esses eventos são eventualmente inseridos em uma tabela de origem onde lhes é atribuído um timestamp, então o atraso seria timedelta(minutes=1).

window_duration

Deve ser > 0

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
# Look back 7 days, offset by 1 minute to account for data ingestion delay
window = RollingWindow(
window_duration=timedelta(days=7),
delay=timedelta(minutes=1)
)

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 terminando 30 minutos antes do tempo de avaliação. Para um evento às 15:00, isso inclui todos os eventos das 13:30 até (mas não incluindo) as 14:30. Isto é útil para considerar atrasos na ingestão de dados.

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

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

Parâmetro

Restrições

window_duration

Deve ser > 0

Parâmetro

Restrições

window_duration

Deve ser > 0

Python
from databricks.feature_engineering.entities import TumblingWindow
from datetime import timedelta

window = TumblingWindow(
window_duration=timedelta(days=7)
)

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

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.

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.

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

window = SlidingWindow(
window_duration=timedelta(days=7),
slide_duration=timedelta(days=1)
)

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 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 parte histórica da janela é computada pelo pipeline de lotes, um recurso sawtooth está pronto para ser disponibilizado logo após o início da materialização, mesmo quando a janela abrange meses ou anos. Uma janela deslizante só fica completa após a sua duração total ter decorrido. A window_duration mínima deve ser superior a dois dias (o limite inferior imposto), mas as janelas sawtooth são recomendadas para durações superiores a 7 dias; para janelas mais curtas, use uma janela deslizante.

nota

Um recurso dente de serra depende de um histórico que já está 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 computada estará incompleta. Antes de decorridos 2 dias completos, o recurso reflete apenas os dados materializados até o momento. Não é recomendado disponibilizar o recurso em produção até que 2 dias completos tenham decorrido. Uma agregação sobre uma janela vazia retorna 0 para Sum e Count, e nulo para Avg, Min e Max.

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 sawtooth suportam as funções de agregação Sum, Avg, Count, Min e Max.

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.
  • Funções de agregação diferentes de Sum, Avg, Count, Min e Max não são compatíveis (por exemplo, First, Last, ApproxCountDistinct e as funções de desvio padrão e variância).
  • 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 (AggregationFunction). O pipeline é executado em um agendamento fixo definido por uma expressão cron Quartz.

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 uma programação cron fixa

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 uma programação cron fixa

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.