Funkciónézetek

Important

Ez a funkció nyilvános előzetes verzióban van. A munkaterület rendszergazdái az Előnézetek lapon szabályozhatják a funkcióhoz való hozzáférést. Lásd: Az Azure Databricks előzetes verziójának kezelése.

A funkciónézetek lehetővé teszik az adatforrásokból származó funkciók definiálására és kiszámítására. A funkciók számos különböző forrásból (Delta-tábla, Kafka Stream és kérelem-idő adatok) és számításokból (időablakos összesítések, egyszerű oszlopkijelölések stb.) határozhatók meg. Ez az útmutató a következő munkafolyamatokat ismerteti:

  • Szolgáltatásfejlesztési munkafolyamat
    • A Unity Catalog olyan funkcióobjektumainak definiálására használható create_feature , amelyek a modell betanításában és a munkafolyamatok kiszolgálásában használhatók.
    • Másik lehetőségként hozza létre helyileg a Feature objektumokat, majd használja a register_feature a tartósításukhoz a Unity Catalogban később. A helyileg létrehozott funkciók a regisztráció előtt használhatók create_training_set .
  • Modell betanítási munkafolyamata
    • A gépi tanulás időponthoz kötött összesített funkcióinak kiszámítására használható create_training_set . A funkciónézetekkel való betanítás részletes dokumentációját a Modellek betanítása funkciónézetekkel című témakörben találja.
  • Funkció materializálása és a munkafolyamat kiszolgálása
    • Miután egy funkciót create_feature definiált, vagy a get_feature használatával lekérte, a funkciót vagy a funkciókészletet anyagíthatja materialize_features egy offline tárhelyre a hatékony újrafelhasználás érdekében, vagy egy online tárhelyre az online kiszolgálás céljából.
    • Használja a create_training_set materializált nézetet offline kötegelt betanítási adatkészlet előkészítéséhez.

Az API részleteiért tekintse meg a Funkciónézetek API-referenciáját.

Követelmények

  • Szerver nélküli számítás vagy klasszikus számítási fürt, amelyen a Databricks Runtime 17.0 ML vagy újabb verziója fut.

  • Telepítenie kell az egyéni Python-csomagot. Futtassa a következő kódsorokat minden alkalommal, amikor jegyzetfüzetet futtat:

    %pip install databricks-feature-engineering>=0.16.0
    dbutils.library.restartPython()
    

Gyors kezdés: példa

A futtatható kezdő lépések jegyzetfüzetéhez lásd a példajegyzetfüzetet.

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

Példajegyzetfüzet

Feature Views – gyorsútmutató notebook

Jegyzetfüzet szerezz

Streamelési funkciók

Használj streaming funkciókat, amikor a funkciók értékeit folyamatosan kell frissíteni, nem pedig köteses ütemezésben. A streaming és batch funkciók ugyanazokat Feature a konstruktorokat, aggregációs funkciókat, valamint a képzési és szolgáltatási munkafolyamatokat használják.

A streaming funkciók nem jelennek meg offline áruházban. A tréninghez és a batch következtetéshez a Databricks a forrásból származó jellemzők értékeit számítja ki.

A streaming funkciók a következő követelményeket követik:

  • Be kell adnod egy online_config. A streaming funkciók nem támogatják offline_config.
  • Nem lehet egyetlen hívásban kombinálni a streaming és a csomagos funkciókat materialize_features . Külön hívást csinálj minden trigger típushoz.
  • transformation_sql nem támogatott a streaming funkciókhoz.
  • A streaming materializáció csak azokat a rekordokat dolgozza fel, amelyek a csővezeték kezdete után érkeznek, és nem tölti ki a történelmi rekordokat. A görgető ablak aggregációk csak az első teljes adatablak érkezése után adnak teljes eredményeket.

Streaming funkció forrásának kiválasztása

Válassz forrást a frissességi igényeid és a meglévő beviteli beállításod alapján:

  • Használjon StreamSource elemet, ha az elsődleges szempont az egy másodpercnél rövidebb frissesség. StreamSource A funkciók 200 milliszekundumos végponttól végpontig tartó késleltetést biztosítanak a P99-ben. Először hozz létre egy adatfolyamot, majd hivatkozz rá egy StreamSource használatával. A streamforrások bemenetként támogatják a Kafkát, és automatikusan karbantartanak egy adatbetöltési Delta-táblát, amely a betanításhoz az adatok historikus másolataként szolgál.
  • Használd a(z) DeltaTableSource-t, ha már van egy alacsony késleltetésű adatbeviteli útvonalad egy Delta-táblába. Várj frissességet, akár több tíz másodperc alatt.
  • Használd a Zerobust a(z) DeltaTableSource feltöltéséhez, ha még nem rendelkezel alacsony késleltetésű betöltési útvonallal. A Zerobus adatbeolvasása néhány tíz másodpercet vesz igénybe, így a jellemzők frissessége várhatóan egy percen belüli lesz.

StreamSource segítségével definiálni egy streaming funkciót

Egy StreamSource streamre a háromrészes neve (catalog.schema.stream_name) alapján hivatkozik. A stream nem a Unity Catalog hozzáférés-szabályozható objektuma, de egy Unity Catalog-sémához tartozik, és a hozzáférést a stream adatbeviteli táblája szabályozza. Az entitás-, idősor- és függvénydefiníciókban szereplő oszlophivatkozásokat a value. vagy a key. előtaggal kell ellátni annak jelzésére, hogy a Kafka-üzenet melyik részéből kell olvasni. A beágyazott mezők pont jelöléssel (például value.user.address.city) támogatottak.

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

Definiáljunk egy streaming funkciót egy DeltaTableSource segítségével

Egy DeltaTableSource elemen definiált jellemző streaming jellemzőként való megvalósításához add át a StreamingMode elemet triggerként a materialize_features elemnek. A jellemződefiníció ugyanazokat az API-kat használja, mint egy DeltaTableSource által támogatott kötegelt jellemző. A delta tábla források támogatják az aggregációs és oszlopkiválasztási funkciókat.

A forrás Delta táblának be kell kapcsolnia az adatfolyam megváltoztatását (CDF) a beállítással delta.enableChangeDataFeed=true.

A következő példa egy aggregációs funkciót definiál és materializál egy Delta táblaforrással.

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

Használj egy Delta táblát, amelyet Zerobus tölt fel

Egy Zerobus által feltöltött Delta tábla streaming funkcióforrásként szolgálhat. A Zerobus nem állítja be automatikusan a(z) delta.enableChangeDataFeed=true elemet. Ezt a tulajdonságot manuálisan kell beállítani a cél Delta táblán, mielőtt streaming funkcióként használnád.

Szűrő feltételek a streaming forrásokon

Használd a(z) filter_condition elemet a sorok aggregálás előtti szűrésére, akár StreamSource, akár DeltaTableSource esetén.

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

Rovatválogatás streaming forrásokból

ColumnSelection a funkciók továbbítják a legfrissebb értéket minden entitáskulcshoz aggregáció nélkül. A betanítás során a jellemzőértékek megfelelnek az adott időpontra vonatkozó pontosságnak.

Az oszlopkiválasztási funkcióknak nincs TTL-e. Egy kiválasztott érték eltávolításához az online áruházból a forrásnak null értéket kell kibocsátania a kiválasztott oszlopra.

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

Hozzáférés a beágyazott mezőkhöz egy StreamSource-ból

Egy StreamSourceesetén a beágyazott JSON mezőket dot jelöléssel lehet elérni (például value.nested_field.amount). Kiszolgáláskor a kérés adattörzse és a válasz a levélcsomópontok neveit használja (például amount a value.amount helyett). A levélcsomópontok neveinek egyedinek kell lenniük az összes entitás-, idősorozat- és jellemzőnév között egy modellen vagy Feature Specen belül, mivel a szolgáltatási végpont a levélneveket használja az értékek irányítására.

Streamelési funkciók időablakai

A streamelési funkciók csak RollingWindow az aggregációk esetében támogatottak. A gördülő ablakok folyamatosan újraszámítódnak a legfrissebb adatok alapján, ami összhangban van a streamingforrások valós idejű jellegével. TumblingWindow és SlidingWindow rögzített múltbeli időintervallumokra vonatkozó kötegelt számításokra szolgál.

Streamelési funkciók – példajegyzetfüzet

Streaming jellemzőnézetek – gyorsútmutató jegyzetfüzet

Jegyzetfüzet szerezz

Modell betanítása és következtetése

A modellek betanításával és a kötegelt következtetés Feature Views használatával történő futtatásával kapcsolatban, beleértve a(z) log_model(), score_batch() és create_training_set() elemet, lásd a(z) Modellek betanítása Feature Views használatával című részt.

Funkció materializálása

Miután definiálja a funkciókat, materializálhatja vagy tárolhatja őket offline vagy online tárolókban, ezáltal hatékonyan újra felhasználhatók a képzési és szolgáltatási munkafolyamatokban. A funkciók megvalósítása után a modelleket processzormodellek kiszolgálásával is kiszolgálhatja. A részletekért lásd: Jellemzőnézetek materializálása.

Bevált gyakorlatok

Funkcióelnevezés

  • Használjon leíró neveket az üzleti szempontból kritikus funkciókhoz.
  • Kövesse a csapatok egységes elnevezési konvencióit.
  • Használjon automatikusan létrehozott neveket a funkciók fejlesztésének megkezdésekor.

Időablakok

  • Az ablakhatárok igazítása üzleti ciklusokhoz (napi, heti).
  • A rövidebb ablakok rögzítik a legutóbbi trendeket, de zajosak lehetnek. A hosszabb ablakok stabilabb funkcióeloszlásokat eredményeznek, de előfordulhat, hogy nem jelennek meg a legutóbbi viselkedési változások. Azt válassza, amelyik a használati esetében a legjobban tükrözi, hogy milyen gyorsan változik a mögöttes jel. Egy 7 napos ablak például kisimítja a napi ingadozásokat, és konzisztens modellbemeneteket hoz létre, míg egy 1 órás ablak gyorsan reagál a viselkedési változásokra, de olyan varianciát eredményezhet, amely rontja a modell teljesítményét. Ha a modell pontossága csökken, amikor az eloszlás eltolódik, használjon hosszabb időt a bemenetek stabilizálásához.
  • Az ejtő és tolóablakok méretezhetőbbek, mint a gördülő (folyamatos) ablakok. A legtöbb használati esethez használjon tolóablakokat.

Performance

  • A data scanningok minimalizálása érdekében egyetlen materialize_features hívásban materializálja az attribútumokat azonos adatforrásból.
  • Ugyanazt a részletességet (például az 1 órás vagy az összes 1 napos diaidőt) használja ugyanazon adatforrás funkcióihoz, hogy jobb csoportosítást lehessen lehetővé tenni a materializálás során.

Entitásoszlopok és szűrőfeltételek

Használja ezt a döntési útmutatót, amikor ugyanazon forrástáblából származó funkciókkal dolgozik:

Használja a entity (a create_feature esetén), amikor különböző összesítési szintekre van szüksége:

  • Ügyfélszintű funkciók (ügyfélenként egy sor): entity=["customer_id"]
  • Ügyfél-kereskedői funkciók (ügyfélenként több sor): entity=["customer_id", "merchant_id"]
  • A különböző összesítési szintek azonosak DeltaTableSourcelehetnek: különböző értékeket adhat meg entity az egyes funkciódefiníciókhoz

Használja a filter_condition (a DeltaTableSourceen), amikor a sorokat ugyanazon az összesítési szinten kell szűrnie:

  • Csak nagy értékű tranzakciók: filter_condition="amount > 100" (továbbra is összesítve ügyfélenként)
  • Csak befejezett rendelések: filter_condition="status = 'completed'" (ügyfélenként továbbra is összesítve)

Hüvelykujjszabály: Ha a módosítás entitásértékenként eltérő számú sort eredményezne, használjon különböző entity értékeket a funkciódefiníciókban. Ha csak szűri, hogy mely sorok járulnak hozzá ugyanahhoz az összesítéshez, használja filter_condition a forrást.

Gyakori minták

Ügyfél-elemzés

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

Trendelemzés

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

Szezonális minták

# 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

  • Az entitások és idősorok oszlopainak neveinek egyeznie kell a betanítási (címkézett) adatkészlet és az API-ban create_training_set használt funkciódefiníciók között.
  • A betanítási label adathalmaz oszlopaként használt oszlopnév nem szerepelhet az s definiálásához Featurehasznált forrástáblákban.
  • Az API támogatja create_feature a függvények (UDAF-k) korlátozott listáját. Lásd : Támogatott függvények.
  • Az entitásoszlopok nem lehetnek típusúak DATE vagy TIMESTAMP.
  • RequestSource csak az ScalarDataType (INTEGER, FLOAT, BOOLEAN, STRING, DOUBLE, LONG, TIMESTAMP, DATE, SHORT) függvényben definiált skaláris adattípusokat támogatja. Az olyan összetett típusok, mint a tömbök, a térképek és a szerkezetek nem támogatottak.
  • RequestSource nem támogatja az aggregációs függvényeket vagy az időablakokat. Csak ColumnSelection függvények használhatók.
  • Az entitásoszlopok neveinek, az időbélyegek oszlopneveinek és a kérelemfunkció-oszlopneveknek globálisan egyedinek kell lenniük a betanítási készletben vagy a végpontot kiszolgáló összes forrásban.
  • score_batch kiszolgáló nélküli számításon előfordulhat, hogy nem sikerül. Ezt úgy kerülheti meg, hogy egy klasszikus számítási fürtöt használ, amely Databricks Runtime 17.0 ML vagy újabb verziót futtat.

A materializálásra vonatkozó korlátozásokért lásd a Korlátozások című témakört.