Hinweis
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, sich anzumelden oder das Verzeichnis zu wechseln.
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, das Verzeichnis zu wechseln.
Important
Dieses Feature befindet sich in der Public Preview. Arbeitsbereichsadministratoren können den Zugriff auf dieses Feature über die Vorschauseite steuern. Siehe Manage Azure Databricks Previews.
Mit Featureansichten können Sie Features aus Datenquellen definieren und berechnen. Features können mithilfe einer Vielzahl von Quellen (Delta-Tabelle, Kafka-Stream und Daten zur Anfragezeit) und Berechnungen (zeitfensterbasierte Aggregationen, einfache Spaltenauswahlen und vieles mehr) definiert werden. In diesem Leitfaden werden die folgenden Workflows behandelt:
-
Workflow für die Featureentwicklung
- Verwenden Sie
create_feature, um Featureobjekte des Unity-Katalogs zu definieren, die in Modellschulungen und beim Servieren von Workflows verwendet werden können. - Alternativ können Sie
Feature-Objekte lokal erstellen undregister_featurespäter im Unity-Katalog speichern. Lokal erstellte Features können vorcreate_training_setder Registrierung verwendet werden.
- Verwenden Sie
-
Modellschulungsworkflow
- Verwenden Sie
create_training_set, um Punkt-in-Zeit aggregierte Merkmale für maschinelles Lernen zu berechnen. Ausführliche Dokumentation zur Schulung mit Featureansichten finden Sie unter "Train models with Feature Views".
- Verwenden Sie
-
Arbeitsablauf zur Feature-Verarbeitung und Bereitstellung
- Nachdem Sie ein Feature mit
create_featuredefiniert oder es mitget_featureabgerufen haben, können Siematerialize_featuresverwenden, um das Feature oder die Gruppe von Features in einem Offline-Store für eine effiziente Wiederverwendung oder einem Online-Store für Onlinedienste zu materialisieren. - Verwenden Sie
create_training_setmit der materialisierten Sicht, um einen Offline-Batch-Trainingsdatensatz vorzubereiten.
- Nachdem Sie ein Feature mit
Details zur API finden Sie in der API-Referenz für Feature Views.
Anforderungen
Serverless-Compute oder ein classic compute-Cluster mit Databricks Runtime 17.0 ML oder höher.
Sie müssen das benutzerdefinierte Python-Paket installieren. Führen Sie bei jedem Ausführen eines Notizbuchs die folgenden Codezeilen aus:
%pip install databricks-feature-engineering>=0.16.0 dbutils.library.restartPython()
Schnellstartbeispiel
Ein ausgeführtes Schnellstartnotizbuch finden Sie unter Beispielnotizbuch.
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(),
)
Beispielnotizbuch
Schnellstart-Notebook für Feature-Views
Streaming-Funktionen
Verwenden Sie Streaming-Funktionen, wenn die Funktionswerte kontinuierlich aktualisiert werden müssen, anstatt nach einem Batch-Schema. Streaming- und Batch-Funktionen verwenden dieselben Feature Konstruktoren, Aggregationsfunktionen sowie Trainings- und Bereitstellungs-Workflows.
Streaming-Funktionen werden nicht in einem Offline-Store bereitgestellt. Für Training und Batch-Inferenz berechnet Databricks Feature-Werte aus der Quelle.
Streaming-Funktionen haben folgende Anforderungen:
- Du musst ein
online_configangeben. Streaming-Funktionen unterstützenoffline_confignicht . - Man kann Streaming- und Batch-Funktionen nicht in einem Anruf
materialize_featureskombinieren. Mach für jeden Auslöser einen separaten Anruf. -
transformation_sqlwird für Streaming-Funktionen nicht unterstützt. - Streaming-Materialisierung verarbeitet nur Datensätze, die nach Beginn der Pipeline eintreffen, und füllt keine historischen Datensätze nach. Rolling-Window-Aggregate liefern vollständige Ergebnisse erst, nachdem das erste vollständige Datenfenster eingetroffen ist.
Wahl einer Streaming-Feature-Quelle
Wählen Sie eine Quelle basierend auf Ihren Frischeanforderungen und Ihrem bestehenden Aufnahmesystem:
- Verwende eine Punkt, in der
StreamSourcesub-Sekunden-Frische Priorität hat.StreamSourceFunktionen bieten eine End-to-End-Latenz von 200 Millisekunden bei P99. Richte zuerst einen Stream ein und referenziere ihn dann mit einemStreamSource. Stream-Quellen unterstützen Kafka als Eingabe und führen automatisch eine Aufnahme-Delta-Tabelle als historische Kopie der Daten für das Training. - Benutze ein,
DeltaTableSourcewenn du bereits einen latenzarmen Aufnahmepfad in eine Delta-Tabelle hast. Erwarten Sie eine Frische in der Größenordnung von mehreren zehn Sekunden. - Verwenden Sie Zerobus, um ein
DeltaTableSourcezu füllen, wenn Sie noch keinen latenzarmen Aufnahmepfad haben. Die Aufnahme von Zerobus dauert etwa Dutzende Sekunden, daher ist mit einer Fresse der Funktionen unter einer Minute zu rechnen.
Definiere eine Streaming-Funktion mit einer StreamSource
Ein StreamSource verweist über seinen dreiteiligen Namen (catalog.schema.stream_name) auf einen Stream. Ein Stream ist kein Objekt, das über Unity Catalog abgesichert werden kann, sondern ist auf ein Schema in Unity Catalog beschränkt, und der Zugriff wird über die Ingestion-Tabelle des Streams gesteuert. Spaltenverweise in Entitäts-, Zeitreihen- und Funktionsdefinitionen müssen mit value. oder key. als Präfix versehen werden, um anzugeben, welcher Teil der Kafka-Nachricht gelesen werden soll. Geschachtelte Felder werden mit Punktnotation unterstützt (z. B 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)),
),
)
Definiere eine Streaming-Funktion mit einer DeltaTableSource
Um ein Merkmal, das auf einem DeltaTableSource als Streaming-Feature definiert ist, zu materialisieren, übergebe StreamingMode als Auslöser an materialize_features. Die Feature-Definition verwendet dieselben APIs wie ein Batch-Feature, das von einer gestützt wird DeltaTableSource. Delta-Tabellenquellen unterstützen Aggregations- und Spaltenauswahlfunktionen.
Die Quell-Delta-Tabelle muss den Change Data Feed (CDF) durch Einstellung delta.enableChangeDataFeed=trueaktiviert haben.
Das folgende Beispiel definiert und materialisiert eine Aggregationsfunktion mit einer Delta-Tabellenquelle.
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(),
)
Verwenden Sie eine Delta-Tabelle, die von Zerobus befüllt wird
Eine Delta-Tabelle, die von Zerobus gefüllt wird, kann als Quelle für Streaming-Funktionen dienen. Zerobus setzt sich nicht automatisch ein delta.enableChangeDataFeed=true . Sie müssen diese Eigenschaft manuell in der Ziel-Delta-Tabelle einstellen, bevor Sie sie als Streaming-Feature-Quelle verwenden.
Filterbedingungen auf Stromquellen
Verwenden Sie filter_condition es, um Zeilen vor der Aggregation für entweder a StreamSource oder DeltaTableSourcezu filtern.
stream_source = StreamSource(
full_name="my_catalog.my_schema.my_stream",
filter_condition="value.event_type = 'purchase'",
)
Spaltenauswahl aus Streaming-Quellen
ColumnSelection Features funktionieren mit Streamingquellen. Die ausgewählte Spalte stellt den neuesten Wert der Quelle für jede Entität dar und berücksichtigt dabei die Genauigkeit des Zeitpunkts.
Spaltenauswahlfunktionen haben keine TTL. Um einen ausgewählten Wert aus dem Online-Shop zu entfernen, muss die Quelle einen Nullwert für die ausgewählte Spalte ausgeben.
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"),
)
Zugriff auf verschachtelte Felder von einer StreamSource
Für ein StreamSourcekann man auf verschachtelte JSON-Felder mit Dot-Notation zugreifen (zum Beispiel value.nested_field.amount). Zur Laufzeit verwenden die Request-Payload und die Antwort Blattknotennamen (z. B. amount anstelle von value.amount). Blattknotennamen müssen auf alle Entitäts-, Zeitreihen- und Merkmalsnamen innerhalb eines Modells oder Merkmalsspecifikations eindeutig sein, da der servierende Endpunkt Blattnamen verwendet, um Werte zu routen.
Zeitfenster für Streamingfunktionen
Streaming-Features unterstützen bei Aggregationen nur RollingWindow. Rollierende Fenster werden fortlaufend auf Basis der neuesten Daten neu berechnet, was dem Echtzeitcharakter von Streamingquellen entspricht.
TumblingWindow und SlidingWindow sind für die Batchberechnung über feste historische Intervalle ausgelegt.
Beispiel-Notebook für Streaming-Features
Schnellstart-Notebook für Streaming-Feature-Views
Modelltraining und Inferenz
Informationen zum Trainieren von Modellen und Ausführen der Batchunterleitung mit Featureansichten, einschließlich log_model(), score_batch()und create_training_set(), finden Sie unter Train-Modelle mit Featureansichten.
Materialisierung von Features
Nachdem Sie Features definiert haben, können Sie diese zur effizienten Wiederverwendung in Schulungs- und Bereitstellungsworkflows in Offline- oder Onlinespeichern materialisieren. Nach der Materialisierung von Features können Sie die Modelle mit CPU-Modellbereitstellung bereitstellen. Ausführliche Informationen finden Sie unter Materialisieren von Featureansichten.
Bewährte Methoden
Featurebenennung
- Verwenden Sie beschreibende Namen für unternehmenskritische Features.
- Halten Sie einheitliche Benennungskonventionen teamübergreifend ein.
- Verwenden Sie automatisch generierte Namen , während Sie mit der Entwicklung von Features beginnen.
Zeitfenster
- Fensterbegrenzungen an Geschäftszyklen ausrichten (täglich, wöchentlich).
- Kürzere Fenster erfassen aktuelle Trends, können aber laut sein. Längere Fenster erzeugen stabilere Featureverteilungen, können aber die jüngsten Verhaltensverschiebungen verpassen. Wählen Sie basierend darauf, wie schnell sich die zugrunde liegenden Signaländerungen für Ihren Anwendungsfall ändern. Beispielsweise glättet ein 7-tägiges Fenster tägliche Schwankungen und erzeugt konsistente Modelleingaben, während ein 1-Stunden-Fenster schnell auf Verhaltensänderungen reagiert, aber Abweichung verursachen kann, die die Modellleistung beeinträchtigt. Wenn die Genauigkeit Ihres Modells beeinträchtigt wird, wenn sich die Verteilung verschiebt, verwenden Sie ein längeres Fenster, um Eingaben zu stabilisieren.
- Kippfenster und Schiebefenster sind besser skalierbar als kontinuierliche Fenster. Beginnen Sie mit gleitenden Fenstern für die meisten Anwendungsfälle.
Performance
- Materialisieren Sie Features aus derselben Datenquelle in einem einzigen
materialize_featuresAufruf, um Datenscans zu minimieren. - Verwenden Sie dieselbe Granularität (z. B. alle 1-Stunden- oder alle 1-Tage-Foliendauern) für Features in derselben Datenquelle, um eine bessere Gruppierung während der Materialisierung zu ermöglichen.
Entitätsspalten im Vergleich zu Filterbedingungen
Verwenden Sie dieses Entscheidungshandbuch beim Arbeiten mit Features aus derselben Quelltabelle:
Verwenden Sie entity (on create_feature), wenn Sie unterschiedliche Aggregationsebenen benötigen:
-
Features auf Kundenebene (eine Zeile pro Kunde):
entity=["customer_id"] -
Kunden-Händler-Funktionen (mehrere Zeilen pro Kunde):
entity=["customer_id", "merchant_id"] -
Verschiedene Aggregationsebenen können dasselbe
DeltaTableSourcegemeinsam nutzen: Geben Sie für jede Featuredefinition unterschiedlicheentityWerte an.
Verwenden Sie filter_condition (on DeltaTableSource), wenn Sie Zeilen auf derselben Aggregationsebene filtern müssen:
-
Nur hochwertige Transaktionen:
filter_condition="amount > 100"(noch aggregiert pro Kunde) -
Nur abgeschlossene Bestellungen:
filter_condition="status = 'completed'"(noch aggregiert pro Kunde)
Faustregel: Wenn Ihre Änderung zu einer anderen Anzahl von Zeilen pro Entitätswert führen würde, verwenden Sie unterschiedliche entity Werte für Ihre Featuredefinitionen. Wenn Sie lediglich die Zeilen filtern möchten, die zur gleichen Aggregation beitragen, wenden Sie filter_condition auf die Quelle an.
Allgemeine Muster
Kundenanalysen
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)))),
]
Trendanalyse
# 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))),
)
Saisonale Muster
# 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))),
)
Einschränkungen
- Die Namen von Entitäts- und Zeitserienspalten müssen zwischen dem Schulungs-Dataset (bezeichnet) und den Featuredefinitionen übereinstimmen, wenn sie in der
create_training_setAPI verwendet werden. - Der Spaltenname, der als
label-Spalte im Schulungsdatensatz verwendet wird, sollte nicht in den Quelltabellen existieren, die zur Definition vonFeatures verwendet werden. - Eine eingeschränkte Liste von Funktionen (UDAFs) wird in der
create_featureAPI unterstützt. Siehe unterstützte Funktionen. - Entitätsspalten dürfen nicht vom Typ
DATEoderTIMESTAMP. -
RequestSourceunterstützt nur skalare Datentypen, die inScalarDataType(INTEGER,FLOAT, ,BOOLEAN,STRINGDOUBLE,LONG,TIMESTAMP, )DATESHORTdefiniert sind. Komplexe Typen wie Arrays, Maps und Strukturen werden nicht unterstützt. -
RequestSourceUnterstützt keine Aggregationsfunktionen oder Zeitfenster. Es können nurColumnSelectionFunktionen verwendet werden. - Die Gruppe von Entitätsspaltennamen, Zeitserienspaltennamen und Anforderungsfunktionsspaltennamen muss global eindeutig über alle Quellen in einem Trainingssatz oder Bereitstellungsendpunkt sein.
-
score_batchBei serverloser Berechnung ist dies möglicherweise nicht erfolgreich. Verwenden Sie dazu einen klassischen Computecluster mit Databricks Runtime 17.0 ML oder höher.
Informationen zu materialisierungsspezifischen Einschränkungen finden Sie unter "Einschränkungen".