Stream beállítása

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-client Python 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

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:

  • SELECT a betöltési táblában olvasási hozzáférést biztosít a Streamhez.
  • MANAGE a 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_timestamp szá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 és kafka_offset oszlopokat, azokkal egyedileg azonosíthatja az egyes rekordokat.
  • Ha nincs megadva deduplikációs oszlop, az alapértelmezett deduplikációs kulcs a key, value és stream_record_timestamp teljes 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:

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

Jegyzetfüzet szerezz