Megjegyzés
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhat bejelentkezni vagy módosítani a címtárat.
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhatja módosítani a címtárat.
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
Featureobjektumokat, majd használja aregister_featurea tartósításukhoz a Unity Catalogban később. A helyileg létrehozott funkciók a regisztráció előtt használhatókcreate_training_set.
- A Unity Catalog olyan funkcióobjektumainak definiálására használható
-
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.
- A gépi tanulás időponthoz kötött összesített funkcióinak kiszámítására használható
-
Funkció materializálása és a munkafolyamat kiszolgálása
- Miután egy funkciót
create_featuredefiniált, vagy aget_featurehasználatával lekérte, a funkciót vagy a funkciókészletet anyagíthatjamaterialize_featuresegy 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_setmaterializált nézetet offline kötegelt betanítási adatkészlet előkészítéséhez.
- Miután egy funkciót
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
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ákoffline_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_sqlnem 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
StreamSourceelemet, ha az elsődleges szempont az egy másodpercnél rövidebb frissesség.StreamSourceA 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á egyStreamSourcehaszná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)
DeltaTableSourcefeltö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
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_featureshí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 megentityaz 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_sethasznált funkciódefiníciók között. - A betanítási
labeladathalmaz oszlopaként használt oszlopnév nem szerepelhet az s definiálásáhozFeaturehasznált forrástáblákban. - Az API támogatja
create_featurea 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
DATEvagyTIMESTAMP. -
RequestSourcecsak azScalarDataType(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. -
RequestSourcenem támogatja az aggregációs függvényeket vagy az időablakokat. CsakColumnSelectionfü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_batchkiszolgá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.