Kommentar
Åtkomst till den här sidan kräver auktorisering. Du kan prova att logga in eller ändra kataloger.
Åtkomst till den här sidan kräver auktorisering. Du kan prova att ändra kataloger.
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_featureför att definiera funktionsobjekt i Unity Catalog som kan användas i modelltränings- och serveringsarbetsflöden. - Du kan också skapa
Featureobjekt lokalt och användaregister_featureför att spara dem i Unity Catalog senare. Lokalt konstruerade funktioner kan användas medcreate_training_setföre registreringen.
- Använd
-
Arbetsflöde för modellträning
- Använd
create_training_setfö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.
- Använd
-
Arbetsflöde för materialisering och servering av funktioner
- När du har definierat en funktion med
create_featureeller hämtat den med kanget_featuredu användamaterialize_featuresfö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_settillsammans med den materialiserade vyn för att förbereda ett offline-batchträningsdataset.
- När du har definierat en funktion med
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
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öroffline_config. - Du kan inte kombinera streaming- och batchfunktioner i ett enda
materialize_featuressamtal. Gör ett separat anrop för varje triggertyp. -
transformation_sqlstö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
StreamSourcenä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 enStreamSource. 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
DeltaTableSourcenä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
DeltaTableSourcenä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
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_featuresanrop 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 olikaentityvä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 definieraFeature. - 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
DATEellerTIMESTAMP. -
RequestSourcestöder endast skalära datatyper som definierats iScalarDataType(INTEGER,FLOAT,BOOLEAN,STRING,DOUBLE,LONG,TIMESTAMP,DATE).SHORTKomplexa typer som matriser, kartor och structs stöds inte. -
RequestSourcestöder inte aggregeringsfunktioner eller tidsfönster. EndastColumnSelectionfunktioner 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_batchkanske 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.