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 use CronSchedule 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

Förutom batchfunktioner från Delta-tabeller kan du definiera funktioner från strömmande källor för användningsfall i realtid. Direktuppspelningsfunktioner använder samma funktionsklass som batchfunktioner – samma Feature konstruktorer, samma sammansättningsfunktioner, samma tränings- och serverarbetsflöden – så uppgradering från batch till realtid kräver minimala kodändringar. När strömningsfunktionerna har materialiserats levererar de färskhet från slutpunkt till slutpunkt (p99-svarstid på 200 ms) direkt till din modell som betjänar slutpunkter.

Om du vill använda direktuppspelningsfunktioner konfigurerar du först en Stream och refererar sedan till den med hjälp av en StreamSource. Strömkällor stöder Kafka som indatakälla och upprätthåller automatiskt en inläsningstabell (Delta) som en historisk kopia av data för träning.

Definiera en direktuppspelningsfunktion

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

Filtrera villkor på StreamSource

Använd filter_condition för att filtrera rader från strömmen före aggregering, precis som på DeltaTableSource.

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

Kolumnval från strömmar

ColumnSelection-funktioner fungerar med streamingkällor. Den valda kolumnen visar det senaste värdet från strömmen för varje entitet, med bibehållen korrekthet vid den aktuella tidpunkten.

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

Få åtkomst till kapslade fält

Du kan komma åt kapslade JSON-fält med hjälp av punkt notation (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 ). Lövnodnamn måste vara unika för alla entitets-, tidsserie- och funktionsutdatakolumner i en modell eller funktionsspecifikation, eftersom den betjänande slutpunkten använder lövnamn för att dirigera 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.