Funktionsvyer

Important

Den här funktionen finns som allmänt tillgänglig förhandsversion. Arbetsyteadministratörer kan styra åtkomsten till den här funktionen från sidan Förhandsversioner . Se Hantera förhandsversioner av Azure Databricks.

Med funktionsvyer kan du definiera och beräkna funktioner från datakällor. Funktioner kan definieras med hjälp av en mängd olika källor (Delta-tabell, Kafka Stream och begärandetidsdata) och beräkningar (tidsfönsteraggregeringar, enkla kolumnval med mera). Den här guiden beskriver följande arbetsflöden:

  • Arbetsflöde för funktionsutveckling
    • Använd create_feature för att definiera funktionsobjekt i Unity Catalog som kan användas i modelltränings- och serveringsarbetsflöden.
    • Du kan också skapa Feature objekt lokalt och använda register_feature för att spara dem i Unity Catalog senare. Lokalt konstruerade funktioner kan användas med create_training_set före registreringen.
  • Arbetsflöde för modellträning
    • Använd create_training_set för att beräkna vid en viss tidpunkt aggregerade egenskaper för maskininlärning. Detaljerad dokumentation om träning med funktionsvyer finns i Träna modeller med funktionsvyer.
  • Arbetsflöde för materialisering och servering av funktioner
    • När du har definierat en funktion med create_feature eller hämtat den med kan get_featuredu använda materialize_features för att materialisera funktionen eller uppsättningen funktioner till en offlinebutik för effektiv återanvändning eller till en onlinebutik för onlineservering.
    • Använd create_training_set tillsammans med den materialiserade vyn för att förbereda ett offline-batchträningsdataset.

Api-information finns i API-referens för funktionsvyer.

Requirements

  • Serverlös beräkning eller ett klassiskt beräkningskluster som kör Databricks Runtime 17.0 ML eller senare.

  • Du måste installera det anpassade Python-paketet. Kör följande kodrader varje gång du kör en anteckningsbok.

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

Snabbstartsexempel

En körbar snabbstartsanteckningsbok finns i Exempelanteckningsbok.

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

Exempelanteckningsbok

Snabbstartsanteckningsbok för funktionsvyer

Hämta anteckningsbok

Direktuppspelningsfunktioner

Använd strömningsfunktioner när funktionsvärden måste uppdateras kontinuerligt istället för på batchschema. Strömnings- och batchfunktioner använder samma Feature konstruktörer, aggregeringsfunktioner samt tränings- och serveringsarbetsflöden.

Streamingfunktioner materialiseras inte i en offlinebutik. För träning och batchinferens beräknar Databricks funktionsvärden från källan.

Streamingfunktioner har följande krav:

  • Du måste ange online_config. Streamingfunktioner har inte stöd för offline_config.
  • Du kan inte kombinera streaming- och batchfunktioner i ett enda materialize_features samtal. Gör ett separat anrop för varje triggertyp.
  • transformation_sql stöds inte för streamingfunktioner.
  • Materialisering via strömning bearbetar endast poster som kommer in efter att pipelinen har startat och efterfyller inte historiska poster. Aggregeringar med rullande fönster ger fullständiga resultat först efter att det första fullständiga fönstret med data har kommit in.

Val av källa för strömningsfunktioner

Välj en källa baserat på dina färskhetskrav och befintliga intagssystem:

  • Använd en StreamSource när uppdatering inom bråkdelar av en sekund prioriteras. StreamSourcefunktionerna ger en p99-latens från början till slut på 200 millisekunder. Först sätter du upp en ström, sedan refererar du till den med hjälp av en StreamSource. Strömningskällor har stöd för Kafka som indata och upprätthåller automatiskt en Delta-tabell för inmatning som en historisk kopia av data för träning.
  • Använd en DeltaTableSource när du redan har en inmatningsväg med låg latens till en Delta-tabell. Räkna med färskhet på tiotals sekunder.
  • Använd Zerobus för att fylla en DeltaTableSource när du inte redan har ett inmatningsflöde med låg latens. Inläsning i Zerobus tar i storleksordningen tiotals sekunder, så förvänta dig att featuredata är aktuella inom mindre än en minut.

Definiera en strömningsfunktion med hjälp av en StreamSource

En StreamSource refererar till en stream med dess tredelade namn (catalog.schema.stream_name). En dataström är inte ett skyddsbart objekt i Unity Catalog, men det är begränsat till ett Unity Catalog-schema och åtkomsten styrs av Streams inmatningstabell. Kolumnreferenser i entitets-, tidsserie- och funktionsdefinitioner måste prefixas med value. eller key. för att ange vilken del av Kafka-meddelandet som ska läsas. Kapslade fält stöds med punkt notation (till exempel 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)),
    ),
)

Definiera en streamingfunktion med en DeltaTableSource

För att materialisera en funktion definierad på en DeltaTableSource som en strömmande funktion, skicka StreamingMode som trigger till materialize_features. Funktionsdefinitionen använder samma API:er som en batchfunktion som backas upp av en DeltaTableSource. Delta-tabellkällor stödjer aggregerings- och kolumnvalsfunktioner.

Käll-Delta-tabellen måste ha ändringsdataflöde (CDF) aktiverat genom att sätta delta.enableChangeDataFeed=true.

Följande exempel definierar och materialiserar en aggregeringsfunktion med en Delta-tabellkälla.

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

Använd en Delta-tabell fylld av Zerobus

En Delta-tabell fylld av Zerobus kan fungera som en källa för strömningsfunktioner. Zerobus ställer inte in delta.enableChangeDataFeed=true automatiskt. Du måste ställa in den här egenskapen manuellt på mål-Delta-tabellen innan du använder den som en källa för strömningsfunktioner.

Filtervillkor för strömmande källor

Använd filter_condition för att filtrera rader före aggregering för antingen en StreamSource eller DeltaTableSource.

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

Val av kolumner från strömmande källor

ColumnSelection funktioner passerar genom det senaste värdet från källan för varje entitetsnyckel utan aggregering. Vid träning respekterar funktionsvärden punkt-i-tid-noggrannhet.

Kolumnvalsfunktioner har ingen TTL. För att ta bort ett valt värde från onlinebutiken måste källan generera ett nullvärde för den valda kolumnen.

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

Åtkomst till nästlade fält från en StreamSource

För en StreamSource, kan du komma åt nästlade JSON-fält med hjälp av punktnotation (till exempel value.nested_field.amount). Vid serveringstillfället använder begärandenyttolasten och svaret lövnodnamn (till exempel amount i stället value.amountför ). Bladnodnamn måste vara unika för alla entitets-, tidsserie- och funktionsnamn inom en modell eller funktionsspecifikation, eftersom den servande ändpunkten använder bladnamn för att routa värden.

Tidsfönster för strömningsfunktioner

Direktuppspelningsfunktioner stöder endast RollingWindow aggregeringar. Rullande fönster omberäknas kontinuerligt över de senaste data, vilket överensstämmer med realtidstypen för strömningskällor. TumblingWindow och SlidingWindow är utformade för batchberäkning över fasta historiska intervall.

Exempelnotebook för streamingfunktioner

Snabbstartsguide för strömmande funktionsvyer

Hämta anteckningsbok

Modellträning och slutsatsdragning

Information om hur du tränar modeller och kör batchinferens med funktionsvyer, inklusive log_model(), score_batch()och create_training_set(), finns i Träna modeller med funktionsvyer.

Materialisering av funktioner

När du har definierat funktioner kan du materialisera dem till offline- eller onlinebutiker för effektiv återanvändning i tränings- och serveringsarbetsflöden. När du har materialiserat funktioner kan du serva modeller med hjälp av CPU-modellbetjäning. Mer information finns i Materialisera funktionsvyer.

Metodtips

Namngivning av funktioner

  • Använd beskrivande namn för affärskritiska funktioner.
  • Följ konsekventa namngivningskonventioner mellan team.
  • Använd automatiskt genererade namn när du börjar utveckla funktioner.

Tidsfönster

  • Justera fönstergränser med konjunkturcykler (dagligen, varje vecka).
  • Kortare fönster fångar de senaste trenderna men kan vara bullriga. Längre fönster ger stabilare funktionsdistributioner men kan missa de senaste beteendeförskjutningarna. Välj baserat på hur snabbt den underliggande signalen ändras för ditt användningsfall. Ett 7-dagarsfönster jämnar till exempel ut de dagliga fluktuationerna och ger konsekventa modellindata, medan ett 1-timmarsfönster reagerar snabbt på beteendeförändringar men kan medföra varians som försämrar modellens prestanda. Om modellens noggrannhet försämras när fördelningen skiftar använder du ett längre fönster för att stabilisera indata.
  • Rullande och skjutbara fönster är mer skalbara än rullande (kontinuerliga) fönster. Börja med skjutfönster för de flesta användningsfall.

Performance

  • Materialisera funktioner från samma datakälla i ett enda materialize_features anrop för att minimera datagenomsökningar.
  • Använd samma kornighet (till exempel alla 1-timmars eller alla 1-dagars tidsintervaller) för funktioner från samma datakälla för att möjliggöra bättre gruppering vid materialisering.

Entitetskolumner jämfört med filtervillkor

Använd den här beslutsguiden när du arbetar med funktioner från samma källtabell:

Använd entity (på create_feature) när du behöver olika aggregeringsnivåer:

  • Funktioner på kundnivå (en rad per kund): entity=["customer_id"]
  • Kund- och handelsfunktioner (flera rader per kund): entity=["customer_id", "merchant_id"]
  • Olika aggregeringsnivåer kan dela samma DeltaTableSource: ange olika entity värden för varje funktionsdefinition

Använd filter_condition (på DeltaTableSource) när du behöver filtrera rader på samma aggregeringsnivå:

  • Endast transaktioner med högt värde: filter_condition="amount > 100" (fortfarande aggregerade per kund)
  • Endast slutförda beställningar: filter_condition="status = 'completed'" (fortfarande aggregerade per kund)

Tumregel: Om ändringen skulle resultera i ett annat antal rader per entitetsvärde använder du olika entity värden i dina funktionsdefinitioner. Om du bara filtrerar vilka rader som bidrar till samma aggregering, använd filter_condition på källan.

Vanliga mönster

Kundanalys

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

Trendanalys

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

Säsongsmönster

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

Limitations

  • Namn på entitets- och tidsseriekolumner måste matcha mellan träningsdatauppsättningen (märkt) och funktionsdefinitionerna när de används i API:et create_training_set .
  • Kolumnnamnet som används som label-kolumn i träningsdatamängden bör inte finnas i de källtabeller som används för att definiera Feature.
  • En begränsad lista över funktioner (UDAFs) stöds i API:et create_feature . Se Funktioner som stöds.
  • Entitetskolumner får inte vara av typen DATE eller TIMESTAMP.
  • RequestSourcestöder endast skalära datatyper som definierats i ScalarDataType (INTEGER, FLOAT, BOOLEAN, STRING, DOUBLE, LONG, TIMESTAMP, DATE). SHORT Komplexa typer som matriser, kartor och structs stöds inte.
  • RequestSource stöder inte aggregeringsfunktioner eller tidsfönster. Endast ColumnSelection funktioner kan användas.
  • Uppsättningen med entitetskolumnnamn, tidseriekolumnnamn och funktionskolumnnamn för begäranden måste vara globalt unika för alla källor i en träningsuppsättning eller en serverslutpunkt.
  • score_batch kanske inte fungerar med serverlös databehandling. Kringgå detta genom att använda ett klassiskt beräkningskluster som kör Databricks Runtime 17.0 ML eller senare.

Information om materialiseringsspecifika begränsningar finns i Begränsningar.