Kommentar
Åtkomst till den här sidan kräver auktorisering. Du kan prova att logga in eller ändra kataloger.
Åtkomst till den här sidan kräver auktorisering. Du kan prova att ändra kataloger.
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-client0.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
Anslutning till Unity Catalog (rekommenderas)
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:
-
SELECTi inmatningstabellen ger läsåtkomst till Stream. -
MANAGEi 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_timestampi återfyllnadskällan. -
Inmatningsmin: Det minsta
stream_record_timestampantalet 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_partitionochkafka_offsetkolumner använder du dem för att unikt identifiera varje post. - Om inga dedupliceringskolumner anges är standarddedupliceringsnyckeln den fullständiga kombinationen av
key,valueochstream_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: