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 Stream egy külső streamelési adatforrást jelöl, például az Apache Kafkát. A streamek a kapcsolat részleteit, a hitelesítést, a sémákat és a betöltési konfigurációt tárolják. A stream létrehozása után a funkciónézet-definíciók használatával hivatkozhat rá valós idejű streamelési funkciók létrehozásához.
A streamek háromrészes névvel (catalog.schema.stream_name) rendelkeznek. A Streamhez való hozzáférést a hozzá tartozó betöltési táblázat szabályozza. Részletekért lásd a betöltési és a visszatöltési adatokat.
Requirements
- Jegyzetfüzet-parancsok futtatásához: kiszolgáló nélküli számítási környezet vagy Databricks Runtime 17.0 ML vagy újabb verziót futtató klasszikus számítási fürt.
- A
feature-engineering-clientPython csomag 0.16.0-s vagy újabb verzióját kell telepíteni.
Stream létrehozása
Új Stream létrehozásához használható create_stream() . A Stream négy konfigurációs összetevőt igényel:
- Forráskonfiguráció: Megadja a streamelési platformot (például Kafka) és a forrásspecifikus részleteket (például a Kafka témakör-előfizetését).
- Kapcsolatkonfiguráció: Megadja, hogyan csatlakozhat és hitelesíthet a streamelési platformhoz, beleértve a bootstrap-kiszolgálókat és a hitelesítő adatokat.
- Sémakonfiguráció: Meghatározza az üzenetkulcsok és -értékek szerkezetét.
- Betöltési konfiguráció: Megadja a streamadatok betöltésének helyét és módját. Részletekért lásd a betöltési és a visszatöltési adatokat.
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
KafkaStreamConfig,
KafkaSubscriptionMode,
StreamConnectionConfig,
DirectSchemas,
SchemaConfig,
IngestionConfig,
IngestionDestination,
StreamBackfillSource,
)
client = FeatureEngineeringClient()
stream = client.create_stream(
name="my_catalog.my_schema.my_stream",
source_config=KafkaStreamConfig(
subscription_mode=KafkaSubscriptionMode(subscribe="events-topic"),
),
connection_config=StreamConnectionConfig(
uc_connection_name="my-kafka-connection"
),
schema_config=DirectSchemas(
payload_schema=SchemaConfig(
json_schema=(
'{'
' "type": "object",'
' "properties": {'
' "transaction_id": {"type": "string"},'
' "user_id": {"type": "string"},'
' "amount": {"type": "number"},'
' "event_time": {"type": "string", "format": "date-time"}'
' }'
'}'
)
),
),
ingestion_config=IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
),
)
Csatlakozás streamforrásokhoz
A streamelési funkciók definiálása előtt csatlakozzon és tesztelje a Stream Lakeflow-folyamat kapcsolatát a Kafka-közvetítővel. Lásd : Streamelés kiszolgáló nélküli számításon és Csatlakozás az Apache Kafkához.
Az AWS által felügyelt streamelés (Amazon MSK) esetében lásd az Amazon MSK kiszolgáló nélküli privát kapcsolatát. A Kafka-hitelesítési lehetőségekről további információt a Hitelesítés című témakörben talál.
Authentication
Unity Catalog-kapcsolat (ajánlott)
A Kafka-fürthöz való hitelesítéshez használjon Unity Catalog-kapcsolatot. Ez a felügyelt hitelesítés ajánlott megközelítése. Kapcsolat létrehozásához lásd: Kapcsolat létrehozása.
connection_config = StreamConnectionConfig(
uc_connection_name="my-kafka-connection"
)
Közvetlen mTLS
Közvetlen mTLS-hitelesítéshez adja meg a Unity Catalog kötetén tárolt keystore- és truststore-fájlokat, a jelszavakra pedig Databricks secret scope-okon keresztül hivatkozzon. A Kafkával történő SSL-hitelesítéssel kapcsolatos további információkért lásd: Ssl használata Azure Databricks a Kafkához való csatlakozáshoz.
from databricks.feature_engineering.entities import (
DirectMtlsConfig,
MtlsConfig,
SecretScopeReference,
)
connection_config = DirectMtlsConfig(
bootstrap_servers="broker1:9092,broker2:9092",
mtls_config=MtlsConfig(
keystore_location="/Volumes/my_catalog/my_schema/my_volume/keystore.jks",
keystore_password_ref=SecretScopeReference(
scope="my_scope", key="keystore_password"
),
key_password_ref=SecretScopeReference(
scope="my_scope", key="key_password"
),
truststore_location="/Volumes/my_catalog/my_schema/my_volume/truststore.jks",
truststore_password_ref=SecretScopeReference(
scope="my_scope", key="truststore_password"
),
),
)
SASL
A SASL-hitelesítés (SASL/SCRAM és SASL/PLAIN) nem támogatott az előzetes verzióban.
Előfizetési módok
Az előfizetési mód azt határozza meg, hogy a Stream hogyan választja ki a Kafka-témaköröket, amelyekből használni szeretné. Három mód támogatott:
| Üzemmód | Description | Példa |
|---|---|---|
subscribe |
Témakörnevek vesszővel tagolt listája | KafkaSubscriptionMode(subscribe="topic1,topic2") |
subscribe_pattern |
Java regex mintára illeszkedő témanevek | KafkaSubscriptionMode(subscribe_pattern="events-.*") |
assign |
A témakör–partíció hozzárendeléseket megadó JSON | KafkaSubscriptionMode(assign='{"my-topic": [0, 1, 2]}') |
Sémakonfiguráció
Definiáljuk az üzenetkulcsok és értékek szerkezetét, hogy a felvételi és funkciódefiníciók az egyes mezők olvashatók legyenek. A Kafka-források payload_schema esetében a Kafka üzenetértékének (a value Kafka kulcs-érték modelljében) felel meg, és key_schema a Kafka üzenetkulcsának felel meg. A payload_schema vagy a key_schema közül legalább az egyiket meg kell adni.
Mindegyik SchemaConfig három formátum egyikét fogadja el, amely megfelel a forrás üzeneteinek sorozatosításának: json_schema, avro_schema, vagy proto_schema. Ha egy kulcshoz vagy hasznos adathoz nincs séma megadva, a rendszer egyszerű sztringként kezeli.
JSON-séma
Biztosíts egy JSON séma stringet .json_schema
schema_config = DirectSchemas(
payload_schema=SchemaConfig(
json_schema=(
'{'
' "type": "object",'
' "properties": {'
' "user_id": {"type": "string"},'
' "amount": {"type": "number"},'
' "event_time": {"type": "string"}'
' }'
'}'
)
),
key_schema=SchemaConfig(
json_schema='{"type": "string"}'
),
)
Avro-séma
Biztosíts egy Avro séma láncot a avro_schema. Az avro logikai típusok támogatottak, beleértve timestamp-millis, date, és decimal.
schema_config = DirectSchemas(
payload_schema=SchemaConfig(
avro_schema=(
'{'
' "type": "record",'
' "name": "Event",'
' "fields": ['
' {"name": "user_id", "type": "string"},'
' {"name": "amount", "type": "double"},'
' {"name": "event_time",'
' "type": {"type": "long", "logicalType": "timestamp-millis"}}'
' ]'
'}'
)
),
)
Protobuf-séma
Adjon meg egy ProtoSchemaSpec to proto_schema címet a protokoll pufferek.proto forrásszövegével és a hasznos raher üzenet nevével. Import ProtoSchemaSpec innen databricks.feature_engineering.entities.
message_nameA teljes minősítésű üzenetnévnek kell lennie, beleértve a szövegben kijelentett címet package is (például com.example.Event, nem Event)..proto Mind a proto2, mind a proto3 szintaxisa támogatott.
google.protobuf.Timestamp és a skalár csomagolás típusok (StringValue, Int32Value, és így tovább) támogatottak, és importjaik automatikusan megoldódnak. Más ismert típusokat, mint Durationpéldául , Struct, és Any, elutasítják; ezeket az értékeket inkább támogatott skalárként vagy üzenetként kódoljuk. A fixed32 skaláris típusok fixed64 , valamint map a nem húros billentyűk szintén nem támogatottak.
from databricks.feature_engineering.entities import ProtoSchemaSpec
schema_config = DirectSchemas(
payload_schema=SchemaConfig(
proto_schema=ProtoSchemaSpec(
schema_text=(
'syntax = "proto3";\n'
'package com.example;\n'
'import "google/protobuf/timestamp.proto";\n'
'message Event {\n'
' string user_id = 1;\n'
' double amount = 2;\n'
' google.protobuf.Timestamp event_time = 3;\n'
'}'
),
message_name="com.example.Event",
)
),
)
Betöltés és visszatöltés
A ingestion_config paraméter konfigurálja a streamadatok rögzítését és tárolását betanításra és kiszolgálásra.
A Streamhez való hozzáférést a betöltési tábla szabályozza:
-
SELECTa betöltési táblában olvasási hozzáférést biztosít a Streamhez. -
MANAGEa betöltési táblában törlési hozzáférést biztosít.
A táblákra vonatkozó jogosultságokkal kapcsolatos további információkért lásd a Table és a Unity Catalog privileges reference című részt.
Adatbetöltési folyamat
Stream létrehozásakor a Databricks elindít egy felügyelt betöltési folyamatot, amely folyamatosan olvas üzeneteket a Kafka-témakörből, és egy Delta-táblába (a betöltési táblába) írja őket. Az adatfolyam a legfrissebb Kafka-offsettől indul, és folyamatosan fut, kizárólag a stream létrehozása után érkező új üzeneteket rögzítve. Ez a betöltési táblázat streamelési funkciókkal való betanításra szolgál. Egy stream törlésekor a betöltési folyamat és a betöltési tábla is törlődik.
Betöltési célhely
A ingestion_destination háromrészes Delta-tábla nevét adja meg, ahol a streamadatok meg vannak írva.
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
)
Adatbetöltési táblaséma
A betöltési táblázat az üzenetadatokat és a metaadatoszlopokat tartalmazza:
| Column | Típus | Description |
|---|---|---|
key |
Eltérő (ettől: key_schema) |
A Kafka üzenetkulcsa, a megadott séma szerint strukturálva. |
value |
Eltérő (ettől: payload_schema) |
A Kafka-üzenet értéke (hasznos adat), a megadott séma szerint strukturálva. |
stream_record_timestamp |
TIMESTAMP |
A rekord időbélyege. Előretöltött adatok esetén ez a Kafka broker beérkezési időbélyege. A visszatöltési adatok esetében ez az ügyfél által megadott. |
kafka_topic |
STRING |
Az a Kafka-témakör, amelyből a rekordot felhasználták. |
kafka_partition |
INT |
Az a Kafka-partíció, amelyből a rekordot beolvasták. |
kafka_offset |
LONG |
A rekord partíción belüli Kafka-offsetje. |
record_source |
STRING |
Vagy "stream" (feltöltés az élő Kafka-streamből), vagy "backfill" (a visszatöltési forrásból). |
Utólagos feltöltés forrása
Mivel az előretöltési folyamat a legújabb Kafka-eltolásból indul ki, nem rögzíti a stream létrehozása előtt létező üzeneteket. A betanítás előzményadat-lefedettségének biztosításához konfiguráljon egy opcionális háttérbetöltési forrást.
Ha egy visszatöltési forrás van konfigurálva, a Databricks lefuttat egy egyszeri MERGE INTO feladatot, amely a visszatöltési sorokat a betöltési táblába másolja a(z) record_source="backfill" használatával. A MERGE művelet csak akkor fut le, ha az átfedésvizsgáló megerősíti, hogy a backfill forrás és az előretöltési adatfolyam időbélyegei átfedésben vannak (lásd Átfedés a backfill és az élő adatfolyam adatai között). Ha az átfedési feltétel 2 napon belül nem teljesül, a MERGE mindenképpen fut, hogy elkerülje a határozatlan ideig történő blokkolást.
A visszatöltési táblának tartalmaznia kell egy stream_record_timestamp oszlopot, amelynek típusa TIMESTAMP, és UTC időzónájú. További Kafka-metaadat-oszlopok (kafka_topic, kafka_partition, kafka_offset) átadásra kerülnek, ha jelen vannak a backfillforrásban, ellenkező esetben pedig NULL értékre vannak állítva.
from databricks.feature_engineering.entities import StreamBackfillSource
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
backfill_source=StreamBackfillSource(
delta_table_name="my_catalog.my_schema.historical_events"
),
)
Átfedés a backfill és az élő stream adatai között
Mielőtt a BACKFILL és a betöltési tábla között futtatja a MERGE parancsot, egy átfedéses ellenőrzés összehasonlítja a két tábla időbélyegeit:
-
Visszatöltés maximális száma: A feltöltési forrásban megadott maximális érték
stream_record_timestamp. -
Betöltési minimum: A betöltési táblában lévő sorok minimális
stream_record_timestampszáma (record_source="stream").
A MERGE akkor megy végbe, ha a visszatöltés legkésőbbi időbélyege legalább 1 órával későbbi, mint a beviteli tábla legkorábbi időbélyege. Ez az átfedés biztosítja, hogy ne legyenek hézagok az adatbetöltési táblában. Ha az átfedési feltétel 2 napon belül nem teljesül, a MERGE mindenképpen fut, hogy elkerülje a határozatlan ideig történő blokkolást.
Mivel a betöltési folyamat a legújabb Kafka-eltolástól indul, csak a stream létrehozása után érkező üzeneteket rögzíti. A háttérbetöltési forrásnak olyan adatokat kell tartalmaznia, amelyek a betöltési időtartományra terjednek ki – nem csak a stream létrehozási idejéig.
Ha például 15:00-kor hoz létre egy adatfolyamot, a forward-fill folyamat 15:00-tól kezdve olvassa az üzeneteket. A feltöltési forrásnak legalább 16:00-ig (az előretöltési kezdés után 1 órával) időbélyeggel rendelkező adatokat kell tartalmaznia az átfedés ellenőrzéséhez. Ez azt jelenti, hogy 16:00 óra után frissítenie kell a backfill táblát, hogy az adatbeviteli táblában ne legyenek hiányok.
Deduplication
A deduplication_columns használatával adhatja meg az oszlopelérési utakat az ismétlődő sorok azonosításához a visszatöltési és az előretöltési adatfolyamok adatai közötti betöltés során. Pont jelölés használata beágyazott mezőkhöz (például "value.user_id").
Az adatok alapján válassza ki a deduplikációs oszlopokat:
- Ha a stream minden rekordja egyedi azonosítót (például) tartalmaz,
value.transaction_idhasználja ezt az oszlopot a deduplikációhoz. - Ha a háttértöltési forrás tartalmaz
kafka_partitionéskafka_offsetoszlopokat, azokkal egyedileg azonosíthatja az egyes rekordokat. - Ha nincs megadva deduplikációs oszlop, az alapértelmezett deduplikációs kulcs a
key,valueésstream_record_timestampteljes kombinációja. Ez nem ajánlott, mivel ez a szigorú feltételek egyeztetése könnyen duplikációkhoz vezethet.
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
deduplication_columns=["value.transaction_id"],
)
Streamek kezelése
Stream lekérése
stream = client.get_stream(name="my_catalog.my_schema.my_stream")
Adatfolyamok listázása
streams = client.list_streams(
catalog_name="my_catalog",
schema_name="my_schema",
max_results=50,
include_schemas=False,
)
Állítsa be a(z) include_schemas=True értéket úgy, hogy tartalmazza a séma teljes részleteit. A sémák nagyok lehetnek, és ez hosszú ideig futó műveletet eredményezhet. A sémák egyenkénti lekéréséhez használja a következőt get_stream: .
Stream törlése
A stream törlése törli a betöltési folyamatot és a betöltési táblát is.
Warning
A törölt streamre hivatkozó modellek és szolgáltatások többé nem férhetnek hozzá a mögöttes streamadatokhoz. Ha szüksége van ezekre az adatokra, de már nincs szüksége a streamre, a törlés előtt készítsen másolatot a betöltési tábláról.
client.delete_stream(name="my_catalog.my_schema.my_stream")
Példajegyzetfüzet
Az alábbi jegyzetfüzetben egy teljes körű példa látható, amely létrehoz egy Streamet, definiálja a streamingfunkciókat, és üzembe helyezi azt egy kiszolgálóvégponton: