API-referens för Feature Views

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.

Åtkomstkontroll

Funktioner är styrningsbara Unity Catalog-objekt. Åtkomst till en funktion styrs av behörigheterna CREATE FEATURE, READ FEATUREoch MANAGE Unity Catalog. Fullständiga beskrivningar finns i Referens för behörigheter för Unity Catalog.

  • CREATE FEATURE – Krävs för att skapa en funktion i ett schema. create_feature och register_feature kräver CREATE FEATURE på det överordnade schemat. Efter principen om minsta behörighet beviljar CREATE FEATURE du på schemanivå. Du kan också bevilja den i en katalog så att du kan skapa funktioner i alla scheman i katalogen.
  • READ FEATURE – Krävs för att läsa en funktion och dess data. get_feature, create_training_setoch läsa materialiserade funktionsdata för träning eller servering kräver READ FEATURE på funktionen. READ FEATURE som beviljats för ett schema eller en katalog gäller för alla aktuella och framtida funktioner som den innehåller.
  • MANAGE – Krävs för att hantera en funktions livscykel och bidrag. Att ta bort en funktion med delete_featureoch materialisera en funktion med materialize_features eller delete_materialized_feature, kräver MANAGE på funktionen.

Alla funktionsåtgärder kräver USE CATALOG också i den överordnade katalogen och USE SCHEMA i det överordnade schemat. Information om hur MANAGE och READ FEATURE gäller för materialisering finns i Behörigheter.

API för funktionsvy

Feature konstruktor och register_feature()

Den rekommenderade metoden är att konstruera ett Feature objekt lokalt och använda register_feature för att spara det i Unity Catalog. Med det här tvåstegsarbetsflödet kan du experimentera med funktioner (inklusive create_training_set) innan du registrerar dem.

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() registrerar en lokalt konstruerad Feature i Unity Catalog.

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
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() validerar, skapar och registrerar omedelbart en funktion i Unity Catalog i ett enda steg. Använd detta när du inte behöver experimentera med funktionen lokalt först.

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

Parameters:

  • source: Datakällan som används i funktionsberäkning (DeltaTableSource, StreamSource, eller RequestSource).
  • function: En AggregationFunction som buntar ihop operatorn (till exempel Sum(input="amount")), indatakolumnen och tidsfönstret. Eller ColumnSelection("column_name") för direktfunktioner.
  • catalog_name: Unity Catalog-katalognamnet för funktionen.
  • schema_name: Unity Catalog-schemanamnet för funktionen.
  • entity: Lista över kolumnnamn som definierar aggregerings- eller uppslagsnycklarna (primära nycklar). Krävs för alla källtyper utom RequestSource. Till exempel ["user_id"] aggregerar eller söker efter per användare.
  • timeseries_column: Tidsstämpelkolumnen som används för aggregering av tidsfönster eller val av senaste värde. Krävs för alla källtyper utom RequestSource.
  • name: Valfritt funktionsnamn. Om det utelämnas genereras automatiskt från indatakolumnen, funktionen och fönstret (till exempel amount_avg_rolling_7d).
  • description: Valfri beskrivning av funktionen.

Returnerar: En validerad funktionsinstans

Höjer: ValueError om verifieringen misslyckas

delete_feature()

Tar bort en funktion från Unity Catalog med dess fullständigt kvalificerade namn.

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

Innan du tar bort en funktion tar du bort eller uppdaterar modeller eller funktionsspecifikationer som refererar till den. Om funktionen har materialiserats tar du först bort den materialiserade funktionen. Se Så här tar du bort en materialiserad funktion.

Namn som genereras automatiskt

När name utelämnas genereras ett namn automatiskt. Genererade namn följer mönstret: {column}_{function}_{window}. Ett exempel:

  • price_avg_rolling_1h (Genomsnittspris på 1 timme)
  • transaction_count_rolling_30d_1d (30 dagars antal transaktioner med 1d fördröjning från händelsetidsstämpeln)

Funktioner som stöds

Sammansättningsfunktioner

Note

Sammansättningsfunktioner omsluts i ett AggregationFunction tillsammans med ett tidsfönster, enligt beskrivningen i tidsfönster. Varje funktion tar en input parameter som anger källkolumnen som ska aggregeras.

Function Description Exempel på användningsfall
Sum(input="column") Totalt antal värden Per användares dagliga appanvändning i minuter
Avg(input="column") Medelvärde av värden Genomsnittligt transaktionsbelopp
Count(input="column") Antal poster Antal inloggningar per användare
Min(input="column") Minsta värde Lägsta puls som registrerats av en bärbar enhet
Max(input="column") Maxvärde Högsta transaktionsbelopp per session
StddevPop(input="column") Populationens standardavvikelse Variabilitet för dagliga transaktioner för alla kunder
StddevSamp(input="column") Exempel på standardavvikelse Variabilitet för klickfrekvenser för annonskampanj
VarPop(input="column") Populationsvarians Spridning av sensoravläsningar för IoT-enheter i en fabrik
VarSamp(input="column") Exempelavvikelse Spridning av filmklassificeringar över en samplad grupp
ApproxCountDistinct(input="column", relativeSD=0.05) Ungefärligt unikt antal Distinkt antal köpta artiklar
ApproxPercentile(input="column", percentile=0.95, accuracy=100) Ungefärlig percentil svarstid för p95
First(input="column") Första värdet Tidsstämpel för första inloggning
Last(input="column") Senaste värde Senaste inköpsbelopp
FirstN(input="column", n=3) Första n värdena som en array De tre första produkterna som ses i en session
LastN(input="column", n=3) Sista n värden som en array De tre senaste stödfallsstatusarna
FirstDistinct(input="column", n=3) Först n distinkta värden som en array De tre första distinkta produktkategorierna som ses
LastDistinct(input="column", n=3) Sista n distinkta värden som en array De tre senaste distinkta handelskategorierna

Note

First, Last, FirstN, LastN, , FirstDistinct, och LastDistinct inkluderar nollvärden som standard. Om du vill hoppa över null-värden lägger du till en filter_condition som uttryckligen exkluderar indatakolumner som är null.

FirstN, LastN, , och LastDistinct använder funktionerna timeseries_column för att ordna inmatningsrader och returnera en array med upp till n värdenFirstDistinct. Parametern n måste vara ett positivt heltal. FirstN och FirstDistinct väljer värden från tidigast till senast. LastN och LastDistinct väljer värden från senaste till tidigaste, och returnerar sedan de valda värdena i tidsstämpelordning. FirstDistinct och LastDistinct ta bort dubbletter när man väljer värden i den riktningen.

Till exempel, om källraderna för en entitet är ordnade efter event_time som ["A", "A", "B", "C", "B", "B"], returnerar följande funktioner:

Function Result
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, , FirstDistinctoch LastDistinct kräver databricks-feature-engineering version 0.17.0 eller senare.

ColumnSelection (genomströmning)

ColumnSelection väljer en enskild kolumn från en källa utan att tillämpa någon aggregering. Den omsluts direkt i parametern function (inte inuti AggregationFunction). Returtypen härleds från källschemat.

Function Description Exempel på användningsfall
ColumnSelection("col") Senaste värdet för en kolumn (ingen sammansättning) Den senaste leverantörskategorin, genomströmning av ett begärandefält

ColumnSelection kan användas med valfri datakälla:

  • DeltaTableSource: Returnerar det senaste värdet per entitetsnyckel via en punkt-i-tid-koppling (ingen lookback-fönsteraggregering).
  • StreamSource: Returnerar det senaste värdet per entitetsnyckel från strömmen (ingen sammansättning av återblicksfönstret).
  • RequestSource: Passerar genom värdet som anges vid inferens (eller extraheras från den märkta DataFrame vid träningstillfället).
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",
)

Exempel: sammansättnings- och kolumnvalsfunktioner

I följande exempel visas funktioner som definierats över samma datakälla.

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

Funktioner med filtervillkor

Med filter_condition parametern kan du filtrera rader från källtabellen innan du beräknar sammansättningar. Det här fungerar som en SQL-sats WHERE som tillämpas innan du grupperar och aggregerar data.

Note

filter_condition filtrerar rader före aggregering, till exempel en SQL-sats WHERE som tillämpades före GROUP BY. Den ändrar inte kornigheten, som alltid definieras av i funktionsdefinitionen entity .

Filter är användbara när du arbetar med stora källtabeller som innehåller en supermängd data som behövs för funktionsberäkning och minimerar behovet av att skapa separata vyer ovanpå dessa tabeller.

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

Datakällor

DeltaTableSource

DeltaTableSource är ett tillfälliga Python objekt som används för att definiera hur funktioner beräknas från en källtabell. Den skapar ingen ny tabell. Den anger konfigurationen för att läsa data och aggregeringsfunktioner.

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
)

Parameters:

  • catalog_name, schema_name, table_name: Identifiera deltatabellen för källan i Unity Catalog.
  • filter_condition: En SQL-sats WHERE som tillämpades före aggregering. Exempel: "status = 'completed'".
  • transformation_sql: Ett SQL-uttryck SELECT som tillämpas på källtabellen. Använd det här alternativet om du vill byta namn på kolumner, gjutna typer eller beräknings härledda kolumner före aggregering. Om det utelämnas markeras alla kolumner (*). Exempel: "user_id, CAST(amount AS DOUBLE) AS amount, event_time".
  • dataframe_schema: Schemat för resulterande DataFrame efter transformeringar, i Spark StructType JSON-format (från df.schema.json()). Krävs om transformation_sql anges. Detta talar om för systemet vilka kolumnnamn och typer som är resultatet av omvandlingen.

När både filter_condition och transformation_sql anges är den resulterande frågan: SELECT {transformation_sql} FROM {table} WHERE {filter_condition}.

Note

( timeseries_column som anges i funktionsdefinitionen, inte på DeltaTableSource) måste vara av typen TimestampType eller DateType. Heltalstyper kan fungera men orsaka förlust i precision för tidsfönsteraggregat.

Exempel: Använda transformation_sql för kolumntransformeringar

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

Exempel: Härleda transformation_sql och dataframe_schema från en PySpark DataFrame

Du kan skriva omvandlingen som en PySpark-fråga och sedan extrahera schemat från den resulterande DataFrame:

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

Note

transformation_sql stöder endast radvisa uttryck (kolumnbyten, gjutningar, aritmetik). Sammansättningsfunktioner som COUNT(*) eller SUM() stöds inte. Använd AggregationFunction i funktionsdefinitionen i stället.

DeltaTableSource.from_sql()

Som en bekvämlighet kan du skapa en DeltaTableSource från en SQL-fråga. Metoden parsar frågan för att automatiskt extrahera tabellnamnet, transformation_sql, och filter_condition.

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

Endast enkla SELECT ... FROM ... [WHERE ...] frågor stöds. Komplex SQL (JOIN, underfrågor, CTE: er, UNION) avvisas. För komplexa frågor skapar du DeltaTableSource direkt med transformation_sql och filter_condition.

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

Iterera med to_dataframe()

Använd source.to_dataframe() för att förhandsgranska de data som ska användas för funktionsberäkning. Detta är användbart när du itererar på filter_condition och transformation_sql tills de ger de förväntade resultaten.

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

Förstå entiteter

Entitetskolumner definierar aggregeringsnivån för dina funktioner. De anges i Feature definitionen, inte på DeltaTableSource. Entiteter avgör:

  • Gruppera data: Funktioner aggregeras per unik kombination av entitetsvärden (ungefär GROUP BY som i SQL)
  • Den primära nyckelstrukturen: Varje unik entitetskombination resulterar i en rad med beräknade funktioner

Exempel: Funktioner på kundnivå

Följande kod aggregerar funktioner på kundnivå (en rad per kund):

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

Exempel: Funktioner på kundnivå

Om du vill aggregera funktioner på en mer detaljerad nivå (en rad per kundarkivkombination) använder du flera entitetskolumner:

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

När du behöver funktioner på olika aggregeringsnivåer (till exempel på kundnivå och kundlagringsnivå) använder du olika entity värden i dina funktionsdefinitioner. Samma DeltaTableSource sak kan delas mellan funktioner med olika entitetskonfigurationer.

StreamSource

StreamSource refererar till en Stream. Dataströmmen innehåller konfiguration av anslutning, autentisering, schema och inmatning för strömningskällan. För Kafka måste kolumnreferenser i funktionsdefinitioner föregås av eller value. ange vilken del av meddelandet som ska läsaskey..

StreamSource(
    full_name: str,                       # Required: Three-part Stream name (catalog.schema.stream)
    filter_condition: Optional[str],      # Optional: SQL WHERE clause applied before aggregation
)

Parameters:

  • full_name: Det fullständiga tredelade namnet på en stream (till exempel "my_catalog.my_schema.my_stream").
  • filter_condition (valfritt): En SQL-sats WHERE som används för att strömma data före aggregering med hjälp av kolumnreferenser med punktprefix (till exempel "value.event_type = 'purchase'").
from databricks.feature_engineering.entities import StreamSource

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

RequestSource

RequestSource definierar ett schema för data som tillhandahålls vid inferens i nyttolasten för begäran i stället för att sökas upp från en fördefinierad tabell. Under träningen extraheras dessa kolumner från den märkta DataFrame som skickas till create_training_set. Under modellservern måste anroparen inkludera dem i nyttolasten för HTTP-begäran.

RequestSource används med ColumnSelection (för att överföra ett värde direkt). Den stöder inte aggregeringsfunktioner eller tidsfönster.

Definiera schemat

Definiera schemat som en lista över FieldDefinition objekt, där var och en anger ett kolumnnamn och en ScalarDataType:

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

Datatyper som stöds

RequestSourcestöder skalära typer som definierats i ScalarDataType: INTEGER, FLOAT, BOOLEAN, STRING, DOUBLE, LONG, TIMESTAMP, , DATE. SHORT Komplexa typer som matriser, kartor och structs stöds inte.

Hur begärandedata är hydratiserade

Kontext Behavior
Utbildning (create_training_set) Kolumner extraheras från den märkta DataFrame. Typer verifieras mot det deklarerade schemat. Felmatchningar orsakar ett fel (ingen implicit typkonvertering).
Servering (modellslutpunkt) Kolumner hämtas från dataframe_records eller dataframe_split i HTTP-begäran. JSON-värden skickas till de deklarerade typerna (t.ex. JSON-nummer → DOUBLE).

Modellsignatur

När en modell loggas med en log_model träningsuppsättning som innehåller RequestSource funktioner läggs kolumnerna RequestSource till i MLflow-modellsignaturen som nödvändiga indata. Det innebär att serverdelsslutpunktens API-schema återspeglar vilka fält som anropare måste ange vid inferens.

API för utbildning och slutsatsdragning

create_training_set och score_batch beräkna korrekta funktionsvärden för tidpunkt på begäran från källdata. För funktioner som stöder offlinematerialisering, till exempel glidande fönsteraggregeringar på deltatabellkällor, förbättrar materialisering av funktioner först till ett offlinearkiv prestandan för båda åtgärderna. När materialiserade offlinefunktioner är tillgängliga läser åtgärderna fördefinierade offlinedata i stället för att omberäkna funktionsvärden från källan. Se Materialisera funktionsvyer för att materialisera funktioner till en offlinebutik.

create_training_set()

Skapar en träningsdatauppsättning med rätt funktionsberäkning vid tidpunkt. Mer information finns i Träna modeller med funktionsvyer.

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

Loggar en modell med funktionsmetadata för ursprungsspårning och automatisk funktionssökning under slutsatsdragning. Mer information finns i Träna modeller med funktionsvyer.

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

Utför batchinferens offline med automatisk funktionssökning. Använder funktionsmetadata som lagras med modellen för att beräkna rätt funktioner vid tidpunkt, vilket säkerställer konsekvens med träningen.

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

Indatadataramen måste innehålla kolumnerna entitet och tidsserie som används under träningen. Funktioner beräknas automatiskt från källdata.

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

Tidsfönster

Funktionsvyer stöder tre olika fönstertyper för att styra lookback-beteendet för tidsfönsterbaserade aggregeringar: rullande, rullande och glidande.

  • Rullande fönster ser tillbaka från händelsetiden. Varaktighet och fördröjning definieras uttryckligen.
  • Rullande fönster är fasta, icke-överlappande tidsfönster. Varje datapunkt tillhör exakt ett fönster.
  • Glidande fönster är överlappande, glidande tidsfönster med ett konfigurerbart glidintervall.

Följande bild visar hur de fungerar.

Rullande, rullande och glidande lookback-fönster.

Rullande fönster

Note

RollingWindow hette tidigare ContinuousWindow. Om du migrerar från en tidigare SDK-version uppdaterar du importen i enlighet med detta.

Rullande fönster är up-to- datum- och realtidsaggregeringar, som vanligtvis används över strömmande data. I strömmande pipelines genererar det rullande fönstret endast en ny rad när innehållet i fönstret med fast längd ändras, till exempel när en händelse kommer in eller lämnar. När en rullande fönsterfunktion används i träningspipelines utförs en korrekt funktionsberäkning för tidpunkt på källdata med varaktigheten för fönstret med fast längd omedelbart före en viss händelses tidsstämpel. Detta hjälper till att förhindra online-offline skevhet eller dataläckage. Funktioner vid tidpunkten T aggregerar händelser från [T − varaktighet, T).

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

I följande tabell visas parametrarna för ett rullande fönster. Start- och sluttiderna för fönstret baseras på följande parametrar:

  • Starttid: evaluation_time - window_duration - delay (inklusive)
  • Sluttid: evaluation_time - delay (exklusiv)
Parameter Constraints
delay (valfritt) Måste vara ≥ 0 (flyttar fönstret bakåt i tiden från tidsstämpeln för utvärderingen). Använd delay för att ta hänsyn till eventuella systemfördröjningar mellan den tidpunkt då händelsen skapas och händelsetidsstämpeln för att förhindra framtida händelseläckage i träningsdatauppsättningar. Om det till exempel finns en fördröjning på en minut mellan den tid då händelser skapas och dessa händelser slutligen hamnar i en källtabell där de tilldelas en tidsstämpel, blir timedelta(minutes=1)fördröjningen .
window_duration Måste vara > 0
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))

Definiera ett rullande fönster med fördröjning med hjälp av koden nedan.

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

Exempel på rullande fönster

  • window_duration=timedelta(days=7): Detta skapar ett 7-dagars tillbakablicksfönster som slutar vid den aktuella utvärderingstiden. För ett evenemang kl. 14:00 på dag 7 inkluderar detta alla händelser från 14:00 på dag 0 upp till (men inte inklusive) 14:00 på dag 7.

  • window_duration=timedelta(hours=1), delay=timedelta(minutes=30): Detta skapar ett 1-timmars lookback-fönster som slutar 30 minuter före utvärderingstiden. För ett evenemang kl. 15:00 inkluderar detta alla händelser från 13:30 till (men inte inklusive) 14:30. Detta är användbart för att ta hänsyn till datainmatningsfördröjningar.

Rullande fönster

För funktioner som definieras med fallande fönster beräknas aggregeringar över ett förutbestämt fönster med fast längd som avancerar med ett förskjutningsintervall, vilket ger icke-överlappande fönster som delar upp tiden fullt ut. Därför bidrar varje händelse i källan till exakt ett fönster. Vid tidpunkten t aggregeras data från fönster som avslutas vid eller före t (exklusivt). Windows börjar vid Unix-epoken.

class TumblingWindow(TimeWindow):
    window_duration: datetime.timedelta

I följande tabell visas parametrarna för ett rullande fönster.

Parameter Constraints
window_duration Måste vara > 0
from databricks.feature_engineering.entities import TumblingWindow
from datetime import timedelta

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

Exempel på rullande fönster

  • window_duration=timedelta(days=5): Detta skapar förutbestämda fönster med fast längd på 5 dagar vardera. Exempel: Fönster nr 1 sträcker sig över dag 0 till dag 4, fönster 2 sträcker sig över dag 5 till dag 9, fönster 3 sträcker sig över dag 10 till dag 14 och så vidare. Mer specifikt innehåller Fönster nr 1 alla händelser med tidsstämplar som börjar på 00:00:00.00 dag 0 upp till (men inkluderar inte) några händelser med tidsstämpel 00:00:00.00 på dag 5. Varje händelse tillhör exakt ett fönster.

Skjutfönster

För funktioner som definieras med glidande fönster beräknas aggregeringar över ett förutbestämt fönster med fast längd som avancerar med ett glidintervall, vilket ger överlappande fönster. Varje händelse i källan kan bidra till funktionsaggregering för flera fönster. Vid tidpunkten t aggregeras data från fönster som avslutas vid eller före t (exklusivt). Windows börjar vid Unix-epoken.

class SlidingWindow(TimeWindow):
    window_duration: datetime.timedelta
    slide_duration: datetime.timedelta

I följande tabell visas parametrarna för ett skjutfönster.

Parameter Constraints
window_duration Måste vara > 0
slide_duration Måste vara > 0 och <window_duration
from databricks.feature_engineering.entities import SlidingWindow
from datetime import timedelta

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

Exempel på skjutfönster

  • window_duration=timedelta(days=5), slide_duration=timedelta(days=1): Detta skapar överlappande 5-dagars fönster som avancerar med 1 dag varje gång. Exempel: Fönster nr 1 sträcker sig över dag 0 till dag 4, fönster 2 sträcker sig över dag 1 till dag 5, fönster 3 sträcker sig över dag 2 till dag 6 och så vidare. Varje fönster innehåller händelser från 00:00:00.00 startdagen upp till (men inte inklusive) 00:00:00.00 på slutdagen. Eftersom windows överlappar varandra kan en enskild händelse tillhöra flera fönster (i det här exemplet tillhör varje händelse upp till 5 olika fönster).

Materialiseringsutlösare

Utlöser kontroll när en materialiseringspipeline körs. Utlösartypen beror på funktionstypen.

CronSchedule

Används CronSchedule för aggregeringsfunktioner (AggregationFunction). Pipelinen körs enligt ett fast schema som definieras av ett Quartz cron-uttryck.

from databricks.feature_engineering.entities import CronSchedule

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

TableTrigger

Användning TableTrigger för ColumnSelection funktioner eller aggregeringsfunktioner (AggregationFunction) som backas upp av en DeltaTableSource. Pipelinen körs när den överordnade Delta-tabellen tar emot en ny incheckning.

För aggregeringsfunktioner är pipelinen strypt så att den inte körs på varje commit. Pipelinen körs högst en gång per halva långfilmens fönsterlängd, men aldrig oftare än var femte minut. Till exempel körs en film med ett tumbling-fönster på 1 timme högst en gång var 30:e minut, eller en film med 8-timmarsfönster som mest en gång var fjärde timme. 5-minutersgolvet gäller när halva fönstret är mindre än så, så fönster som är 10 minuter eller mindre löper högst var 5:e minut. Aggregeringsfunktioner vars fönster är under 5 minuter kan inte användas TableTrigger, använd istället en strömningstrigger.

from databricks.feature_engineering.entities import TableTrigger

trigger = TableTrigger()

StreamingMode

Använd StreamingMode för funktioner som backas upp av en StreamSource. Pipelinen körs som en pipeline för kontinuerlig direktuppspelning.

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

Välja en utlösare

Varje funktion använder en trigger; Alternativen efter funktionstyp är:

Funktionstyp Trigger När den körs
Sammansättning (AggregationFunction) från DeltaTableSource CronSchedule Enligt ett fast cron-schema
Sammansättning (AggregationFunction) från DeltaTableSource TableTrigger I varje källtabellincheckning
ColumnSelection (från DeltaTableSource) TableTrigger I varje källtabellincheckning
Funktioner från StreamSource StreamingMode Kontinuerlig direktuppspelning

Du kan inte materialisera funktioner som kräver olika utlösartyper i ett enda materialize_features anrop. Utfärda separata anrop i stället.

Migrera betafunktioner till offentlig förhandsversion

Den offentliga förhandsversionen av funktionsvyer introducerar förstklassiga funktionsentiteter i Unity Catalog, som styrs av CREATE FEATURE behörigheterna och READ FEATURE och kräver databricks-feature-engineering version 0.16.0 eller senare. Funktioner som skapas under betaversionen (med version 0.15.0) lagras som Unity Catalog-funktioner och stöder inte alla funktioner för offentlig förhandsversion. Återskapa dina betafunktioner med version 0.16.0 för att få långsiktigt stöd för offentlig förhandsversion. Funktioner måste tas bort och återskapas, inte bara materialiseras på nytt.

Mer information om funktioner finns i Funktionsvyer.

Vad du behöver göra

  • Uppgradera till 0.16.0. Det här är den klientversion som krävs för funktioner för offentlig förhandsversion (batch och strömning).
  • Återskapa dina funktioner. Betafunktionsvyer måste tas bort och återskapas, inte materialiseras på nytt, eftersom de inte stöder alla funktioner för offentlig förhandsversion.
  • Migrera innan fönstret stängs. Befintliga betafunktioner måste migreras före den 22 juli 2026.

Identifiera funktioner för betaversion och offentlig förhandsversion

Funktioner för offentlig förhandsversion visas som ett funktionsobjekt i Unity Catalog, till exempel i Katalogutforskaren. Betafunktioner visas som en funktion med en YAML-definition. Alla funktioner som representeras som en funktion är en betafunktion som du behöver migrera.

Migrera betafunktioner

Migrering av en betafunktion har tre delar:

  • Återskapa funktionen som en funktion för offentlig förhandsversion.
  • Materialisera funktionen igen, så att dess offline- och onlinetabeller återskapas under den nya funktionen.
  • När du har verifierat de migrerade funktionerna tar du bort betafunktionerna och deras materialiseringar.

Återskapa funktionerna

Använd list_beta_feature_views för att hitta dina betafunktioner, Feature.clone() för att skapa en oregistrerad kopia och register_feature för att registrera varje kopia som en offentlig förhandsversionsfunktion. Kloning rensar registreringen, katalogen och schemat så att funktionen kan registreras igen.

För att undvika namnkollisioner registrerar du migrerade funktioner med ett annat namn eller i ett annat schema än betafunktionerna. I följande exempel registreras varje funktion i det ursprungliga schemat med ett _migrated namnsuffix.

# Update this to the catalog whose beta Feature Views you want to migrate.
CATALOG_TO_MIGRATE = "main"

from databricks.feature_engineering import FeatureEngineeringClient

fe = FeatureEngineeringClient()

# 1. Find every beta Feature View in the catalog. Returns Feature objects,
#    scanned across all schemas in the catalog.
beta_features = fe.list_beta_feature_views(catalog_name=CATALOG_TO_MIGRATE)

# Keep each beta feature paired with its migrated counterpart for the next steps.
migrations = []
for beta_feature in beta_features:
    catalog_name, schema_name, leaf_name = beta_feature.full_name.split(".")
    # 2. Clone the feature as an unregistered copy, renamed with a "_migrated" suffix.
    cloned = beta_feature.clone(new_name=f"{leaf_name}_migrated")
    # 3. Re-register the clone as a Public Preview feature.
    migrated = fe.register_feature(
        feature=cloned,
        catalog_name=catalog_name,
        schema_name=schema_name,
    )
    migrations.append((beta_feature, migrated))

Materialisera om de migrerade funktionerna

Om en betafunktion materialiserades materialiseras dess offentliga förhandsversionsmotsvarighet på nytt så att dess offline- och onlinetabeller återskapas under den nya funktionen. Ange offline- och onlinebutikskonfigurationerna för den migrerade funktionen och rekonstruera utlösaren från betafunktionens befintliga materialisering.

from databricks.feature_engineering.entities import (
    CronSchedule,
    OfflineStoreConfig,
    OnlineStoreConfig,
    TableTrigger,
)

for beta_feature, migrated in migrations:
    # Inspect the beta feature's existing materializations to see what to rebuild and
    # to reconstruct the same trigger.
    trigger = None
    needs_offline = needs_online = False
    for mf in fe.list_materialized_features(feature_name=beta_feature.full_name):
        needs_online = needs_online or bool(mf.is_online)
        needs_offline = needs_offline or not mf.is_online
        # Rebuild the trigger from the materialized feature.
        if mf.cron_schedule_trigger is not None:
            trigger = CronSchedule(
                quartz_cron_expression=mf.cron_schedule_trigger.cron_expression,
                timezone_id="UTC",  # Materialized schedules run in UTC.
            )
        elif mf.table_trigger is not None:
            trigger = TableTrigger()
        elif mf.streaming_mode is not None:
            # Streaming features use StreamingMode, which can be reused as-is.
            trigger = mf.streaming_mode
    if not (needs_offline or needs_online):
        continue  # The beta feature was never materialized.

    catalog_name, schema_name, _ = migrated.full_name.split(".")
    fe.materialize_features(
        features=[migrated],
        offline_config=OfflineStoreConfig(
            catalog_name=catalog_name,
            schema_name=schema_name,
            table_name_prefix="migrated_features",
        )
        if needs_offline
        else None,
        online_config=OnlineStoreConfig(
            catalog_name=catalog_name,
            schema_name=schema_name,
            table_name_prefix="migrated_features",
            online_store_name="my_online_store",
        )
        if needs_online
        else None,
        trigger=trigger,
    )

Note

Om varje funktion materialiseras i sitt eget materialize_features anrop skapas en separat pipeline. För att minska beräkningskostnaden grupperar du funktioner som delar ett offline- och onlinemål och utlöser till ett enda materialize_features anrop genom att skicka dem tillsammans i features.

Ta bort betafunktionerna

Varning

Ta bort betafunktioner och deras materialiseringar först när du har kontrollerat att de migrerade funktionerna och deras materialiserade data är korrekta. Borttagningen går inte att ångra.

När du har verifierat de migrerade funktionerna tar du bort varje betafunktions materialiseringar och sedan själva betafunktionen.

for beta_feature, _ in migrations:
    # Delete the beta feature's materializations first.
    mfs = list(fe.list_materialized_features(feature_name=beta_feature.full_name))
    offline_mfs = [mf for mf in mfs if not mf.is_online]
    if offline_mfs:
        # Aggregation features pair an offline and online table; deleting the offline
        # materialized feature removes its paired online table too.
        for mf in offline_mfs:
            fe.delete_materialized_feature(materialized_feature=mf)
    else:
        # Online-only features (ColumnSelection, streaming) have no offline pair; delete
        # the online materialized feature directly.
        for mf in mfs:
            fe.delete_materialized_feature(materialized_feature=mf)
    # Then delete the beta feature definition.
    fe.delete_feature(full_name=beta_feature.full_name)