Notitie
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen u aan te melden of de directory te wijzigen.
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen de mappen te wijzigen.
Important
Deze functie bevindt zich in openbare preview-versie. Werkruimtebeheerders kunnen de toegang tot deze functie beheren vanaf de pagina Previews . Zie Azure Databricks previews beheren.
Met functieweergaven kunt u functies van gegevensbronnen definiëren en berekenen. Functies kunnen worden gedefinieerd met behulp van verschillende bronnen (Delta-tabel, Kafka Stream en aanvraagtijdgegevens) en berekeningen (tijdvensteraggregaties, eenvoudige kolomselecties en meer). In deze handleiding worden de volgende werkstromen behandeld:
-
Werkstroom voor functieontwikkeling
- Hiermee
create_featuredefinieert u functieobjecten voor Unity Catalog die kunnen worden gebruikt in modeltraining en het leveren van werkstromen. - U kunt objecten ook lokaal maken
Featureen gebruikenregister_featureom ze later in Unity Catalog op te slaan. Lokaal samengestelde functies kunnen worden gebruikt metcreate_training_setvóór de registratie.
- Hiermee
-
Werkstroom voor modeltraining
- Gebruik
create_training_setdeze functie om geaggregeerde functies van een bepaald tijdstip voor machine learning te berekenen. Raadpleeg Modellen trainen met Feature Views voor gedetailleerde documentatie over het trainen met Feature Views.
- Gebruik
-
Kenmerk materialisatie en voorziening werkstroom
- Na het definiëren van een kenmerk met
create_featureof het ophalen ervan met behulp vanget_feature, kunt u hetmaterialize_featureskenmerk of de set kenmerken materialiseren naar een offline-winkel voor efficiënt hergebruik of naar een online-winkel voor online-aanbieding. - Gebruik
create_training_setmet de gemaakte weergave om een gegevensset voor offline batchtraining voor te bereiden.
- Na het definiëren van een kenmerk met
Zie de API-naslaginformatie voor functieweergaven voor API-details.
Requirements
Serverloze berekening of een klassiek rekencluster met Databricks Runtime 17.0 ML of hoger.
U moet het aangepaste Python-pakket installeren. Voer de volgende regels code uit telkens wanneer je een notebook uitvoert:
%pip install databricks-feature-engineering>=0.16.0 dbutils.library.restartPython()
Quickstart-voorbeeld
Voor een uitvoerbaar quickstart-notitieblok, zie Voorbeeldnotitieblok.
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(),
)
voorbeeldnotitieblok
Snelstartnotebook voor featureweergaven
Streaming-functies
Gebruik streamingfuncties wanneer featurewaarden continu moeten worden bijgewerkt in plaats van op een batchschema. Streaming- en batchfunctionaliteit gebruiken dezelfde Feature constructors, aggregatiefuncties en trainings- en uitvoeringsworkflows.
Streamingfuncties materialiseren zich niet in een offline winkel. Voor training en batch-inferentie berekent Databricks featurewaarden op basis van de brongegevens.
Streamingfuncties hebben de volgende vereisten:
- Je moet
online_configopgeven. streamingfuncties ondersteunen geenoffline_config. - Je kunt streaming- en batchfuncties niet in één
materialize_featuresgesprek combineren. Maak een apart telefoontje voor elk type trigger. -
transformation_sqlwordt niet ondersteund voor streamingfuncties. - Streaming materialisatie verwerkt alleen records die na de start van de pipeline aankomen en vult geen historische records op. Rolling-window-aggregaten geven volledige resultaten pas op nadat het eerste volledige datavenster is binnengekomen.
Een bron kiezen voor streamingfuncties
Kies een bron op basis van je versheidsbehoeften en bestaande innamesysteem:
- Gebruik een
StreamSourcewanneer actualiteit binnen minder dan een seconde de hoogste prioriteit heeft.StreamSourceFuncties bieden een end-to-end latentie van 200 milliseconden voor de P99. Stel eerst een stream op, en verwijs die vervolgens met eenStreamSource. Stroombronnen ondersteunen Kafka als invoer en onderhouden automatisch een inname-Delta-tabel als historische kopie van de data voor training. - Gebruik een
DeltaTableSourceals je al een invoerpad met lage latentie naar een Delta-tabel hebt. Reken op een actualiteit in de orde van tientallen seconden. - Gebruik Zerobus om een
DeltaTableSourcete vullen wanneer je nog geen ingestiepad met lage latentie hebt. Zerobus-inname duurt enkele tientallen seconden, dus kun je rekenen op feature-actualiteit van minder dan een minuut.
Definieer een streamingfunctie met behulp van een StreamSource
Een StreamSource verwijst naar een Stream op basis van de driedelige naam (catalog.schema.stream_name). Een Stream is geen beveiligbaar object voor Unity Catalog, maar is beperkt tot een Unity Catalog-schema en toegang wordt beheerd door de opnametabel van Stream. Kolomverwijzingen in entiteits-, tijdreeks- en functiedefinities moeten worden voorafgegaan door value. of key. om aan te geven welk deel van het Kafka-bericht moet worden gelezen. Geneste velden worden ondersteund met punt notatie (bijvoorbeeld 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)),
),
)
Definieer een streamingfunctie met behulp van een DeltaTableSource
Om een kenmerk te materialiseren dat op een DeltaTableSource als een streaming-feature is gedefinieerd, geef StreamingMode als de trigger door aan materialize_features. De feature-definitie gebruikt dezelfde API's als een batchfeature die wordt ondersteund door een DeltaTableSource. Delta-tabelbronnen ondersteunen aggregatie- en kolomselectiefuncties.
De Delta-brontabel moet feed voor wijzigingsgegevens (CDF) ingeschakeld hebben door delta.enableChangeDataFeed=true in te stellen.
Het volgende voorbeeld definieert en materialiseert een aggregatiefunctie met een Delta-tabel bron.
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(),
)
Gebruik een Delta-tabel die wordt ingevuld door Zerobus
Een Delta-tabel die door Zerobus wordt ingevuld kan dienen als bron van streamingfuncties. Zerobus stelt delta.enableChangeDataFeed=true niet automatisch in. Je moet deze eigenschap handmatig instellen in de target Delta-tabel voordat je deze als bron van streamingfuncties gebruikt.
Filtervoorwaarden voor streamingbronnen
Gebruik filter_condition om rijen te filteren vóór aggregatie voor ofwel een StreamSource of DeltaTableSource.
stream_source = StreamSource(
full_name="my_catalog.my_schema.my_stream",
filter_condition="value.event_type = 'purchase'",
)
Kolomselectie uit streamingbronnen
ColumnSelection Functies geven de nieuwste waarde van de bron voor elke entiteitssleutel door zonder aggregatie. Tijdens de training voldoen feature-waarden aan punt-in-tijd-nauwkeurigheid.
Kolomselectiefuncties hebben geen TTL. Om een geselecteerde waarde uit de online winkel te verwijderen, moet de bron een nulwaarde voor de geselecteerde kolom uitgeven.
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"),
)
Geneste velden uit een StreamSource openen
Voor een StreamSource, kun je geneste JSON-velden benaderen met behulp van dot-notatie (bijvoorbeeld value.nested_field.amount). Tijdens het uitvoeren van de aanvraag gebruiken de payload en de respons de namen van leaf-knooppunten (bijvoorbeeld amount in plaats van value.amount). Bladknooppuntnamen moeten uniek zijn ten opzichte van alle entiteits-, tijdreeks- en featurenamen binnen een model of Feature Spec, omdat het serving-eindpunt bladnamen gebruikt om waarden te routeren.
Tijdvensters voor streamingfuncties
Streamingfuncties bieden alleen RollingWindow ondersteuning voor aggregaties. Schuivende vensters worden voortdurend opnieuw berekend op basis van de meest recente gegevens, wat aansluit bij het realtimekarakter van streaminggegevensbronnen.
TumblingWindow en SlidingWindow zijn ontworpen voor batchberekening via vaste historische intervallen.
Voorbeeldnotebook voor streamingfuncties
Notebook voor snelstart van streamingfeatureweergaven
Modeltraining en inferentie
Als u modellen wilt trainen en batchdeductie wilt uitvoeren met functieweergaven, waaronder log_model(), score_batch()en create_training_set(), raadpleegt u Modellen trainen met functieweergaven.
Kenmerk-materialisatie
Nadat u functies hebt gedefinieerd, kunt u ze naar offline of online winkels brengen voor efficiënt hergebruik in trainings- en servicewerkstromen. Nadat u functies hebt gerealiseerd, kunt u modellen bedienen met behulp van cpu-modellen. Zie Feature Views materialiseren voor meer informatie.
Beste praktijken
Naamgeving van functies
- Gebruik beschrijvende namen voor bedrijfskritieke functies.
- Volg consistente naamconventies tussen teams.
- Gebruik automatisch gegenereerde namen wanneer u begint met het ontwikkelen van functies.
Tijdvensters
- Venstergrenzen uitlijnen met bedrijfscycli (dagelijks, wekelijks).
- Kortere vensters leggen recente trends vast, maar kunnen luidruchtig zijn. Langere vensters produceren stabielere functiedistributies, maar kunnen recente gedragsverschuivingen missen. Kies op basis van hoe snel het onderliggende signaal verandert voor uw use-case. Een venster van 7 dagen verzacht bijvoorbeeld dagelijkse schommelingen en produceert consistente modelinvoer, terwijl een venster van 1 uur snel reageert op gedragswijzigingen, maar kan afwijking veroorzaken die de modelprestaties verslechtert. Als de nauwkeurigheid van uw model verslechtert wanneer de verdeling verschuift, gebruikt u een langer venster om de invoer te stabiliseren.
- Tumbling en glijdende vensters zijn schaalbaarder dan rollende vensters. Begin met schuifvensters voor de meeste gebruiksvoorbeelden.
Performance
- Materialiseer functies van dezelfde gegevensbron in één
materialize_featuresaanroep om gegevensscans te minimaliseren. - Gebruik dezelfde granulariteit (bijvoorbeeld alle met een duur van 1 uur of 1 dag) voor kenmerken in dezelfde gegevensbron om betere groepering mogelijk te maken tijdens de materialisatie van gegevens.
Entiteitskolommen versus filtervoorwaarden
Gebruik deze handleiding voor beslissingen bij het werken met functies uit dezelfde brontabel:
Gebruik entity (aan create_feature) wanneer u verschillende aggregatieniveaus nodig hebt:
-
Functies op klantniveau (één rij per klant):
entity=["customer_id"] -
Klanten-verkoper functies (meerdere rijen per klant):
entity=["customer_id", "merchant_id"] -
Verschillende aggregatieniveaus kunnen hetzelfde
DeltaTableSourcedelen: verschillendeentitywaarden opgeven voor elke functiedefinitie
Gebruik filter_condition (aan DeltaTableSource) wanneer u rijen op hetzelfde aggregatieniveau wilt filteren:
-
Alleen transacties met hoge waarde:
filter_condition="amount > 100"(nog steeds geaggregeerd per klant) -
Alleen voltooide orders:
filter_condition="status = 'completed'"(nog steeds samengevoegd per klant)
Vuistregel: Als uw wijziging resulteert in een ander aantal rijen per entiteitswaarde, gebruikt u verschillende entity waarden voor uw functiedefinities. Als u alleen filtert welke rijen bijdragen aan dezelfde aggregatie, gebruik dan filter_condition op de bron.
Algemene patronen
Klantanalyse
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))),
)
Seizoensgebonden patronen
# 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
- Namen van entiteits- en tijdreekskolommen moeten overeenkomen tussen de trainingsgegevensset (gelabeld) en de functiedefinities wanneer deze worden gebruikt in de
create_training_setAPI. - De kolomnaam die wordt gebruikt als de
labelkolom in de trainingsgegevensset mag niet voorkomen in de brontabellen die worden gebruikt voor het definiëren vanFeatures. - Een beperkte lijst met functies (UDAF's) wordt ondersteund in de
create_featureAPI. Zie Ondersteunde functies. - Entiteitskolommen kunnen niet van het type
DATEzijn ofTIMESTAMP. -
RequestSourceondersteunt alleen scalaire gegevenstypen die zijn gedefinieerd inScalarDataType(INTEGER,FLOAT,BOOLEAN,STRING,DOUBLE,LONG,TIMESTAMP,DATE).SHORTComplexe typen, zoals matrices, kaarten en structs, worden niet ondersteund. -
RequestSourcebiedt geen ondersteuning voor aggregatiefuncties of tijdvensters. AlleenColumnSelectionfuncties kunnen worden gebruikt. - De set van kolomnamen van entiteiten, tijdreekskolomnamen en kolomnamen van aanvraagfuncties moet globaal uniek zijn over alle bronnen in een trainingsset of een uitvoerend eindpunt.
-
score_batchkan niet slagen op serverloze berekeningen. U kunt dit omzeilen met behulp van een klassiek rekencluster met Databricks Runtime 17.0 ML of hoger.
Zie Beperkingen voor materialisatiespecifieke beperkingen.