Konfigurera en dataström

Important

Den här funktionen finns som allmänt tillgänglig förhandsversion. Arbetsyteadministratörer kan styra åtkomsten till den här funktionen från sidan Förhandsversioner . Se Hantera förhandsversioner av Azure Databricks.

En ström representerar en extern strömmande datakälla, till exempel Apache Kafka. Strömmar lagrar anslutningsinformation, autentisering, scheman och inmatningskonfiguration. När en ström har skapats kan du referera till den med hjälp av funktionsvydefinitioner för att skapa realtidsströmningsfunktioner.

Strömmar har tredelade namn (catalog.schema.stream_name). Åtkomst till en Stream styrs av dess associerade inmatningstabell. Mer information finns i Inmatning och återfyllnad .

Requirements

  • För att köra notebook-kommandon: serverlöst eller ett klassiskt beräkningskluster som kör Databricks Runtime 17.0 ML eller senare.
  • Python-paketversion feature-engineering-client 0.16.0 eller senare måste installeras.

Skapa en strömning

Använd create_stream() för att skapa en ny Stream. En Stream kräver fyra konfigurationskomponenter:

  • Källkonfiguration: Anger strömningsplattformen (till exempel Kafka) och källspecifik information (till exempel ämnesprenumeration för Kafka).
  • Anslutningskonfiguration: Anger hur du ansluter och autentiserar till strömningsplattformen, inklusive bootstrap-servrar och autentiseringsuppgifter.
  • Schemakonfiguration: Definierar strukturen för meddelandenycklar och -värden.
  • Inmatningskonfiguration: Anger var och hur dataström matas in. Mer information finns i Inmatning och återfyllnad .
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"
        ),
    ),
)

Anslutning till strömkällor

Innan du definierar strömningsfunktioner ansluter och testar du en strömmande Lakeflow-pipelineanslutning till din Kafka-mäklare. Se Direktuppspelning på serverlös beräkning och Anslut till Apache Kafka.

Information om AWS-hanterad direktuppspelning (Amazon MSK) finns i Serverlös privat anslutning till Amazon MSK. Mer information om Alternativ för Kafka-autentisering finns i Autentisering.

Authentication

Använd en Unity Catalog-anslutning för att autentisera till ditt Kafka-kluster. Det här är den rekommenderade metoden för hanterad autentisering. Information om hur du skapar en anslutning finns i Skapa en anslutning.

connection_config = StreamConnectionConfig(
    uc_connection_name="my-kafka-connection"
)

Direkt mTLS

För direkt mTLS-autentisering anger du keystore- och truststore-filer som lagras på en Unity Catalog-volym, med lösenord som anges via Databricks secret scopes. Mer information om SSL-autentisering med Kafka finns i Använda SSL för att ansluta Azure Databricks till Kafka.

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

SASL-autentisering (både SASL/SCRAM och SASL/PLAIN) stöds inte under förhandsversionen.

Prenumerationslägen

Prenumerationsläget anger hur Stream väljer Kafka-ämnen som den ska konsumera från. Tre lägen stöds:

Läge Description Exempel
subscribe Kommaavgränsad lista med ämnesnamn KafkaSubscriptionMode(subscribe="topic1,topic2")
subscribe_pattern Java regex-mönster för matchning av ämnesnamn KafkaSubscriptionMode(subscribe_pattern="events-.*")
assign JSON som anger tilldelningar för ämnespartition KafkaSubscriptionMode(assign='{"my-topic": [0, 1, 2]}')

Schemakonfiguration

Definiera strukturen för meddelandenycklar och värden så att inläsning och funktionsdefinitioner kan läsa enskilda fält. För Kafka-källor payload_schema motsvarar det Kafka-meddelandevärdet ( value i Kafkas nyckelvärdesmodell) och key_schema motsvarar Kafka-meddelandenyckeln. Minst en av payload_schema eller key_schema måste tillhandahållas.

Varje SchemaConfig format accepterar ett av tre format, som matchar hur källan serialiserar sina meddelanden: json_schema, avro_schema, eller proto_schema. Om inget schema anges för en nyckel eller nyttolast behandlas det som en enkel sträng.

JSON-schema

Ange en JSON-schemasträng till 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-schema

Ange en Avro-schemasträng till avro_schema. Avro-logiska typer stöds, inklusive timestamp-millis, date, och 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-schema

Lämna en ProtoSchemaSpec till med proto_schemaprotokollbuffertens.proto källtext och namnet på nyttolastmeddelandet. Importera ProtoSchemaSpec från databricks.feature_engineering.entities.

message_name måste vara det fullt kvalificerade meddelandenamnet, inklusive det package deklarerade i texten .proto (till exempel com.example.Event, inte Event). Både proto2- och proto3-syntax stöds.

google.protobuf.Timestamp och skalär-wrappertyperna (StringValue, Int32Value, och så vidare) stöds, och deras importer löses automatiskt. Andra välkända typer, såsom Duration, Struct, och Any, avvisas; koda istället dessa värden som en stödd skalär eller meddelande. Och fixed64 skalärtyperna fixed32 och map med icke-strängnycklar stöds inte heller.

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

Inmatning och återfyllnad

Parametern ingestion_config konfigurerar hur dataströmdata samlas in och lagras för träning och servering.

Åtkomst till en stream styrs av inmatningstabellen:

  • SELECT i inmatningstabellen ger läsåtkomst till Stream.
  • MANAGE i inmatningstabellen ger borttagningsåtkomst.

Mer information om tabellbehörigheter finns under Tabell och referens för behörigheter i Unity Catalog.

Inmatningspipeline

När en dataström skapas startar Databricks en hanterad inmatningspipeline som kontinuerligt läser meddelanden från Kafka-ämnet och skriver dem i en Delta-tabell (inmatningstabellen). Pipelinen startar från den senaste Kafka-offseten och körs kontinuerligt och fångar endast upp nya meddelanden som kommer in efter att strömmen har skapats. Den här inmatningstabellen används för träning med strömningsfunktioner. När en dataström tas bort tas även inmatningspipelinen och inmatningstabellen bort.

Inmatningsmål

ingestion_destination Anger det tredelade Delta-tabellnamnet där dataströmmen skrivs.

ingestion_config = IngestionConfig(
    ingestion_destination=IngestionDestination(
        delta_table_name="my_catalog.my_schema.events_ingestion"
    ),
)

Tabellschema för inmatning

Inmatningstabellen innehåller meddelandedata tillsammans med metadatakolumner:

Column Type Description
key Varierar (från key_schema) Kafka-meddelandenyckeln, strukturerad enligt det schema som du angav.
value Varierar (från payload_schema) Kafka-meddelandevärdet (nyttolasten), strukturerat enligt det schema som du angav.
stream_record_timestamp TIMESTAMP Postens tidsstämpel. För framåtfyllnadsdata är detta Kafka-brokerns inmatningstidsstämpel. För bakfyllnadsdata tillhandahålls dessa av kunden.
kafka_topic STRING Kafka-ämnet som posten förbrukades från.
kafka_partition INT Den Kafka-partition som posten förbrukades från.
kafka_offset LONG Postens Kafka-offset inom sin partition.
record_source STRING Antingen "stream" (framåtfyllning från den aktiva Kafka-strömmen) eller "backfill" (från återfyllnadskällan).

Återfyllnadskälla

Eftersom pipelinen för vidarebefordran startar från den senaste Kafka-förskjutningen samlar den inte in meddelanden som fanns innan dataströmmen skapades. Konfigurera en valfri återfyllnadskälla för att tillhandahålla historisk datatäckning för träning.

När en återfyllnadskälla har konfigurerats kör Databricks ett engångsjobb MERGE INTO som kopierar återfyllnadsrader till inmatningstabellen med record_source="backfill". MERGE körs först efter att överlappningskontrollen bekräftar att återfyllnadskällan och dataströmmen för vidarebefordran har överlappande tidsstämplar (se Överlappning mellan återfyllnads- och liveströmdata). Om överlappningsvillkoret inte uppfylls inom 2 dagar körs MERGE ändå för att undvika blockering på obestämd tid.

Tabellen för återfyllnad måste innehålla en stream_record_timestamp kolumn av typen TIMESTAMP i UTC-tidszonen. Andra Kafka-metadatakolumner (kafka_topic, kafka_partition, kafka_offset) vidarebefordras om de finns i återfyllnadskällan, eller annars sätts de till NULL.

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

Överlappning mellan återfyllnads- och liveströmdata

Innan du kör en MERGE mellan återfyllningen och inmatningstabellen jämför en överlappningskontroll tidsstämplarna i de två tabellerna:

  • Max för återfyllnad: Maximalt stream_record_timestamp i återfyllnadskällan.
  • Inmatningsmin: Det minsta stream_record_timestamp antalet rader (record_source="stream") i inmatningstabellen.

MERGE fortsätter när återfyllningens senaste tidsstämpel överskrider inmatningstabellens tidigaste tidsstämpel med minst 1 timme. Den här överlappningen säkerställer att det inte finns några luckor i inmatningstabellen. Om överlappningsvillkoret inte uppfylls inom 2 dagar körs MERGE ändå för att undvika blockering på obestämd tid.

Eftersom inmatningspipelinen startar från den senaste Kafka-offseten fångar den endast upp meddelanden som anländer efter att dataströmmen har skapats. Din återfyllnadskälla måste innehålla data som sträcker sig in i intagningstidsintervallet – inte bara fram till tidpunkten då strömmen skapades.

Om du till exempel skapar en ström kl. 15:00 börjar pipelinen för vidarebefordran att läsa meddelanden från 15:00 och framåt. Din återfyllnadskälla måste innehålla data med tidsstämplar till minst 16:00 (1 timme efter start av framåtfyllning) för att uppfylla överlappningskontrollen. Det innebär att du bör uppdatera din återfyllnadstabell efter 16:00 för att säkerställa att inmatningstabellen inte har några luckor.

Deduplication

Använd deduplication_columns för att ange kolumnsökvägar för att identifiera duplicerade rader vid inmatning av backfill-data och strömmande forward-fill-data. Använd punkt notation för kapslade fält (till exempel "value.user_id").

Välj dedupliceringskolumner baserat på dina data:

  • Om varje post i dataströmmen innehåller en unik identifierare (till exempel value.transaction_id), använder du den kolumnen för deduplicering.
  • Om din återfyllnadskälla innehåller kafka_partition och kafka_offset kolumner använder du dem för att unikt identifiera varje post.
  • Om inga dedupliceringskolumner anges är standarddedupliceringsnyckeln den fullständiga kombinationen av key, valueoch stream_record_timestamp. Detta rekommenderas inte eftersom den här strikta villkorsmatchningen enkelt kan leda till dubbletter.
ingestion_config = IngestionConfig(
    ingestion_destination=IngestionDestination(
        delta_table_name="my_catalog.my_schema.events_ingestion"
    ),
    deduplication_columns=["value.transaction_id"],
)

Hantera strömmar

Hämta en dataström

stream = client.get_stream(name="my_catalog.my_schema.my_stream")

Visa strömmar

streams = client.list_streams(
    catalog_name="my_catalog",
    schema_name="my_schema",
    max_results=50,
    include_schemas=False,
)

Ange include_schemas=True för att inkludera fullständig schemainformation. Scheman kan vara stora och detta kan resultera i en tidskrävande åtgärd. Om du vill hämta scheman individuellt i stället använder du get_stream.

Ta bort en dataström

Att ta bort en dataström tar även bort både dess inmatningspipeline och inmatningstabell.

Varning

Modeller eller funktioner som refererar till den borttagna dataströmmen har inte längre åtkomst till underliggande dataström. Skapa en kopia av inmatningstabellen före borttagning om du behöver dessa data men inte längre behöver dataströmmen.

client.delete_stream(name="my_catalog.my_schema.my_stream")

Exempelanteckningsbok

Ett heltäckande exempel som skapar en Stream, definierar streamingfunktioner och driftsätter till en slutpunkt för servering finns i följande notebook:

Snabbstartsguide för strömmande funktionsvyer

Hämta anteckningsbok