Hinweis
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, sich anzumelden oder das Verzeichnis zu wechseln.
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, das Verzeichnis zu wechseln.
Von Bedeutung
Dieses Feature befindet sich in der Public Preview. Arbeitsbereichsadministratoren können den Zugriff auf dieses Feature über die Vorschauseite steuern. Siehe Manage Azure Databricks Previews.
Ein Stream stellt eine externe Streamingdatenquelle dar, z. B. Apache Kafka. Datenströme speichern Verbindungsdetails, Authentifizierung, Schemas und Aufnahmekonfiguration. Nachdem ein Datenstrom erstellt wurde, können Sie mithilfe von Feature View-Definitionen darauf Bezug nehmen, um Streaming-Features in Echtzeit zu erstellen.
Datenströme haben dreiteilige Namen (catalog.schema.stream_name). Der Zugriff auf einen Stream wird durch die zugeordnete Erfassungstabelle gesteuert. Siehe -Aufnahme und Abgleich für Details.
Requirements
- Zum Ausführen von Notizbuchbefehlen: Serverless oder ein klassischer Computecluster mit Databricks Runtime 17.0 ML oder höher.
- Das
feature-engineering-clientPython-Paket Version 0.17.0 oder höher muss installiert werden.
Herstellen einer Verbindung mit Streamquellen
Bevor Sie Streamingfunktionen definieren, stellen Sie eine Streaming-Lakeflow-Pipeline-Verbindung zu Ihrem Kafka-Broker her und testen Sie diese. Feature Store basiert auf serverlosem SDP, was bedeutet, dass Sie einen Mechanismus benötigen, um Ihre klassischen Rechenressourcen (Broker oder Endpoint) mit den serverlosen Rechenressourcen von Databricks zu verbinden. Dies geschieht über Produkte wie privatelink oder indem Ihre klassische Rechenleistung vom öffentlichen Internet aus zugänglich ist.
Einen Stream erstellen
Verwenden Sie create_stream(), um einen neuen Stream zu erstellen. Ein Stream erfordert vier Konfigurationskomponenten:
- Quellkonfiguration: Spezifiziert die Streaming-Plattform und quellenspezifische Details, wie z. B. das Themenabonnement für eine Kafka-Quelle.
- Verbindungskonfiguration: Gibt an, wie Eine Verbindung mit der Streamingplattform hergestellt und authentifiziert wird, einschließlich Bootstrap-Server und Anmeldeinformationen.
- Schemakonfiguration: Definiert die Struktur von Nachrichtenschlüsseln und -werten.
- Aufnahmekonfiguration: Gibt an, wo und wie Streamdaten aufgenommen werden. Siehe -Aufnahme und Abgleich für Details.
Für die quellenspezifische source_config und Verbindungsstruktur sowie ein vollständiges create_stream() Beispiel siehe Apache Kafka. Das Schema und die Aufnahmeoptionen werden quellenübergreifend geteilt.
Apache Kafka
Um Daten aus Apache Kafka zu streamen, verwenden Sie KafkaStreamConfig als Quellkonfiguration und eine Unity-Catalog-Verbindung zur Authentifizierung. Siehe Streaming mit serverloser Rechenleistung und Mit Apache Kafka verbinden zur Kafka-Konnektivität.
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"
),
),
)
Kafka-Abonnementmodi
Der Abonnementmodus gibt an, von welchen Kafka-Themen der Stream Daten konsumiert. Es werden drei Modi unterstützt:
| Modus | Description | Example |
|---|---|---|
subscribe |
Durch Trennzeichen getrennte Liste der Themennamen | KafkaSubscriptionMode(subscribe="topic1,topic2") |
subscribe_pattern |
Themennamen, die mit einem Java-Regex-Muster übereinstimmen | KafkaSubscriptionMode(subscribe_pattern="events-.*") |
assign |
JSON-Code zur Angabe von Themenpartitionszuweisungen | KafkaSubscriptionMode(assign='{"my-topic": [0, 1, 2]}') |
Kafka-Authentifizierung
Unity-Katalogverbindung (empfohlen)
Verwenden Sie eine Unity-Katalogverbindung, um sich bei Ihrem Kafka-Cluster zu authentifizieren. Dies ist der empfohlene Ansatz für die verwaltete Authentifizierung. Informationen zum Erstellen einer Verbindung finden Sie unter Erstellen einer Verbindung. Der Ersteller des Streams muss auf der Verbindung über USE CONNECTION verfügen. Jeder Nutzer, der Features mit dem Stream als Quelle erstellt, muss außerdem USE CONNECTION für die Verbindung haben.
connection_config = StreamConnectionConfig(
uc_connection_name="my-kafka-connection"
)
Die Verbindung unterstützt sowohl IAM (Service Credential) als auch SASL-Authentifizierung.
IAM (Dienstanmeldeinformationen)
Authentifizieren Sie sich mit einer Unity Catalog-Dienstanmeldeinformation, beispielsweise um eine Verbindung mit Amazon MSK mithilfe von IAM herzustellen. Informationen zum Erstellen von Anmeldeinformationen finden Sie unter Anmeldeinformationen erstellen. Setzen Sie den Namen der Service-Zugangsdaten mit der Option credential :
CREATE CONNECTION IF NOT EXISTS `my-kafka-connection`
TYPE KAFKA
OPTIONS (
bootstrap_servers '<bootstrap_servers>',
credential '<service_credential>'
)
Zusätzlich zu USE CONNECTION für die Verbindung, benötigen Identitäten, die die Dienstanmeldeinformationen verwenden, auch ACCESS. Erteilen Sie ACCESS für die referenzierten Dienstanmeldeinformationen an den Creator des Streams und für jede Identität, die Funktionen mit dem Stream materialisiert. Siehe Erteilen von Berechtigungen zum Verwenden von Dienstanmeldeinformationen für den Zugriff auf einen externen Clouddienst.
SASL
Die SASL-Authentifizierung verwendet Benutzernamen und Passwörter. Legen Sie sasl_mechanism auf eine der folgenden Optionen fest:
PLAINSCRAM-SHA-256SCRAM-SHA-512
Geben Sie die Anmeldedaten mit den Optionen user und password an. Die Verbindung speichert diese Zugangsdaten sicher.
Das folgende Beispiel verwendet SASL/SCRAM. Für SASL/PLAIN setzen Sie sasl_mechanism auf PLAIN.
CREATE CONNECTION IF NOT EXISTS `my-kafka-connection`
TYPE KAFKA
OPTIONS (
bootstrap_servers '<bootstrap_servers>',
sasl_mechanism 'SCRAM-SHA-512',
user '<username>',
password '<password>'
)
direktes mTLS
Stellen Sie für die direkte mTLS-Authentifizierung Schlüsselspeicher- und Truststore-Dateien bereit, die auf einem Unity-Katalogvolume gespeichert sind, mit Kennwörtern, auf die über geheime Databricks-Bereiche verwiesen wird. Weitere Informationen zur SSL-Authentifizierung mit Kafka finden Sie unter Verwendung von SSL zum Herstellen einer Verbindung mit Azure Databricks mit 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"
),
),
)
Schema-Konfiguration
Definieren Sie die Struktur von Nachrichtenschlüsseln und -werten, sodass Ingestion und Feature-Definitionen einzelne Felder lesen können. Für Kafka-Quellen entspricht payload_schema dem Wert der Kafka-Nachricht (dem value im Schlüssel-Wert-Modell von Kafka), und key_schema entspricht dem Schlüssel der Kafka-Nachricht. Mindestens einer von payload_schema oder key_schema muss angegeben werden.
Jedes SchemaConfig akzeptiert eines von drei Formaten, die übereinstimmen, wie die Quelle ihre Nachrichten serialisiert: json_schema, avro_schema, oder proto_schema. Wenn kein Schema für einen Schlüssel oder eine Nutzlast bereitgestellt wird, wird es als einfache Zeichenfolge behandelt.
Die Codebeispiele in diesem Abschnitt verwenden Schemata, die inline mit deklariert sind DirectSchemas, wobei das Schema als Zeichenkette bereitgestellt wird. Um Schemata mit einer externen Schemaregistrierung zu verwalten, siehe Schema-Registry für Details.
JSON-Schema
Geben Sie eine JSON Schema-Zeichenfolge für json_schema an.
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
Gib eine Avro-Schema-Zeichenfolge für avro_schema. Avro-logische Typen werden unterstützt, darunter timestamp-millis, date, und 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
Geben Sie ein ProtoSchemaSpec an proto_schema mit dem Quelltext des Protokollpuffers.proto und dem Namen der Payload-Nachricht an. Importiere ProtoSchemaSpec aus databricks.feature_engineering.entities.
message_name muss der vollständig qualifizierte Name der Nachricht sein, einschließlich der im Text .proto deklarierten package (zum Beispiel com.example.Event, nicht Event). Sowohl proto2- als auch proto3-Syntax werden unterstützt.
google.protobuf.Timestamp und die skalaren Wrapper-Typen (StringValue, Int32Value, und so weiter) werden unterstützt, deren Importe automatisch aufgelöst werden. Andere bekannte Typen wie Duration, Struct, und Any, werden abgelehnt; kodieren Sie diese Werte stattdessen als unterstützte Skalare oder Nachricht. Auch die Skalartypen fixed32 und fixed64map mit Nicht-String-Tasten werden nicht unterstützt.
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",
)
),
)
Entschlüsselung von Daten mit Schemata
Databricks decodiert jede Nachricht mit den Spark-Funktionen from_json, from_avro und from_protobuf. Folgende Verhaltensweisen gelten, egal ob Sie das Schema inline deklarieren oder es aus einer Schema-Registry lösen:
- Fehlgeformte Platten. Die Dekodierung erfolgt im Modus
PERMISSIVE, sodass ein Datensatz, der nicht mit seinem Schema übereinstimmt, als Nullwert dekodiert wird, anstatt zum Fehlschlagen des Streams zu führen. - Avro-Gewerkschaften. Eine Vereinigung mehrerer Datensatztypen dekodiert zu einer Struktur mit einem Feld pro Datensatztyp, von denen jeder nach seinem Avro-Datensatz benannt ist.
- Protobuf-Typen. Unsignierte ganze Zahlen dekodieren zu einem breiteren vorzeichenmäßigen Typ (zum Beispiel
uint32aufBIGINTunduint64aufDECIMAL(20,0)), Enum-Felder dekodieren auf ihren Stringnamen, und skalare Wrapper-Typen (zum BeispielStringValueundInt32Value) dekodieren auf eine nullierbare Spalte des gewickelten Typs.
Schemaregistrierung
Schema-Register speichern und versiegeln Schemata, die von Streaming-Produzenten und -Konsumenten verwendet werden, und setzen Kompatibilitätsregeln durch, während sich diese Schemata weiterentwickeln. Wenn eine externe Schema-Registry konfiguriert ist, liest Feature Store das Schema aus der Registry und verwendet es, um die Streaming-Nachricht zu dekodieren. Du deklarierst das Schema nicht inline im Stream, wenn du eine Schema-Registry verwendest.
Die Unterstützung für das Schema-Register hat folgende Einschränkungen:
- Unterstützt nur für Kafka-Streams.
- Es wird nur das Confluent Schema Registry unterstützt.
- Nur die Avro - und Protobuf-Formate werden unterstützt. Um JSON-Nachrichten zu lesen, deklarieren Sie stattdessen das Schema inline. Siehe JSON-Schema.
- Jeder Stream ist mit genau einem Confluent-Subjekt für den Nachrichtenwert und einem für den Nachrichtenschlüssel (sofern vorhanden) verbunden. Stream-Themen mit mehreren Schema-Datensätzen sind keine unterstützte Konfiguration. Wenn Ihr Stream mit Themen verbunden ist, die mehrere Schemata enthalten, werden Datensätze, die nicht mit dem Schema des angegebenen Subjekts übereinstimmen, als null dekodiert.
Verbinden Sie sich mit einem Schema-Register
Stellen Sie die Registry-Verbindungsdetails als Optionen auf der Kafka Unity Catalog-Verbindung bereit und speichern Sie das Geheimnis der Registrierungs-API in einem Databricks-Geheimnisumfang. Die Run-as-Identität des Streams muss über die Berechtigung READ für den Secret-Bereich verfügen, da die Ingestion-Pipeline das Geheimnis zur Laufzeit liest. Wie man eine Verbindung erstellt und konfiguriert, siehe Create a Connection.
Füge die schema_registry_url, schema_registry_api_key und schema_registry_api_secret-Optionen zur für die Authentifizierung verwendeten Verbindung hinzu. Das folgende Beispiel erzeugt eine Kafka-Verbindung, die sich mit einem Unity Catalog Service-Credential beim Broker authentifiziert und mit einem API-Schlüssel beim Register tätig ist:
CREATE CONNECTION IF NOT EXISTS `my-kafka-connection`
TYPE KAFKA
OPTIONS (
bootstrap_servers '<bootstrap_servers>',
credential '<service_credential>',
schema_registry_url 'https://<registry-host>',
schema_registry_api_key '<registry_api_key>',
schema_registry_api_secret secret('<scope>', '<key>')
)
Setze sowohl die schema_registry_api_secret-Option in der Kafka-Verbindung als auch die Secret-Scope-Referenz im Stream auf dasselbe Geheimnis.
Erstellen Sie einen Strom, der ein Schema-Register verwendet
Übergeben Sie SchemaRegistryConfig als schema_config. Beziehen Sie sich auf das Registry-API-Geheimnis mit api_secret_ref, und identifizieren Sie das Thema und das Format mit payload_schema_locator für den Nachrichtenwert oder key_schema_locator für den Nachrichtenschlüssel. Mindestens ein Locator muss bereitgestellt werden.
Beachten Sie hier die Unterschiede im Vergleich zu den direkten Schema-Beispielen im Abschnitt Schema-Konfiguration . Wenn Sie eine Schemaregistrierung verwenden, geben Sie das Schema im Stream nicht inline für schema_config an. Stattdessen gibst du ein SchemaRegistryConfig an, das das Schema in der Registry identifiziert.
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
KafkaStreamConfig,
KafkaSubscriptionMode,
StreamConnectionConfig,
SchemaRegistryConfig,
SchemaLocator,
SchemaLocatorConfluentSchema,
SchemaLocatorFormat,
SecretScopeReference,
IngestionConfig,
IngestionDestination,
)
client = FeatureEngineeringClient()
stream = client.create_stream(
name="my_catalog.my_schema.my_stream",
source_config=KafkaStreamConfig(
subscription_mode=KafkaSubscriptionMode(subscribe="transactions"),
),
connection_config=StreamConnectionConfig(
uc_connection_name="my-kafka-connection"
),
schema_config=SchemaRegistryConfig(
api_secret_ref=SecretScopeReference(
scope="my_scope", key="sr_api_secret"
),
payload_schema_locator=SchemaLocator(
confluent_schema=SchemaLocatorConfluentSchema(
subject="transactions-value"
),
format=SchemaLocatorFormat.FORMAT_AVRO,
),
),
ingestion_config=IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.transactions_ingestion"
),
),
)
Ein Confluent-Subjekt ist der benannte Anwendungsbereich, unter dem die Versionshistorie eines Schemas registriert und die Kompatibilität erzwungen wird. Setzen Sie subject auf den Namen des relevanten Scopes, der üblicherweise aus der Subjektnamensstrategie bestimmt wird:
-
TopicNameStrategy (default, leitet das Subjekt aus dem Topic-Namen ab):
<topic>-valuefür den Wert und<topic>-keyfür den Schlüssel. Zum Beispiel verwendet das Wertschema für das Thematransactionsdas Subjekttransactions-value. -
RecordNameStrategy (leitet das Subjekt vom Datensatznamen des Schemas ab, unabhängig vom Topic): der vollqualifizierte Datensatzname, beispielsweise
com.example.Payment. Dies ist der Namespace und Name des Datensatzes für Avro bzw. das Paket und der Name der Nachricht für Protobuf. -
TopicRecordNameStrategy (kombiniert Thema- und Datensatznamen):
<topic>-<fully-qualified-record-name>, zum Beispieltransactions-com.example.Payment.
format ist erforderlich. Stellen Sie es auf SchemaLocatorFormat.FORMAT_AVRO oder SchemaLocatorFormat.FORMAT_PROTOBUF ein, damit es der Art entspricht, wie das Thema serialisiert wird.
Schemaentwicklung
Die Ingestion-Pipeline löst das aktuelle Schema des Subjekts beim Start auf. Wenn Sie eine neue abwärtskompatible Schema-Version auf dem Subjekt im Schema-Register registrieren, verwendet die laufende Pipeline weiterhin die Version, mit der sie gestartet war.
Da Databricks die Ingestion-Pipeline als serverlose Lakeflow-Pipeline verwaltet, startet die Pipeline periodisch neu. Beim nächsten Neustart wird die neue Schema-Version übernommen. Es kann bis zu einer Woche dauern, bis neue oder geänderte Felder in der Eingabetabelle erscheinen.
Wie die Pipeline mit Datensätzen umgeht, die nicht mit dem aktuell verwendeten Schema übereinstimmen, siehe Decoding Data Using Schemas.
Aufnahme und Nachfüllung
Der ingestion_config-Parameter konfiguriert, wie Streamdaten für Training und Bereitstellung erfasst und gespeichert werden.
Der Zugriff auf einen Stream wird durch die Ingestionstabelle gesteuert:
-
SELECTin der Erfassungstabelle gewährt Lesezugriff auf den Stream. -
MANAGEin der Erfassungstabelle gewährt Löschzugriff.
Weitere Informationen zu Tabellenberechtigungen finden Sie unter Tabelle und Referenz zu den Unity Catalog-Berechtigungen.
Erfassungspipeline
Wenn ein Stream erstellt wird, startet Databricks eine verwaltete Ingestionspipeline, die kontinuierlich Nachrichten aus dem Quell-Stream liest und sie in eine Delta-Tabelle (die Ingestionstabelle) schreibt. Die Pipeline startet ab der neuesten Position in der Quelle und läuft kontinuierlich weiter, wobei nur neue Nachrichten erfasst werden, die nach der Erstellung des Streams eintreffen. Diese Ingestionstabelle wird für das Trainieren mit Streaming-Features verwendet. Wenn ein Datenstrom gelöscht wird, werden auch die Erfassungspipeline und die Erfassungstabelle gelöscht.
Aufnahmeziel
Mit ingestion_destination wird der aus drei Teilen bestehende Name der Delta-Tabelle angegeben, in die Streamdaten geschrieben werden.
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
)
Schema der Erfassungstabelle
Die Aufnahmetabelle enthält die Nachrichtendaten zusammen mit Metadaten-Spalten. Die gemeinsamen Spalten sind bei jeder Quelle vorhanden; die kafka_* Spalten sind nur für einen Kafka-Stream vorhanden.
| Spalte | Typ | Source | Description |
|---|---|---|---|
key |
Variiert je nach key_schema |
Common | Der Nachrichtenschlüssel, strukturiert nach dem von dir bereitgestellten Schema. |
value |
Variiert je nach payload_schema |
Common | Der Nachrichtenwert (Payload), strukturiert nach dem von dir bereitgestellten Schema. |
stream_record_timestamp |
TIMESTAMP |
Common | Der Zeitstempel des Datensatzes. Für Forward-Fill-Daten ist dies der Quell-Ingest-Zeitstempel. Für Backfill-Daten wird dies vom Kunden bereitgestellt. |
record_source |
STRING |
Common | Entweder "stream" (Vorwärts-Ausfüllung aus dem Live-Stream) oder "backfill" (aus der Abgleichsquelle). |
kafka_topic |
STRING |
Kafka | Das Kafka-Thema, aus dem der Datensatz konsumiert wurde. |
kafka_partition |
INT |
Kafka | Die Kafka-Partition, aus der der Datensatz konsumiert wurde. |
kafka_offset |
LONG |
Kafka | Der Kafka-Offset des Datensatzes innerhalb seiner Partition. |
Quelle für die Nachbefüllung
Da die Forward-Fill-Pipeline bei der neuesten Position in der Quelle startet, erfasst sie keine Nachrichten, die vor der Erstellung des Streams vorhanden waren. Um historische Datenabdeckung für Schulungen bereitzustellen, konfigurieren Sie eine optionale Backfill-Quelle.
Wenn eine Backfill-Quelle konfiguriert ist, führt Databricks einen einmaligen MERGE INTO Job aus, der Backfill-Zeilen in die Ingestion-Tabelle mit record_source="backfill" kopiert. Der MERGE wird erst ausgeführt, nachdem die Überlappungsprüfung bestätigt hat, dass die Backfill-Quelle und der Forward-Fill-Stream überlappende Zeitstempel haben (siehe Überlappung zwischen Backfill- und Live-Stream-Daten). Wenn die Überlappungsbedingung nicht innerhalb von 2 Tagen erfüllt wird, wird der MERGE trotzdem durchgeführt, um eine Blockierung auf unbestimmte Zeit zu verhindern.
Die Backfill-Tabelle muss eine stream_record_timestamp-Spalte vom Typ TIMESTAMP in der UTC-Zeitzone enthalten. Andere Metadatenspalten werden übernommen, wenn sie in der Backfill-Quelle vorhanden sind, andernfalls werden sie auf NULL gesetzt. Für Kafka sind dies kafka_topic, kafka_partition und kafka_offset.
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"
),
)
Überlappung zwischen Backfill- und Live-Stream-Daten
Vor dem Ausführen eines MERGE zwischen der Abgleichs- und der Erfassungstabelle vergleicht eine Überlappungsprüfung die Zeitstempel in den beiden Tabellen:
-
Backfill-Maximum: Der maximale Wert
stream_record_timestampin der Backfill-Quelle. -
Erfassung min.: Die Mindestanzahl
stream_record_timestampvon Zeilen (record_source="stream") in der Erfassungstabelle.
Der MERGE wird fortgesetzt, wenn der letzte Zeitstempel des Abgleichs mindestens eine Sunde später ist als der früheste Zeitstempel der Erfassungstabelle. Diese Überlappung stellt sicher, dass in der Aufnahmetabelle keine Lücken vorhanden sind. Wenn die Überlappungsbedingung nicht innerhalb von 2 Tagen erfüllt wird, wird der MERGE trotzdem durchgeführt, um eine Blockierung auf unbestimmte Zeit zu verhindern.
Da die Aufnahmepipeline an der neuesten Position in der Quelle beginnt, erfasst sie nur Nachrichten, die nach der Erstellung des Streams eintreffen. Ihre Backfill-Quelle muss Daten enthalten, die bis in den Ingestionszeitbereich hineinreichen – nicht nur bis zum Zeitpunkt der Stream-Erstellung.
Wenn Sie z. B. um 15:00 Uhr einen Datenstrom erstellen, beginnt die Pipeline für das Vorwärts-Ausfüllen, Nachrichten ab 15:00 Uhr zu lesen. Ihre Backfill-Quelle muss Daten mit Zeitstempeln bis mindestens 16:00 Uhr (1 Stunde nach dem Start der Vorwärtsfüllung) enthalten, um die Überlappungsprüfung zu bestehen. Dies bedeutet, dass Sie die Backfill-Tabelle nach 14:00 Uhr aktualisieren sollten, um sicherzustellen, dass die Aufnahmetabelle keine Lücken aufweist.
Deduplication
Verwenden Sie deduplication_columns, um Spaltenpfade anzugeben, mit denen doppelte Zeilen bei der Datenerfassung zwischen Abgleichs- und Vorwärts-Ausfüll-Streamdaten identifiziert werden. Verwenden Sie die Punktnotation für geschachtelte Felder (z. B "value.user_id". ).
Auswählen von Deduplizierungsspalten basierend auf Ihren Daten:
- Wenn jeder Datensatz in Ihrem Stream eine eindeutige Kennung enthält (z. B.
value.transaction_id), verwenden Sie diese Spalte zur Deduplizierung. - Wenn Ihre Backfill-Quelle die Spalten
kafka_partitionundkafka_offsetenthält, verwenden Sie diese, um jeden Datensatz eindeutig zu identifizieren. - Wenn keine Deduplizierungsspalten angegeben werden, ist der Standard-Deduplizierungsschlüssel die vollständige Kombination aus
key, undvaluestream_record_timestamp. Dies wird nicht empfohlen, da dieser strenge Kriterienabgleich leicht zu Duplikaten führen kann.
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
deduplication_columns=["value.transaction_id"],
)
Streams verwalten
Einen Stream abrufen
stream = client.get_stream(name="my_catalog.my_schema.my_stream")
Datenströme auflisten
streams = client.list_streams(
catalog_name="my_catalog",
schema_name="my_schema",
max_results=50,
include_schemas=False,
)
include_schemas=True so festlegen, dass vollständige Schemadetails eingeschlossen werden. Schemata können groß sein, was zu einem lange laufenden Vorgang führen kann. Um stattdessen Schemas einzeln abzurufen, verwenden Sie get_stream.
Löschen Sie einen Stream
Beim Löschen eines Datenstroms werden auch die Erfassungspipeline und die Erfassungstabelle gelöscht.
Warning
Alle Modelle oder Features, die auf den gelöschten Datenstrom verweisen, haben keinen Zugriff mehr auf die zugrunde liegenden Datenstromdaten. Erstellen Sie eine Kopie der Aufnahmetabelle vor dem Löschen, wenn Sie diese Daten benötigen, aber den Datenstrom nicht mehr benötigen.
client.delete_stream(name="my_catalog.my_schema.my_stream")
Beispiel-Notebook
Ein End-to-End-Beispiel, das einen Stream erstellt, Streamingfeatures definiert und auf einem Bereitstellungsendpunkt bereitstellt, finden Sie im folgenden Notebook: