Notatka
Dostęp do tej strony wymaga autoryzacji. Może spróbować zalogować się lub zmienić katalogi.
Dostęp do tej strony wymaga autoryzacji. Możesz spróbować zmienić katalogi.
Ważna
Ta funkcja jest dostępna w publicznej wersji testowej. Administratorzy obszaru roboczego mogą kontrolować dostęp do tej funkcji ze strony Podglądy . Zobacz Zarządzanie wersjami zapoznawczami usługi Azure Databricks.
Widoki funkcji umożliwiają definiowanie i obliczanie funkcji ze źródeł danych. Cechy można definiować przy użyciu różnych źródeł (tabela Delta, strumień Kafka i dane dostępne w czasie żądania) oraz różnych metod obliczeń (agregacje w oknach czasowych, proste selekcje kolumn i inne). W tym przewodniku omówiono następujące przepływy pracy:
- Przepływ pracy tworzenia funkcji
- Użyj
create_featuredo definiowania obiektów funkcji w katalogu Unity, które mogą być używane w przepływach pracy trenowania i serwowania modeli. - Alternatywnie skonstruuj
Featureobiekty lokalnie i używajregister_featureich do przechowywania ich w Unity Catalog później. Funkcje skonstruowane lokalnie mogą być używane zcreate_training_setprzed rejestracją.
- Użyj
-
Przepływ pracy trenowania modelu
- Użyj
create_training_setdo obliczania funkcji zagregowanych w danym punkcie czasu na potrzeby uczenia maszynowego. Aby uzyskać szczegółową dokumentację dotyczącą trenowania przy użyciu widoków funkcji, zobacz Trenowanie modeli za pomocą widoków funkcji.
- Użyj
-
Materializacja funkcji i obsługa przepływu pracy
- Po zdefiniowaniu funkcji za pomocą
create_featurelub pobraniu jej przy użyciuget_feature, można użyćmaterialize_featuresdo zmaterializowania funkcji lub zestawu funkcji w magazynie offline w celu wydajnego ponownego wykorzystania lub do sklepu online w celu obsługi online. - Użyj
create_training_setwidoku zmaterializowanego, aby przygotować wsadowy zestaw danych do trenowania offline.
- Po zdefiniowaniu funkcji za pomocą
Aby uzyskać szczegółowe informacje o interfejsie API, zobacz Dokumentacja interfejsu API widoków funkcji.
Wymagania
Bezserwerowe obliczenia lub klasyczny klaster obliczeniowy z uruchomionym środowiskiem Databricks Runtime 17.0 ML lub nowszym.
Należy zainstalować niestandardowy pakiet języka Python. Uruchom następujące wiersze kodu za każdym razem, gdy uruchamiasz notes:
%pip install databricks-feature-engineering>=0.16.0 dbutils.library.restartPython()
Przykład szybkiego startu
Aby skorzystać z uruchamialnego notesu szybkiego startu, zobacz Przykładowy notes.
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(),
)
Przykładowy notatnik
Notatnik szybkiego startu dotyczący widoków cech
Funkcje przesyłania strumieniowego
Korzystaj z funkcji streamingowych, gdy wartości cech muszą się aktualizować w trybie ciągłym, zamiast w harmonogramie wsadowym. Cechy strumieniowe i wsadowe korzystają z tych samych Feature definicji, funkcji agregujących oraz procesów uczenia i wdrażania.
Funkcje streamingu nie pojawiają się w sklepie offline. Na potrzeby trenowania i wnioskowania wsadowego Databricks oblicza wartości cech na podstawie danych źródłowych.
Funkcje streamingowe spełniają następujące wymagania:
- Musisz podać
online_config. Funkcje przesyłania strumieniowego nie obsługująoffline_config. - Nie można łączyć funkcji strumieniowych i wsadowych w jednym wywołaniu
materialize_features. Wykonaj osobne wywołanie dla każdego typu wyzwalacza. -
transformation_sqlnie jest obsługiwane w przypadku funkcji przesyłania strumieniowego. - Materializacja strumieniowa przetwarza tylko rekordy, które napływają po uruchomieniu potoku danych, i nie uzupełnia wstecznie rekordów historycznych. Agregaty w oknie rolującym zwracają pełne wyniki dopiero po nadejściu pierwszego pełnego okna danych.
Wybór źródła funkcji streamingowych
Wybierz źródło w zależności od wymagań dotyczących świeżości oraz istniejącej konfiguracji pozyskiwania danych:
- Użyj
StreamSource, gdy priorytetem jest odświeżanie krótsze niż sekunda.StreamSourceFunkcje zapewniają opóźnienie P99 end-to-end wynoszące 200 milisekund. Najpierw skonfiguruj strumień, a następnie odwołaj się do niego za pomocąStreamSource. Źródła strumieniowe obsługują platformę Kafka jako źródło danych wejściowych i automatycznie utrzymują tabelę Delta do pozyskiwania danych jako historyczną kopię danych na potrzeby trenowania. - Użyj elementu
DeltaTableSource, jeśli masz już ścieżkę pozyskiwania danych o niskich opóźnieniach do tabeli Delta. Spodziewaj się świeżości rzędu kilkudziesięciu sekund. - Użyj Zerobus, aby uzupełnić
DeltaTableSource, jeśli nie masz jeszcze ścieżki pozyskiwania danych o niskim opóźnieniu. Pozyskiwanie danych przez Zerobus trwa rzędu kilkudziesięciu sekund, więc należy oczekiwać aktualności cech na poziomie poniżej minuty.
Zdefiniuj funkcję streamingową za pomocą StreamSource
Element StreamSource odwołuje się do strumienia według jego trzyczęściowej nazwy (catalog.schema.stream_name). Strumień nie jest obiektem podlegającym zabezpieczeniom w Unity Catalog, ale należy do schematu Unity Catalog, a dostęp do niego jest kontrolowany przez tabelę pozyskiwania danych strumienia. Odwołania do kolumn w definicjach encji, szeregów czasowych i funkcji muszą być poprzedzone prefiksem value. lub key., aby określić, którą część komunikatu Kafka należy odczytać. Zagnieżdżone pola są obsługiwane przy użyciu notacji kropkowej (na przykład 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)),
),
)
Zdefiniuj funkcję streamingową za pomocą DeltaTableSource
Aby zmaterializować funkcję zdefiniowaną na DeltaTableSource jako funkcję strumieniową, przekaż StreamingMode jako wyzwalacz do materialize_features. Definicja funkcji korzysta z tych samych interfejsów API co funkcja wsadowa oparta na DeltaTableSource. Źródła tabel delta wspierają agregację i funkcje wyboru kolumn.
Źródłowa tabela Delta musi mieć włączony strumień danych zmian (change data feed) (CDF) poprzez ustawienie delta.enableChangeDataFeed=true.
Poniższy przykład definiuje i realizuje cechę agregacji za pomocą źródła tabeli Delta.
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(),
)
Użyj tabeli Delta wypełnionej przez Zerobus
Tabela Delta wypełniona przez Zerobus może służyć jako źródło funkcji streamingowych. Zerobus nie ustawia delta.enableChangeDataFeed=true automatycznie. Musisz ręcznie ustawić tę właściwość w docelowej tabeli Delta, zanim użyjesz jej jako źródła funkcji streamingowej.
Warunki filtrowania dla źródeł strumieniowych
Użyj filter_condition, aby filtrować wiersze przed agregacją dla StreamSource lub DeltaTableSource.
stream_source = StreamSource(
full_name="my_catalog.my_schema.my_stream",
filter_condition="value.event_type = 'purchase'",
)
Wybór kolumn z źródeł streamingowych
ColumnSelection Atrybuty przekazują ze źródła najnowszą wartość dla każdego klucza encji bez agregacji. Podczas uczenia modelu wartości cech zachowują dokładność względem danego momentu w czasie.
Funkcje wyboru kolumn nie mają TTL. Aby usunąć wybraną wartość ze sklepu internetowego, źródło musi wygenerować wartość zerową dla wybranej kolumny.
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"),
)
Uzyskiwanie dostępu do zagnieżdżonych pól w StreamSource
W przypadku StreamSource można uzyskać dostęp do zagnieżdżonych pól JSON za pomocą notacji kropkowej (na przykład value.nested_field.amount). Podczas obsługi treść żądania i odpowiedź używają nazw węzłów liści (na przykład amount zamiast value.amount). Nazwy węzłów liściowych muszą być unikalne dla wszystkich nazw podmiotów, szeregów czasowych i cech w modelu lub specyfikacji cech, ponieważ obsługujący punkt końcowy używa nazw liści do trasowania wartości.
Okna czasowe dotyczące funkcji przesyłania strumieniowego
Funkcje przesyłania strumieniowego obsługują tylko agregacje RollingWindow. Okna przesuwne są stale przeliczane na podstawie najnowszych danych, co odpowiada charakterowi źródeł strumieniowych działających w czasie rzeczywistym.
TumblingWindow i SlidingWindow są przeznaczone do obliczeń wsadowych w stałych interwałach historycznych.
Przykładowy notatnik funkcji przesyłania strumieni
Notatnik szybkiego startu dotyczący widoków cech strumieniowych
Trenowanie i wnioskowanie modelu
Aby wytrenować modele i uruchomić wnioskowanie wsadowe za pomocą widoków funkcji, w tym log_model(), score_batch()i create_training_set(), zobacz Trenowanie modeli za pomocą widoków funkcji.
Materializacja atrybutów
Po zdefiniowaniu funkcji można zmaterializować je w trybie offline lub w sklepach online w celu wydajnego ponownego użycia w trenowaniu i obsługiwaniu przepływów pracy. Po zamaterializowaniu funkcji można udostępniać modele przy użyciu serwowania modelu na CPU. Aby uzyskać szczegółowe informacje, zobacz Materialize Feature Views (Zmaterializowanie widoków funkcji).
Najlepsze rozwiązania
Nazewnictwo funkcji
- Użyj nazw opisowych dla funkcji krytycznych dla działania firmy.
- Przestrzegaj spójnych konwencji nazewnictwa w różnych zespołach.
- Użyj automatycznie generowanych nazw podczas tworzenia funkcji.
Okna czasu
- Dopasuj granice okien do cykli biznesowych (codziennie, co tydzień).
- Krótsze okna przechwytują najnowsze trendy, ale mogą być hałaśliwe. Dłuższe okna generują bardziej stabilne dystrybucje funkcji, ale mogą przegapić ostatnie zmiany behawioralne. Wybierz zależnie od tego, jak szybko zmienia się bazowy sygnał dla twojego przypadku użycia. Na przykład 7-dniowe okno wygładza codzienne wahania i generuje spójne dane wejściowe do modelu, podczas gdy 1-godzinne okno reaguje szybko na zmiany behawioralne, ale może wprowadzać wariancję, która pogarsza jakość działania modelu. Jeśli dokładność modelu spada, gdy rozkład zmienia się, użyj dłuższego okna, aby ustabilizować dane wejściowe.
- Okna wirujące i przesuwne są bardziej skalowalne niż okna kroczące (ciągłe). Zacznij od okien przesuwnych dla większości przypadków użycia.
Performance
- Zmaterializuj funkcje z tego samego źródła danych w jednym
materialize_featureswywołaniu, aby zminimalizować skanowania danych. - Użyj tego samego stopnia szczegółowości (na przykład wszystkich 1-godzinnych lub wszystkich 1-dniowych czasów trwania slajdów) dla funkcji w tym samym źródle danych, aby umożliwić lepsze grupowanie podczas materializacji.
Kolumny encji a warunki filtrowania
Skorzystaj z tego przewodnika decyzyjnego podczas pracy z funkcjami z tej samej tabeli źródłowej:
Użyj entity (na create_feature), jeśli potrzebujesz różnych poziomów agregacji:
-
Funkcje na poziomie klienta (jeden wiersz na klienta):
entity=["customer_id"] -
Funkcje klientów i sprzedawców (wiele wierszy dla każdego klienta):
entity=["customer_id", "merchant_id"] -
Różne poziomy agregacji mogą współdzielić te same
DeltaTableSourcewartości: określ różneentitywartości w każdej definicji funkcji
Użyj filter_condition (na DeltaTableSource), gdy musisz filtrować wiersze na tym samym poziomie agregacji:
-
Tylko transakcje o wysokiej wartości:
filter_condition="amount > 100"(nadal agregowane na klienta) -
Tylko ukończone zamówienia:
filter_condition="status = 'completed'"(nadal agregowane według klienta)
Ogólna zasada: Jeśli zmiana spowoduje powstanie innej liczby wierszy dla wartości encji, użyj różnych wartości entity w definicjach funkcji. Jeśli po prostu filtrujesz wiersze, które mają wpływ na tę samą agregację, użyj filter_condition na źródle.
Często używane wzorce
Analiza klientów
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)))),
]
Analiza trendów
# 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))),
)
Wzorce sezonowe
# 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))),
)
Ograniczenia
- Nazwy kolumn bytów i szeregów czasowych muszą być zgodne między oznaczonym zestawem danych treningowego a definicjami cech używanymi w interfejsie
create_training_setAPI. - Nazwa kolumny wykorzystywana jako
labelw zestawie danych treningowych nie powinna istnieć w tabelach źródłowych używanych do definiowaniaFeature. - Ograniczona lista funkcji (UDAFs) jest obsługiwana w interfejsie
create_featureAPI. Zobacz Obsługiwane funkcje. - Kolumny jednostek nie mogą być typu
DATEaniTIMESTAMP. -
RequestSourceobsługuje tylko typy danych skalarnych zdefiniowane w programieScalarDataType(INTEGER,FLOAT,BOOLEAN,STRING,DOUBLE,LONG,TIMESTAMP,DATE).SHORTTypy złożone, takie jak tablice, mapy i struktury, nie są obsługiwane. -
RequestSourcenie obsługuje funkcji agregacji ani okien czasowych. Można używać tylkoColumnSelectionfunkcji. - Zestaw nazw kolumn jednostki, nazwy kolumn szeregów czasowych i nazwy kolumn funkcji zapytania muszą być globalnie unikatowe we wszystkich źródłach w zbiorze treningowym lub punkcie serwisowym.
-
score_batchmoże nie działać w środowisku bezserwerowym. Aby obejść ten problem, użyj klastra classic compute działającego w środowisku Databricks Runtime 17.0 ML lub nowszym.
W celu zapoznania się z ograniczeniami specyficznymi dla materializacji, zobacz Ograniczenia.