Dokumentacja referencyjna interfejsu API Feature Views

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.

Kontrola dostępu

Funkcje są zarządzalnymi obiektami wykazu aparatu Unity. Dostęp do funkcji jest kontrolowany przez CREATE FEATUREuprawnienia katalogu , READ FEATUREi MANAGE aparatu Unity. Aby uzyskać pełne opisy, zobacz Dokumentacja uprawnień wykazu aparatu Unity.

  • CREATE FEATURE — Wymagane do utworzenia funkcji w schemacie. create_feature i register_feature wymagaj CREATE FEATURE w schemacie nadrzędnym. Zgodnie z zasadą najniższych uprawnień przyznaj CREATE FEATURE na poziomie schematu; możesz również przyznać go w wykazie, aby umożliwić tworzenie funkcji w dowolnym schemacie w tym wykazie.
  • READ FEATURE — wymagane do odczytania funkcji i jej danych. get_feature, create_training_seti odczytywanie zmaterializowanych danych funkcji na potrzeby trenowania lub obsługi wymaganej READ FEATURE dla tej funkcji. READ FEATURE przyznane w schemacie lub wykazie ma zastosowanie do wszystkich bieżących i przyszłych funkcji, które zawiera.
  • MANAGE — Wymagane do zarządzania cyklem życia i dotacji funkcji. Usuwanie funkcji za pomocą delete_featurefunkcji i materializowanie funkcji za pomocą materialize_features funkcji lub delete_materialized_featurewymagaj MANAGE jej.

Wszystkie operacje funkcji wymagają USE CATALOG również w katalogu nadrzędnym i USE SCHEMA w schemacie nadrzędnym. Aby dowiedzieć się, jak MANAGE i READ FEATURE zastosować do materializacji, zobacz Uprawnienia.

Interfejs API widoku funkcji

Feature konstruktor i register_feature()

Zalecanym podejściem Feature jest utworzenie obiektu lokalnie i użycie register_feature go do utrwalania go w wykazie aparatu Unity. Ten dwuetapowy przepływ pracy umożliwia eksperymentowanie z funkcjami (w tym create_training_set) przed ich zarejestrowaniem.

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() rejestruje lokalnie skonstruowany Feature w wykazie aparatu Unity.

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() weryfikuje, konstruuje i natychmiast rejestruje funkcję w wykazie aparatu Unity w jednym kroku. Użyj tej funkcji, jeśli nie musisz najpierw eksperymentować z funkcją lokalnie.

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

Parametry:

  • source: źródło danych używane w obliczeniach funkcji (DeltaTableSource, StreamSource, lub RequestSource).
  • function: Element AggregationFunction , który łączy operator (na przykład Sum(input="amount")), kolumnę wejściową i przedział czasu. Lub ColumnSelection("column_name") w przypadku funkcji przekazywania.
  • catalog_name: nazwa wykazu wykazu aparatu Unity dla funkcji.
  • schema_name: nazwa schematu wykazu aparatu Unity dla funkcji.
  • entity: Lista nazw kolumn definiujących klucze agregacji lub wyszukiwania (klucze podstawowe). Wymagane dla wszystkich typów źródłowych z wyjątkiem RequestSource. Na przykład ["user_id"] agreguje lub wyszukuje poszczególnych użytkowników.
  • timeseries_column: kolumna znacznika czasu używana na potrzeby agregacji przedziału czasu lub wyboru najnowszej wartości. Wymagane dla wszystkich typów źródłowych z wyjątkiem RequestSource.
  • name: opcjonalna nazwa funkcji. Jeśli pominięto, automatycznie wygenerowane z kolumny wejściowej, funkcji i okna (na przykład amount_avg_rolling_7d).
  • description: opcjonalny opis funkcji.

Zwraca: Zweryfikowane wystąpienie funkcji

Zgłasza: ValueError, jeśli sprawdzanie poprawności nie powiedzie się

delete_feature()

Usuwa funkcję z wykazu aparatu Unity przy użyciu w pełni kwalifikowanej nazwy.

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

Przed usunięciem funkcji usuń lub zaktualizuj wszystkie modele lub specyfikacje funkcji odwołujące się do tej funkcji. Jeśli funkcja została zmaterializowana, najpierw usuń zmaterializowaną funkcję. Zobacz Jak usunąć zmaterializowaną funkcję.

Nazwy generowane automatycznie

Po name pominięciu nazwa jest generowana automatycznie. Wygenerowane nazwy są zgodne ze wzorcem: {column}_{function}_{window}. Przykład:

  • price_avg_rolling_1h (średnia cena 1-godzinna)
  • transaction_count_rolling_30d_1d (30-dniowa liczba transakcji z 1d opóźnieniem od znacznika czasu zdarzenia)

Obsługiwane funkcje

Funkcje agregacji

Note

Funkcje agregacji są opakowane razem AggregationFunction z przedziałem czasu, zgodnie z opisem w oknach czasowych. Każda funkcja przyjmuje parametr określający kolumnę źródłową input do agregowania.

Function Description Przykładowy przypadek użycia
Sum(input="column") Suma wartości Dzienne użycie aplikacji przez użytkownika w minutach
Avg(input="column") Średnia wartości Średnia kwota transakcji
Count(input="column") Liczba rekordów Liczba logowań na użytkownika
Min(input="column") Wartość minimalna Najniższa częstotliwość serca zarejestrowana przez urządzenie do noszenia
Max(input="column") Wartość maksymalna Najwyższa kwota transakcji na sesję
StddevPop(input="column") Odchylenie standardowe populacji Zmienność dziennej kwoty transakcji we wszystkich klientach
StddevSamp(input="column") Odchylenie standardowe próbki Zmienność stawek kliknięć kampanii reklamowej
VarPop(input="column") Wariancja populacji Rozprzestrzenianie odczytów czujników dla urządzeń IoT w fabryce
VarSamp(input="column") Przykładowa wariancja Rozprzestrzenianie ocen filmów w grupie próbkowanej
ApproxCountDistinct(input="column", relativeSD=0.05) Przybliżona liczba unikalnych Unikatowa liczba zakupionych przedmiotów
ApproxPercentile(input="column", percentile=0.95, accuracy=100) Przybliżony percentyl Opóźnienie odpowiedzi p95
First(input="column") Pierwsza wartość Znacznik czasu pierwszego logowania
Last(input="column") Ostatnia wartość Najnowsza kwota zakupu
FirstN(input="column", n=3) Pierwsze n wartości jako tablica Pierwsze trzy produkty obejrzane podczas sesji
LastN(input="column", n=3) Wartości końcowe n jako tablica Trzy najnowsze statusy spraw o wsparcie
FirstDistinct(input="column", n=3) Najpierw n różne wartości jako tablica Pierwsze trzy odrębne kategorie produktów obejrzane
LastDistinct(input="column", n=3) Ostatnie n różne wartości jako tablica Trzy najnowsze, odrębne kategorie handlowe

Note

First, Last, FirstN, LastN, FirstDistinct, , oraz LastDistinct domyślnie zawierają wartości zerowe. Aby pominąć wartości null, dodaj obiekt filter_condition , który jawnie wyklucza kolumny wejściowe, które mają wartość null.

FirstN, LastN, , oraz LastDistinct używają cech timeseries_column do porządkowania wierszy wejściowych i zwracania tablicy zawierającej wartości do wartości nFirstDistinct. Parametr musi n być dodatnią liczbą całkowitą. FirstN oraz FirstDistinct wybierać wartości od najwcześniejszego do najpóźniejszego. LastN oraz LastDistinct wybierz wartości od najpóźniejszej do najwcześniejszej i zwróć wybrane wartości w kolejności znaczników czasu. FirstDistinct oraz usuwanie LastDistinct wartości zduplikowanych podczas wybierania wartości w tym kierunku.

Na przykład, jeśli wiersze źródłowe dla danej jednostki są uporządkowane przez event_time , ["A", "A", "B", "C", "B", "B"]następujące funkcje zwracają:

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, FirstDistinct, oraz LastDistinct wymagają databricks-feature-engineering wersji 0.17.0 lub nowzej.

ColumnSelection (przekazywanie)

ColumnSelection wybiera jedną kolumnę ze źródła bez stosowania żadnej agregacji. Jest on owinięty bezpośrednio w parametrze function (a nie wewnątrz AggregationFunction). Zwracany typ jest wnioskowany ze schematu źródłowego.

Function Description Przykładowy przypadek użycia
ColumnSelection("col") Najnowsza wartość kolumny (bez agregacji) Najnowsza kategoria dostawcy, przekazywanie pola żądania

ColumnSelection można używać z dowolnym źródłem danych:

  • DeltaTableSource: zwraca najnowszą wartość na klucz jednostki za pośrednictwem sprzężenia do punktu w czasie (bez agregacji okna wyszukiwania).
  • StreamSource: Zwraca najnowszą wartość na klucz encji ze strumienia (brak agregacji okna wstecznego).
  • RequestSource: przechodzi przez wartość podaną w czasie wnioskowania (lub wyodrębniony z oznaczonej ramki danych w czasie trenowania).
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",
)

Przykład: funkcje agregacji i wyboru kolumn

W poniższym przykładzie przedstawiono funkcje zdefiniowane w tym samym źródle danych.

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

Funkcje z warunkami filtrowania

Parametr filter_condition umożliwia filtrowanie wierszy z tabeli źródłowej przed obliczeniami agregacji. Ta funkcja jest klauzulą SQL WHERE , która jest stosowana przed grupowaniem i agregowaniem danych.

Note

filter_condition filtruje wiersze przed agregacją, podobnie jak klauzula SQL WHERE zastosowana przed GROUP BY. Nie zmienia stopnia szczegółowości, który jest zawsze definiowany przez entity definicję funkcji.

Filtry są przydatne podczas pracy z dużymi tabelami źródłowymi, które zawierają nadzbiór danych potrzebnych do obliczeń funkcji, i zminimalizować potrzebę tworzenia oddzielnych widoków na podstawie tych tabel.

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

Źródła danych

DeltaTableSource

DeltaTableSource jest efemerycznym obiektem Python używanym do definiowania sposobu obliczania funkcji z tabeli źródłowej. Nie tworzy nowej tabeli. Określa konfigurację odczytywania danych i agregowania funkcji.

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
)

Parametry:

  • catalog_name, , schema_nametable_name: zidentyfikuj źródłową tabelę delty w wykazie aparatu Unity.
  • filter_condition: klauzula SQL WHERE zastosowana przed agregacją. Przykład: "status = 'completed'".
  • transformation_sql: wyrażenie SQL SELECT zastosowane do tabeli źródłowej. Służy do zmieniania nazw kolumn, typów rzutowania lub kolumn pochodnych obliczeń przed agregacją. W przypadku pominięcia wszystkie kolumny są zaznaczone (*). Przykład: "user_id, CAST(amount AS DOUBLE) AS amount, event_time".
  • dataframe_schema: schemat wynikowej ramki danych po przekształceniach w formacie Spark StructType JSON (z df.schema.json()). Wymagane w przypadku transformation_sql dostarczenia. Informuje system o nazwach kolumn i typach, które wynikają z transformacji.

Po ustawieniu obu filter_condition elementów i transformation_sql wynikowe zapytanie to: SELECT {transformation_sql} FROM {table} WHERE {filter_condition}.

Note

Wartość (określona timeseries_column w definicji funkcji, a nie na DeltaTableSource) musi być typu TimestampType lub DateType. Typy liczb całkowitych mogą działać, ale powodują utratę dokładności dla agregacji przedziału czasu.

Przykład: Używanie transformation_sql dla przekształceń kolumn

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

Przykład: wyprowadzanie transformation_sql i dataframe_schema z ramki danych PySpark

Przekształcenie można napisać jako zapytanie PySpark, a następnie wyodrębnić schemat z wynikowej ramki danych:

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

Wspierane transformation_sql wyrażenia

Te same zasady dotyczą transformation_sql na i DeltaTableSourceStreamSource.

transformation_sql obsługuje dowolne wyrażenia wierszowe; operacje oceniane niezależnie dla każdego wiersza. Nie zmieniają liczby wierszy ani zgodności jeden do jednego ze źródłami. Wyrażenia wierszowe obejmują zmiany nazw kolumn, rzuty, operacje arytmetyczne i inne.

Operacje zmieniające kształt lub liczbę wierszy nie są obsługiwane, takie jak agregacje takie jak SUM() lub COUNT(). Zamiast tego użyj AggregationFunction definicji funkcji.

DeltaTableSource.from_sql()

Dla wygody możesz utworzyć element DeltaTableSource na podstawie zapytania SQL. Metoda analizuje zapytanie, aby automatycznie wyodrębnić nazwę tabeli, transformation_sqli filter_condition.

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

Obsługiwane są tylko proste SELECT ... FROM ... [WHERE ...] zapytania. Złożone sieci SQL (JOIN, podzapytania, KARTY CT, UNION) są odrzucane. W przypadku złożonych zapytań skonstruuj DeltaTableSource bezpośrednio za pomocą transformation_sql elementów i 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",
)

Iterowanie za pomocą polecenia to_dataframe()

Użyj polecenia source.to_dataframe() , aby wyświetlić podgląd danych, które będą używane do obliczeń funkcji. Jest to przydatne w przypadku iteracji filter_condition i transformation_sql do momentu wygenerowania oczekiwanych wyników.

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

Informacje o jednostkach

Kolumny jednostek definiują poziom agregacji funkcji. Są one określone w Feature definicji, a nie na DeltaTableSource. Jednostki określają:

  • Sposób grupowania danych: funkcje są agregowane na unikatową kombinację wartości jednostek (podobnie jak GROUP BY w języku SQL)
  • Struktura klucza podstawowego: każda unikatowa kombinacja jednostek powoduje wyświetlenie jednego wiersza obliczonych funkcji

Przykład: funkcje na poziomie klienta

Poniższy kod agreguje funkcje na poziomie klienta (jeden wiersz na klienta):

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

Przykład: funkcje na poziomie sklepu klienta

Aby agregować funkcje na bardziej szczegółowym poziomie (jeden wiersz na kombinację magazynu klientów), użyj wielu kolumn jednostek:

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

Jeśli potrzebujesz funkcji na różnych poziomach agregacji (na przykład na poziomie klienta i na poziomie magazynu klienta), użyj różnych entity wartości w definicjach funkcji. Te same DeltaTableSource funkcje mogą być współużytkowane przez różne konfiguracje jednostek.

StreamSource

StreamSource odwołuje się do strumienia. Strumień zawiera konfigurację połączenia, uwierzytelniania, schematu i pozyskiwania dla źródła przesyłania strumieniowego. W przypadku platformy Kafka odwołania do kolumn w definicjach funkcji muszą być poprzedzone prefiksem value. lub key. wskazać, która część komunikatu ma być odczytywana.

StreamSource(
    full_name: str,                       # Required: Three-part Stream name (catalog.schema.stream)
    filter_condition: Optional[str] = None,      # Optional: SQL WHERE clause applied before aggregation
    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
)

Parametry:

  • full_name: pełna trzyczęściowa nazwa strumienia (na przykład "my_catalog.my_schema.my_stream").
  • filter_condition (opcjonalnie): klauzula SQL WHERE stosowana do strumieniowego przesyłania danych przed agregacją przy użyciu odwołań do kolumn z prefiksem kropki (na przykład "value.event_type = 'purchase'").
  • transformation_sql (opcjonalnie): Wyrażenie SQL SELECT stosowane przed agregacją lub wyborem kolumn, z użyciem odniesień do key struktur i value z prefiksami kropki. Obsługuje te same wyrażenia wierszowe co DeltaTableSource. Jeśli zostanie pominięty, źródło korzysta ze wszystkich kolumn ().*
  • dataframe_schema: Schemat JSON Spark StructType dla prognozowanego wyjścia. Wymagane, jeśli ustawisz transformation_sql.
from databricks.feature_engineering.entities import StreamSource

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

Wyprowadzamy dataframe_schema , uruchamiając projekcję na tablicy pobierania strumienia, która odsłania key struktury i value .

transformation_sql = (
    "value.amount * value.conversion_rate AS converted_amount, "
    "struct(value.user_id AS user_id, value.event_time AS time) AS event"
)

ingestion_table = "my_catalog.my_schema.events_ingestion"
dataframe_schema = spark.sql(
    f"SELECT {transformation_sql} FROM {ingestion_table} LIMIT 0"
).schema.json()

stream_source = StreamSource(
    full_name="my_catalog.my_schema.my_stream",
    transformation_sql=transformation_sql,
    dataframe_schema=dataframe_schema,
)

RequestSource

RequestSource Definiuje schemat dla danych, które są udostępniane w czasie wnioskowania w ładunku żądania, a nie wyszukiwania ze wstępnie zmaterializowanej tabeli. Podczas trenowania te kolumny są wyodrębniane z oznaczonej ramki danych przekazanej do elementu create_training_set. Podczas obsługi modelu obiekt wywołujący musi uwzględnić je w ładunku żądania HTTP.

RequestSource jest używany z elementem ColumnSelection (do bezpośredniego przekazywania wartości). Nie obsługuje ona funkcji agregacji ani okien czasowych.

Definiowanie schematu

Zdefiniuj schemat jako listę FieldDefinition obiektów, z których każda określa nazwę kolumny i :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),
    ]
)

Obsługiwane typy danych

RequestSourceobsługuje typy skalarne zdefiniowane w : ScalarDataType, , INTEGER, FLOATBOOLEANSTRINGDOUBLELONG, . TIMESTAMPDATESHORT Typy złożone, takie jak tablice, mapy i struktury, nie są obsługiwane.

Jak dane żądań są nawodnione

Kontekst Behavior
Szkolenie (create_training_set) Kolumny są wyodrębniane z oznaczonej ramki danych. Typy są weryfikowane względem zadeklarowanego schematu. Niezgodności powodują błąd (bez niejawnego rzutowania).
Obsługa (punkt końcowy modelu) Kolumny są pobierane z dataframe_records żądania HTTP lub dataframe_split w żądaniu HTTP. Wartości JSON są rzutowane na zadeklarowane typy (np. liczba JSON → DOUBLE).

Podpis modelu

Gdy model jest rejestrowany przy użyciu log_model zestawu szkoleniowego zawierającego RequestSource funkcje, RequestSource kolumny są dodawane do podpisu modelu MLflow jako wymagane dane wejściowe. Oznacza to, że schemat interfejsu API obsługującego punkt końcowy odzwierciedla pola, które obiekty wywołujące muszą podać w czasie wnioskowania.

Interfejs API trenowania i wnioskowania

create_training_set i score_batch obliczanie prawidłowych wartości funkcji punktu w czasie na żądanie z danych źródłowych. W przypadku funkcji obsługujących materializację w trybie offline, takich jak agregacje okien przesuwnych w źródłach tabel różnicowych, materializacja funkcji najpierw w magazynie offline zwiększa wydajność obu operacji. Gdy zmaterializowane funkcje trybu offline są dostępne, operacje odczytują wstępnie skompilowane dane offline zamiast ponownie skompilować wartości funkcji ze źródła. Zobacz Materialize feature Views to materialize features to materialize features to an offline store (Materialize feature Views to materialize features to an offline store

create_training_set()

Tworzy zestaw danych szkoleniowych z prawidłowymi obliczeniami funkcji do punktu w czasie. Aby uzyskać szczegółowe informacje, zobacz Trenowanie modeli za pomocą widoków funkcji.

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

Rejestruje model z metadanymi funkcji na potrzeby śledzenia pochodzenia i automatycznego wyszukiwania funkcji podczas wnioskowania. Aby uzyskać szczegółowe informacje, zobacz Trenowanie modeli za pomocą widoków funkcji.

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

Wykonuje wnioskowanie wsadowe w trybie offline z automatycznym wyszukiwaniem funkcji. Używa metadanych funkcji przechowywanych w modelu do obliczania prawidłowych funkcji w punkcie w czasie, zapewniając spójność z trenowania.

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

Ramka danych wejściowych musi zawierać kolumny jednostki i czasowników używane podczas trenowania. Funkcje są automatycznie obliczane na podstawie danych źródłowych.

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

Okna czasu

Widoki cech obsługują cztery typy okien, aby kontrolować zachowanie retrospektywy dla agregacji opartych na oknach czasowych. Dostępne typy okien zależą od źródła funkcji:

  • Funkcje źródłowe streamingu mogą korzystać z okien rolowanych i piłozębnych.

  • Funkcje zbiorowe mogą korzystać z okien z ruchomością, przewracaniem i przesuwaniem się.

  • Okna stopniowe wyglądają wstecz od czasu zdarzenia. Czas trwania i opóźnienie są jawnie zdefiniowane.

  • Okna ustalane są stałe, nienakładające się okna czasowe. Każdy punkt danych należy do dokładnie jednego okna.

  • Okna przesuwne to nakładające się, kroczące okna czasowe z konfigurowalnym interwałem przesunięcia.

  • Okna piłozębne utrzymują długie okno retroaktywne świeże nad strumieniem strumieniowym, korzystając z hybrydowej partii i ścieżki strumieniowania. Zobacz okno Sawtooth.

Poniższa ilustracja przedstawia typy okien typu obrotowego, ślizgowego, rolowanego i ząbkowego.

Przewracające się, ślizgające się, toczące się i ząbkowate okna wsteczne.

Okno stopniowe

Note

RollingWindow wcześniej nosił nazwę ContinuousWindow. W przypadku migracji z wcześniejszej wersji zestawu SDK należy odpowiednio zaktualizować importy.

Okna stopniowe są up-to—agregacje daty i czasu rzeczywistego, zwykle używane za pośrednictwem danych przesyłanych strumieniowo. W potokach przesyłania strumieniowego okno stopniowe emituje nowy wiersz tylko wtedy, gdy zawartość okna o stałej długości zmienia się, na przykład gdy zdarzenie wchodzi lub opuszcza. Gdy funkcja okna kroczącego jest używana w potokach trenowania, dokładne obliczenie funkcji punktu w czasie jest wykonywane na danych źródłowych przy użyciu czasu trwania okna o stałej długości bezpośrednio poprzedzających znacznik czasu określonego zdarzenia. Pomaga to zapobiec niesymetryczności danych w trybie online lub wycieku danych w trybie offline. Cechy w czasie T agregują zdarzenia z przedziału [T − długość trwania, T).

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

W poniższej tabeli wymieniono parametry okna kroczącego. Czasy rozpoczęcia i zakończenia okna są oparte na następujących parametrach:

  • Godzina rozpoczęcia: evaluation_time - window_duration - delay (włącznie)
  • Godzina zakończenia: evaluation_time - delay (wyłączność)
Parameter Ograniczenia
delay (opcjonalny) Musi być ≥ 0 (przesuwa okno do tyłu w czasie od znacznika czasu oceny). Służy delay do uwzględnienia dowolnego opóźnienia systemu między utworzeniem zdarzenia a sygnaturą czasową zdarzenia, aby zapobiec wyciekowi przyszłych zdarzeń do zestawów danych szkoleniowych. Jeśli na przykład między czasem utworzenia zdarzeń a tym zdarzeniem zostanie ostatecznie umieszczone w tabeli źródłowej, w której przypisano znacznik czasu, opóźnienie wynosiłoby timedelta(minutes=1)na przykład .
window_duration Musi mieć wartość > 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))

Zdefiniuj okno stopniowe z opóźnieniem przy użyciu poniższego kodu.

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

Przykłady okien kroczących

  • window_duration=timedelta(days=7): To utworzy 7-dniowe okno retrospektywne kończące się w momencie bieżącej oceny. W przypadku wydarzenia o godzinie 14:00 w dniu 7, obejmuje to wszystkie wydarzenia od 23:00 w dniu 0 do (ale nie w tym) 2:00 w dniu 7.

  • window_duration=timedelta(hours=1), delay=timedelta(minutes=30): Utworzy to 1-godzinne okno retrospektywne kończące się 30 minut przed czasem oceny. W przypadku wydarzenia o godzinie 15:00 obejmuje to wszystkie zdarzenia od 13:30 do (ale nie obejmuje) 2:30. Jest to przydatne do uwzględnienia opóźnień pozyskiwania danych.

Okno skokowe

W przypadku funkcji zdefiniowanych przy użyciu niezachodzących okien stałoczasowych, agregacje są obliczane w wstępnie określonym oknie o stałej długości, które przesuwa się według zdefiniowanego interwału, generując nienakładające się okna, które w pełni podzielają czas. W rezultacie każde zdarzenie w źródle przypisuje się dokładnie do jednego okna. Funkcje w momencie t agregują dane z okien kończących się na lub przed t (wyłącznie). System Windows zaczyna się od epoki systemu Unix.

class TumblingWindow(TimeWindow):
    window_duration: datetime.timedelta

W poniższej tabeli wymieniono parametry okna przesuwnego.

Parameter Ograniczenia
window_duration Musi mieć wartość > 0
from databricks.feature_engineering.entities import TumblingWindow
from datetime import timedelta

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

Przykład okna przeskakującego

  • window_duration=timedelta(days=5): Powoduje to utworzenie wstępnie określonych okien o stałej długości 5 dni. Przykład: okno nr 1 obejmuje od 0 do dnia 4, okno #2 obejmuje od 5 do dnia 9, okno #3 obejmuje dzień 10 do 14 dnia itd. W szczególności okno nr 1 zawiera wszystkie zdarzenia ze znacznikami czasu rozpoczynającymi się od 00:00:00.00 dnia 0 do (ale z wyłączeniem) wszystkich zdarzeń ze znacznikami 00:00:00.00 czasu w dniu 5. Każde zdarzenie należy do dokładnie jednego okna.

Okno przesuwane

Dla cech definiowanych za pomocą okien przesuwanych agregacje są obliczane nad oknem przesuwającym się o przedział przesuwania. Przesuwne okno może mieć stały czas trwania lub cały okres życia. Okna o stałym czasie trwania nakładają się, więc każde zdarzenie źródłowe może przyczynić się do agregacji cech dla wielu okien. Okno życiowe obejmuje wszystkie zdarzenia źródłowe przed końcem okna. Funkcje w momencie t agregują dane z okien kończących się na lub przed t (wyłącznie). Windows jest dostosowany do ery Uniksa.

class SlidingWindow(TimeWindow):
    window_duration: Optional[datetime.timedelta]
    slide_duration: datetime.timedelta

W poniższej tabeli wymieniono parametry okna przesuwanego.

Parameter Ograniczenia
window_duration Musi być dodatnia dla okna o stałej trwaniu. Ustaw na None okno na całe życie.
slide_duration Musi być dodatnia. Dla okna o stałej długości musi być również krótsza niż 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)
)

Przykład okna przesuwanego

  • window_duration=timedelta(days=5), slide_duration=timedelta(days=1): Powoduje to nakładanie się 5-dniowych okien, które przesuwają się o 1 dzień każdego dnia. Przykład: okno #1 obejmuje od 0 do 4 dnia, okno #2 obejmuje dzień 1 do dnia 5, okno #3 obejmuje dzień 2 do 6. dnia itd. Każde okno zawiera zdarzenia od 00:00:00.00 dnia rozpoczęcia do dnia końcowego, z wyłączeniem 00:00:00.00. Ponieważ okna nakładają się na siebie, pojedyncze zdarzenie może należeć do wielu okien (w tym przykładzie każde zdarzenie należy do maksymalnie 5 różnych okien).

Okno życia

Ustaw window_duration=None tak, aby tworzyć okno na całe życie. Na każdej granicy slajdu cecha agreguje wszystkie zdarzenia źródłowe dla jednostki z znacznikami czasowymi wyprzedzającymi tę granicę. Na przykład jednodniowy slajd generuje wartość skumulowaną raz dziennie.

Okna życia są obsługiwane wyłącznie przez SlidingWindow. RollingWindow i TumblingWindow wymagają skończonego window_duration.

Note

Okna dożywotnie wymagają databricks-feature-engineering wersji klienta wspierającej window_duration=None włączenie przestrzeni roboczej. Wcześniejsze wersje klienckie nie obsługują tej składni.

from datetime import timedelta
from databricks.feature_engineering.entities import (
    AggregationFunction,
    DeltaTableSource,
    Feature,
    SlidingWindow,
    Sum,
)

lifetime_spend = Feature(
    source=DeltaTableSource(
        catalog_name="main",
        schema_name="store",
        table_name="transactions",
    ),
    entity=["user_id"],
    timeseries_column="transaction_time",
    function=AggregationFunction(
        Sum(input="amount"),
        SlidingWindow(
            window_duration=None,
            slide_duration=timedelta(days=1),
        ),
    ),
    name="lifetime_spend",
)

Okno piłozębne

Ważna

SawtoothWindow jest w fazie beta.

Okno piłozębne to agregacja wspierająca bardzo świeże aktualizacje ostatnich wydarzeń, wraz z codziennym zagęszczaniem danych historycznych. Jego boczna (starsza) krawędź przesuwa się w stałych, codziennych krokach, podczas gdy krawędź prowadząca pozostaje na bieżąco z najnowszymi wydarzeniami, więc efektywna długość okna "przesuwa się" w ciągu każdego dnia. Większość okna jest obsługiwana z danych w tabeli pobierania strumienia, a tylko dwa najnowsze dni pochodzą z transmisji na żywo. To kompromis, który efektywnie oblicza długoterminowe okna (skalujące się do lat), jednocześnie pozostając responsywnym na świeże aktualizacje.

Okno piłozębne: przednia krawędź śledzi najnowsze wydarzenia, podczas gdy krawędź spływowa przesuwa się codziennie, więc zasłonięte okno

Okna piłozębne są realizowane na hybrydowej ścieżce partii i strumienia. Potok wsadowy utrzymuje większość okna, podczas gdy potok strumieniowy utrzymuje najnowsze dane świeże w czasie rzeczywistym. Oba są łączone podczas odczytu, więc dla modelu lub obsługującego konsumenta jest to jedna funkcja.

Ponieważ historyczna część okna jest obliczana przez pipeline wsadowy, funkcja piłowata jest gotowa do obsługi krótko po rozpoczęciu materializacji, nawet gdy okno obejmuje miesiące lub lata. Okno ruchome jest ukończone dopiero po upływie pełnego czasu trwania okna. Minimum window_duration musi być dłuższe niż dwa dni (wymuszona dolna granica), ale okna piłozębne zaleca się na okresy dłuższe niż 7 dni; dla krótszych okien należy użyć okna rolowanego .

Note

Kształt piłozębny opiera się na historii, która już istnieje. Tabela pobierania strumienia musi zawierać dane obejmujące co najmniej cały czas trwania okna, inaczej obliczene okno jest niepełne. Przed upływem 2 pełnych dni funkcja odzwierciedla jedynie dotychczas zmaterializowane dane. Nie zaleca się podawania filmu w trakcie produkcji przed upływem 2 pełnych dni. Agregacja nad pustym oknem daje 0 dla Sum i Count, oraz null dla Avg, Min, oraz .Max

Aby sprawdzić, czy funkcja piłowata jest gotowa, otwórz widok funkcji w Eksploratorze katalogu. W sekcji Materializowane cechy uzupełnianie wsad kończy się, gdy ostatni czas materializacji funkcji minie i jej status pokazuje sukces. Część strumieniowa jest realizowana przez deklaratywny pipeline Lakeflow. Po przejściu weryfikacji Feature View funkcja zmaterializowana łączy się z tym pipeline'em, gdzie możesz monitorować jej status działania.

Okna piłozębne wymagają a StreamSource a materializują się z .StreamingMode

class SawtoothWindow(TimeWindow):
    window_duration: datetime.timedelta

Krawędzie okna piłowatego poruszają się inaczej niż w szybie rolnej: krawędź prowadząca śledzi najnowsze zdarzenie, podczas gdy krawędź spoczynkowa przesuwa się raz dziennie, a nie ciągle. Każdego dnia, o stałym progu 18:00 UTC, krawędź tylna przesuwa się do granicy UTC-północ tego dnia. W efekcie efektywne okno jest nieco dłuższe niż window_duration i rośnie w ciągu dnia, zanim w następnym momencie wraca do normy. Szkolenia i obsługa korzystają z tego samego progu 18:00 UTC, więc szkolenia offline i obsługa online pozostają spójne.

Parameter Ograniczenia
window_duration Musi być dłużej niż dwa dni. Dozwolony jest czas trwania niebędący liczbą dni (na przykład timedelta(days=3, minutes=15)), ale okno jest nadal aktualizowane z dzienną szczegółowością.

Okna piłozębne obsługują Sumfunkcje , Avg, Count, Min, oraz Maxagregacji.

Przykład okna piłowatego

Poniższy przykład pokazuje 7-dniowy bilans transakcji użytkownika. Krawędź prowadząca śledzi bieżące zdarzenie, podczas gdy krawędź spływowa przesuwa się do przodu dzień po dniu. Jeśli chodzi o wydarzenia 10 marca, okno sięga około 3 marca. W miarę upływu 10 marca krawędź natarcia ciągle się przesuwa, podczas gdy krawędź spoczynkowa utrzymuje, więc zakryty rozpięty rozpięty rośnie. Następnie, na początku 11 marca, krawędź podkładu przesuwa się do około 4 marca. Skuteczne okno to zawsze nieco dłużej niż siedem dni. Dwa ostatnie dni są serwowane z transmisji na żywo, a wcześniejsze z stołu wlewającego strumienia.

from databricks.feature_engineering.entities import SawtoothWindow
from datetime import timedelta

# 7-day window kept continuously fresh with streaming data
window = SawtoothWindow(window_duration=timedelta(days=7))

Ograniczenia okna piłowatego

  • Parametr delay nie jest obsługiwany.
  • Funkcje agregacji inne niż , , , , i Max nie są obsługiwane (na przykład First, Last, ApproxCountDistinct, oraz funkcje odchylenia standardowego i wariancji). MinCountAvgSum
  • Okna piłozębne wymagają .StreamSource A DeltaTableSource nie jest wspierany.

Wyzwalacze materializacji

Wyzwala kontrolkę po uruchomieniu potoku materializacji. Typ wyzwalacza zależy od typu funkcji.

CronSchedule

Służy CronSchedule do agregacji funkcji (AggregationFunction). Potok jest uruchamiany zgodnie z ustalonym harmonogramem zdefiniowanym przez wyrażenie kronu kwarcowego.

from databricks.feature_engineering.entities import CronSchedule

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

TableTrigger

Zastosowanie TableTrigger dla ColumnSelection cech lub cech agregacji (AggregationFunction) wspieranych przez .DeltaTableSource Potok jest uruchamiany za każdym razem, gdy nadrzędna tabela delty otrzymuje nowe zatwierdzenie.

W przypadku cech agregacji pipeline jest ograniczany, aby nie uruchamiał się przy każdym commite. Potok uruchamiany jest co najwyżej raz na połowę długości okna funkcji, ale nigdy częściej niż co 5 minut. Na przykład funkcja z 1-godzinnym oknem obrotowym uruchamia się co najwyżej raz na 30 minut, lub funkcja z oknem 8-godzinnym działa co najwyżej raz na 4 godziny. Minimalna granica 5 minut obowiązuje, gdy połowa okna jest krótsza niż to, więc okna trwające 10 minut lub mniej są uruchamiane maksymalnie raz na 5 minut. Funkcje agregacji, których okno to mniej niż 5 minut, nie mogą używać TableTrigger, używając zamiast tego wyzwalacza strumieniowego.

from databricks.feature_engineering.entities import TableTrigger

trigger = TableTrigger()

StreamingMode

Służy StreamingMode do obsługi funkcji wspieranych przez element StreamSource. Potok jest uruchamiany jako potok ciągłego przesyłania strumieniowego.

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

Wybieranie wyzwalacza

Każda funkcja korzysta z jednego wyzwalacza; Opcje według typu funkcji to:

Typ funkcji Trigger Po uruchomieniu
Agregacja (AggregationFunction) z DeltaTableSource CronSchedule Zgodnie z ustalonym harmonogramem cron
Agregacja (AggregationFunction) z DeltaTableSource TableTrigger W każdym zatwierdzeniu tabeli źródłowej
ColumnSelection (z DeltaTableSource) TableTrigger W każdym zatwierdzeniu tabeli źródłowej
Funkcje z StreamSource StreamingMode Ciągłe przesyłanie strumieniowe

Nie można materializować funkcji, które wymagają różnych typów wyzwalaczy w jednym materialize_features wywołaniu. Zamiast tego wydaj oddzielne wywołania.