Visualizações de funcionalidades

Importante

Esse recurso está em Visualização Pública. Os administradores do workspace podem controlar o acesso a esse recurso na página Visualizações . Consulte Gerenciar visualizações do Azure Databricks.

As Exibições de Recursos permitem que você defina e compute recursos de fontes de dados. Os recursos podem ser definidos usando uma variedade de fontes (tabela Delta, transmissão do Kafka e dados no momento da solicitação) e computações (agregações com janelas de tempo, seleções de coluna simples etc.). Este guia aborda os seguintes fluxos de trabalho:

  • Fluxo de trabalho de desenvolvimento de recursos
    • Use create_feature para definir objetos de recurso do Catálogo do Unity que podem ser usados no treinamento de modelos e no fornecimento de fluxos de trabalho.
    • Como alternativa, construa Feature objetos localmente e use-os register_feature para persisti-los no Catálogo do Unity posteriormente. Os recursos construídos localmente podem ser usados com create_training_set antes do registro.
  • Treinamento de modelo fluxo de trabalho
    • Use create_training_set para calcular recursos agregados pontuais para machine learning. Para obter uma documentação detalhada sobre o treinamento com exibições de recursos, consulte Treinar modelos com exibições de recursos.
  • Fluxo de trabalho de materialização de recursos e serviço
    • Depois de definir um recurso com create_feature ou recuperá-lo usando get_feature, você pode usar materialize_features para materializar o recurso ou conjunto de recursos em um repositório offline para reutilização eficiente ou para um repositório online para serviço online.
    • Use create_training_set com a exibição materializada para preparar um conjunto de dados de treinamento em lote offline.

Para obter detalhes da API, consulte a referência da API de Exibições de Recursos.

Requisitos

  • Computação sem servidor ou um cluster de computação clássico executando o Databricks Runtime 17.0 ML ou superior.

  • Você deve instalar o pacote do Python personalizado. Execute as seguintes linhas de código sempre que executar um notebook:

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

Exemplo de início rápido

Para um notebook de início rápido executável, consulte Exemplo de notebook.

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

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

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

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

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

fe = FeatureEngineeringClient()

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

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

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

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

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

    # training code

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

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

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

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

Notebook de exemplo

Notebook de início rápido de Visualizações de recursos

Obter laptop

Recursos de streaming

Use recursos de streaming quando os valores dos recursos precisarem ser atualizados continuamente, em vez de em um cronograma em lote. Recursos de transmissão e lote usam os mesmos Feature construtores, funções de agregação e fluxos de trabalho de treinamento e serviço.

Recursos de streaming não se materializam em uma loja offline. Para treinamento e inferência em lote, o Databricks calcula valores de características a partir da fonte.

Os recursos de streaming têm os seguintes requisitos:

  • Você deve fornecer um online_config. Recursos de streaming não suportam offline_config.
  • Você não pode combinar recursos de streaming e de processamento em lote em uma única chamada materialize_features. Faça uma chamada separada para cada tipo de gatilho.
  • transformation_sql não é suportado para recursos de streaming.
  • A materialização em transmissão processa somente os registros que chegam depois do início do pipeline e não faz o preenchimento retroativo dos registros históricos. Agregados de janela rolativa retornam resultados completos somente após a chegada da primeira janela completa de dados.

Escolha de uma fonte de recurso de streaming

Escolha uma fonte baseada nas suas necessidades de frescura e na configuração de ingestão existente:

  • Use um StreamSource quando a atualização em menos de um segundo for a prioridade. StreamSource Os recursos proporcionam uma latência ponta a ponta P99 de 200 milissegundos. Primeiro configure um Stream; depois, referencie-o usando um StreamSource. As fontes de fluxo suportam Kafka como entrada e mantêm automaticamente uma tabela Delta de ingestão como cópia histórica dos dados para treinamento.
  • Use um DeltaTableSource quando você já tiver um caminho de ingestão de baixa latência para uma tabela Delta. Espere uma atualização na ordem de dezenas de segundos.
  • Use Zerobus para preencher uma DeltaTableSource quando você ainda não tiver um caminho de ingestão de baixa latência. A ingestão do Zerobus leva na ordem de dezenas de segundos, então espere uma atualização dos atributos em menos de um minuto.

Defina um recurso de streaming usando um StreamSource

Um(a) StreamSource faz referência a um Stream por seu nome de três partes (catalog.schema.stream_name). Um Stream não é um objeto protegível do Catálogo do Unity, mas tem como escopo um esquema do Catálogo do Unity e o acesso é regido pela tabela de ingestão do Stream. As referências de coluna em definições de entidade, de timeseries e de função devem ser prefixadas com value. ou key. para indicar qual parte da mensagem Kafka deve ser lida. Há suporte para campos aninhados usando notação de ponto (por exemplo, value.user.address.city).

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

client = FeatureEngineeringClient()

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

feature = Feature(
    name="user_purchase_sum",
    source=stream_source,
    entity=["value.user_id"],
    timeseries_column="value.event_time",
    function=AggregationFunction(
        operator=Sum(input="value.amount"),
        time_window=RollingWindow(window_duration=timedelta(hours=1)),
    ),
)

Defina um recurso de streaming usando um DeltaTableSource

Para materializar uma característica definida em a DeltaTableSource como uma funcionalidade de streaming, passe StreamingMode como gatilho para materialize_features. A definição de característica usa as mesmas APIs que uma funcionalidade batch respaldada por um DeltaTableSource. As fontes da tabela Delta suportam recursos de agregação e seleção de colunas.

A tabela Delta de origem deve ter o feed de dados de alterações (CDF) habilitado com a configuração delta.enableChangeDataFeed=true.

O exemplo a seguir define e materializa uma característica de agregação com uma fonte de tabela Delta.

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

client = FeatureEngineeringClient()

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

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

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

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

Use uma tabela Delta preenchida pelo Zerobus

Uma tabela Delta preenchida pelo Zerobus pode servir como fonte de recurso de streaming. Zerobus não define delta.enableChangeDataFeed=true automaticamente. Você deve definir essa propriedade manualmente na tabela Delta de destino antes de usá-la como fonte de recurso de streaming.

Condições do filtro em fontes de streaming

Use filter_condition para filtrar linhas antes da agregação para um StreamSource ou DeltaTableSource.

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

Seleção de colunas a partir de fontes de streaming

ColumnSelection Os atributos transmitem o valor mais atual da fonte para cada chave de entidade, sem agregação. Durante o treinamento, os valores dos atributos respeitam a precisão temporal.

Os recursos de seleção de colunas não possuem TTL. Para remover um valor selecionado da loja online, a fonte deve emitir um valor nulo para a coluna selecionada.

from databricks.feature_engineering.entities import ColumnSelection

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

Acesse campos aninhados a partir de um StreamSource

Para um StreamSource, você pode acessar campos JSON aninhados usando notação de pontos (por exemplo, value.nested_field.amount). Em tempo de serviço, o conteúdo e a resposta da solicitação usam nomes de nó folha (por exemplo, amount em vez de value.amount). Os nomes dos nós folha devem ser únicos entre todos os nomes de entidades, séries temporais e recursos em um modelo ou Especificação de Recursos, pois o ponto de extremidade do serviço de disponibilização usa os nomes dos nós folha para rotear valores.

Janelas de tempo para recursos de streaming

Os recursos de streaming oferecem suporte apenas a RollingWindow para agregações. As janelas deslizantes recalculam continuamente sobre os dados mais recentes, o que se alinha à natureza em tempo real das fontes de dados em streaming. TumblingWindow e SlidingWindow são projetados para computação em lote em intervalos históricos fixos.

Notebook de exemplo das funcionalidades de transmissão

Bloco de anotações de início rápido exibições de recursos de transmissão

Obter laptop

Treinamento e inferência de modelo

Para treinar modelos e executar a inferência em lotes com Exibições de Recursos, incluindo log_model(), score_batch()e create_training_set(), consulte Treinar modelos com Exibições de Recursos.

Materialização de recursos

Depois de definir recursos, você pode materializá-los em repositórios offline ou online para reutilização eficiente em fluxos de trabalho de treinamento e serviço. Depois de materializar recursos, você pode fornecer modelos usando o Serviço de Modelo de CPU. Para obter detalhes, consulte Materializar exibições de recursos.

Práticas recomendadas

Nomenclatura de funcionalidades

  • Use nomes descritivos para recursos críticos para os negócios.
  • Siga as convenções de nomenclatura consistentes entre as equipes.
  • Use nomes gerados automaticamente à medida que você começa a desenvolver recursos.

Janelas de horas

  • Alinhar limites de janela com ciclos de negócios (diários, semanais).
  • Janelas mais curtas capturam tendências recentes, mas podem ser barulhentas. Janelas mais longas produzem distribuições de recursos mais estáveis, mas podem perder mudanças comportamentais recentes. Escolha com base na rapidez com que o sinal subjacente muda para seu caso de uso. Por exemplo, uma janela de 7 dias suaviza as flutuações diárias e produz entradas de modelo consistentes, enquanto uma janela de 1 hora reage rapidamente às alterações comportamentais, mas pode introduzir variação que degrada o desempenho do modelo. Se a precisão do modelo diminuir quando a distribuição mudar, use uma janela mais longa para estabilizar as entradas.
  • Janelas em cascata e deslizantes são mais escaláveis do que janelas sem interrupção (contínuas). Comece com janelas deslizantes para a maioria dos casos de uso.

Performance

  • Materialize recursos da mesma fonte de dados em uma única chamada materialize_features para minimizar as verificações de dados.
  • Use a mesma granularidade (por exemplo, todas as durações de slide de 1 hora ou 1 dia) para recursos na mesma fonte de dados para habilitar um melhor agrupamento durante a materialização.

Colunas de entidade versus condições de filtro

Use este guia de decisão ao trabalhar com recursos da mesma tabela de origem:

Use entity (em create_feature) quando precisar de diferentes níveis de agregação:

  • Recursos no nível do cliente (uma linha por cliente): entity=["customer_id"]
  • Recursos de cliente-comerciante (várias linhas por cliente): entity=["customer_id", "merchant_id"]
  • Diferentes níveis de agregação podem compartilhar o mesmo DeltaTableSource: especifique valores diferentes entity em cada definição de recurso

Use filter_condition (em DeltaTableSource) quando precisar filtrar linhas no mesmo nível de agregação:

  • Somente transações de alto valor: filter_condition="amount > 100" (ainda agregadas por cliente)
  • Somente pedidos concluídos: filter_condition="status = 'completed'" (ainda agregados por cliente)

Regra geral: Se a alteração resultar em um número diferente de linhas por valor de entidade, use valores diferentes entity em suas definições de recurso. Se você estiver apenas filtrando quais linhas contribuem para a mesma agregação, use filter_condition na origem.

Padrões comuns

Análise de clientes

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

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

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

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

Análise de tendência

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

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

Padrões sazonais

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

Limitações

  • Os nomes das colunas de entidade e de séries temporais devem corresponder entre o conjunto de dados de treinamento (rotulado) e as definições de recursos quando usados na API create_training_set.
  • O nome da coluna usado como a coluna label no conjunto de dados de treinamento não deve existir nas tabelas de origem, usadas para definir Features.
  • Há suporte para uma lista limitada de funções (UDAFs) na create_feature API. Consulte as funções com suporte.
  • As colunas de entidade não podem ser do tipo DATE ou TIMESTAMP.
  • RequestSource dá suporte apenas a tipos de dados escalares definidos em ScalarDataType (INTEGER, FLOAT, BOOLEAN, STRING, DOUBLE, LONG, TIMESTAMP, DATE, SHORT). Não há suporte para tipos complexos, como matrizes, mapas e structs.
  • RequestSource não dá suporte a funções de agregação ou janelas de tempo. Somente ColumnSelection funções podem ser usadas.
  • O conjunto de nomes de colunas de entidade, nomes de colunas de séries temporais e nomes de colunas de recursos de solicitação deve ser globalmente exclusivo em todas as fontes de um conjunto de treinamento ou ponto de extremidade de serviço.
  • score_batch talvez não tenha êxito na computação sem servidor. Contorne isso usando um cluster de computação clássica executando o Databricks Runtime 17.0 ML ou posterior.

Para limitações específicas de materialização, consulte Limitações.