Aller au contenu principal

Configurer un Stream

info

Aperçu

Cette fonctionnalité est en Aperçu public. Les administrateurs du Workspace peuvent contrôler l'accès à cette fonctionnalité à partir de la page Previews . Consultez Gérer les aperçus Databricks.

Un Stream représente une source de données streaming externe, telle qu’Apache Kafka. Les Streams stockent les détails de connexion, l'authentification, les schémas et la configuration d'ingestion. Une fois qu'un Stream est créé, vous pouvez le référencer à l'aide des définitions de Feature View pour créer des fonctionnalités de streaming en temps réel.

Les Stream ont des noms en trois parties (catalog.schema.stream_name). L'accès à un Stream est régi par sa table d'ingestion associée. Consultez Ingestion et remplissage pour plus de détails.

Exigences

  • Pour exécuter des commandes de notebook : serverless ou un cluster de compute classique exécutant Databricks Runtime 17.0 ML ou une version ultérieure.
  • Le package Python feature-engineering-client version 0.16.0 ou supérieure doit être installé.

Créer un Stream

Utilisez create_stream() pour créer un nouveau Stream. Un Stream nécessite quatre composants de configuration :

  • Configuration de la source : Spécifie la plateforme de streaming (par exemple, Kafka) et les détails spécifiques à la source (tels que l'abonnement aux rubriques pour Kafka).
  • Configuration de la connexion : spécifie comment se connecter et s’authentifier à la plateforme de streaming, y compris les serveurs d’amorçage et les informations d’identification.
  • Configuration du schéma : Définit la structure des clés et des valeurs de message.
  • Configuration de l'ingestion : Spécifie où et comment les données Stream sont ingérées. Consultez Ingestion et remplissage pour plus de détails.
Python
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"
),
),
)

Connexion à des sources de Stream

Avant de définir les fonctionnalités de streaming, connectez et testez une connexion de LakeFlow Pipelines de streaming à votre courtier Kafka. Consultez Streaming sur compute serverless et Se connecter à Apache Kafka.

Pour le streaming géré par AWS (Amazon MSK), consultez Connectivité privée Serverless à Amazon MSK. Pour plus de détails sur les options d'authentification Kafka, consultez Authentification.

Authentification

Connexion Unity Catalog (recommandée)

Utilisez une connexion Unity Catalog pour vous authentifier auprès de votre cluster Kafka. Il s'agit de l'approche recommandée pour l'authentification managée. Pour créer une connexion, voir Créer une connexion.

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

mTLS direct

Pour l'authentification mTLS directe, fournissez les fichiers de keystore et de truststore stockés sur un volume Unity Catalog, avec des mots de passe référencés via les Secret Scopes Databricks. Pour plus d'informations sur l'authentification SSL avec Kafka, voir Utiliser SSL pour connecter Databricks à Kafka.

Python
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

L'authentification SASL (à la fois SASL/SCRAM et SASL/PLAIN) n'est pas prise en charge pendant l'aperçu.

Modes d'abonnement

Le mode d'abonnement spécifie comment le Stream sélectionne les rubriques Kafka à consommer. Trois modes sont pris en charge :

Mode

Description

Exemple

subscribe

Liste de noms de rubriques séparés par des virgules

KafkaSubscriptionMode(subscribe="topic1,topic2")

subscribe_pattern

Noms de sujets correspondant aux modèles d'expressions régulières Java.

KafkaSubscriptionMode(subscribe_pattern="events-.*")

assign

JSON spécifiant les attributions de partitions de rubrique

KafkaSubscriptionMode(assign='{"my-topic": [0, 1, 2]}')

Mode

Description

Exemple

subscribe

Liste de noms de rubriques séparés par des virgules

KafkaSubscriptionMode(subscribe="topic1,topic2")

subscribe_pattern

Noms de sujets correspondant aux modèles d'expressions régulières Java.

KafkaSubscriptionMode(subscribe_pattern="events-.*")

assign

JSON spécifiant les attributions de partitions de rubrique

KafkaSubscriptionMode(assign='{"my-topic": [0, 1, 2]}')

Configuration du schéma

Définissez la structure des clés et des valeurs de message en utilisant le format JSON Schema. Pour les sources Kafka, payload_schema correspond à la valeur du message Kafka (la value dans le modèle clé-valeur de Kafka) et key_schema correspond à la clé du message Kafka. Au moins un de payload_schema ou key_schema doit être fourni.

Python
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"}'
),
)

Si aucun schéma n'est fourni pour une clé ou une charge utile, elle est traitée comme une simple chaîne de caractères.

Ingestion et remplissage

Le parameter ingestion_config configure la manière dont les données de Stream sont capturées et stockées pour l'entraînement et le service.

L'accès à un Stream est régi par la table d'ingestion :

  • SELECT sur la table d'ingestion accorde un accès en lecture au Stream.
  • MANAGE sur la table d'ingestion accorde l'accès en suppression.

Pour plus d'informations sur les privilèges de table, consultez Table et Référence des privilèges Unity Catalog.

Pipeline d'ingestion

Lorsqu'un Stream est créé, Databricks start un pipeline d'ingestion géré qui lit en continu les messages du sujet Kafka et les écrit dans une table Delta (la table d'ingestion). Le pipeline start à partir du dernier offset Kafka et s'exécute en continu, capturant uniquement les nouveaux messages qui arrivent après la création du Stream. Cette table d'ingestion est utilisée pour l'entraînement avec des fonctionnalités de streaming. Lorsqu'un Stream est supprimé, son pipeline d'ingestion et sa table d'ingestion sont également supprimés.

Destination d'ingestion

Le ingestion_destination spécifie le nom de table Delta en trois parties où les données de stream sont écrites.

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

Schéma de table d’ingestion

La table d'ingestion contient les données des messages ainsi que les colonnes de métadonnées :

Colonne

Type

Description

key

Varie (à partir de key_schema)

La clé de message Kafka, structurée selon le schéma que vous avez fourni.

value

Varie (à partir de payload_schema)

La valeur du message Kafka (payload), structurée selon le schéma que vous avez fourni.

stream_record_timestamp

TIMESTAMP

L'Timestamp du record. Pour les données de remplissage avant, il s'agit de l'ingestion du broker Kafka Timestamp. Pour les données de remplissage rétroactif, celles-ci sont fournies par le client.

kafka_topic

STRING

Le sujet Kafka d'où l'enregistrement a été consommé.

kafka_partition

INT

La partition Kafka à partir de laquelle l'enregistrement a été consommé.

kafka_offset

LONG

Le décalage Kafka de l'enregistrement au sein de sa partition.

record_source

STRING

Soit "stream" (remplissage vers l'avant à partir du Stream Kafka en direct) ou "backfill" (à partir de la source de remplissage rétroactif).

Colonne

Type

Description

key

Varie (à partir de key_schema)

La clé de message Kafka, structurée selon le schéma que vous avez fourni.

value

Varie (à partir de payload_schema)

La valeur du message Kafka (payload), structurée selon le schéma que vous avez fourni.

stream_record_timestamp

TIMESTAMP

L'Timestamp du record. Pour les données de remplissage avant, il s'agit de l'ingestion du broker Kafka Timestamp. Pour les données de remplissage rétroactif, celles-ci sont fournies par le client.

kafka_topic

STRING

Le sujet Kafka d'où l'enregistrement a été consommé.

kafka_partition

INT

La partition Kafka à partir de laquelle l'enregistrement a été consommé.

kafka_offset

LONG

Le décalage Kafka de l'enregistrement au sein de sa partition.

record_source

STRING

Soit "stream" (remplissage vers l'avant à partir du Stream Kafka en direct) ou "backfill" (à partir de la source de remplissage rétroactif).

Source de remplissage

Étant donné que le pipeline de remplissage avant start du dernier décalage Kafka, il ne capture pas les messages qui existaient avant la création du Stream. Pour fournir une couverture des données historiques pour l'entraînement, configurez une source de remplissage historique facultative.

Lorsqu’une source de remplissage est configurée, Databricks exécute un job MERGE INTO unique qui copie les lignes de remplissage dans la table d’ingestion avec record_source="backfill". Le Merge ne s'exécute qu'après que le vérificateur de chevauchement confirme que la source de remplissage rétrospectif et le Stream de remplissage prospectif ont des Timestamp qui se chevauchent (voir Chevauchement entre les données de remplissage rétrospectif et de Stream en direct). Si la condition de chevauchement n'est pas remplie dans les 2 jours, le Merge s'exécute quand même pour éviter un blocage indéfini.

La table de remplissage rétroactif doit inclure une colonne stream_record_timestamp de type TIMESTAMP dans le fuseau horaire UTC. D'autres colonnes de métadonnées Kafka (kafka_topic, kafka_partition, kafka_offset) sont transmises si elles sont présentes sur la source de remplissage rétroactif, ou définies sur NULL dans le cas contraire.

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

Chevauchement entre le remplissage historique et les données de Live Stream en direct

Avant d'exécuter un MERGE entre le remplissage rétrospectif et la table d'ingestion, une vérification de chevauchement compare les Timestamp sur les deux tables :

  • Remplissage max : le stream_record_timestamp maximum dans la source de remplissage.
  • Ingestion min : Le minimum stream_record_timestamp de lignes (record_source="stream") dans la table d'ingestion.

Le MERGE se poursuit lorsque le dernier Timestamp du remplissage excède le Timestamp le plus ancien de la table d'ingestion d'au moins 1 heure. Ce chevauchement garantit l'absence de lacunes dans la table d'ingestion. Si la condition de chevauchement n'est pas remplie dans les 2 jours, le MERGE s'exécute quand même pour éviter un blocage indéfini.

Étant donné que le pipeline d'ingestion démarre à partir du dernier offset Kafka, il ne capture que les messages arrivant après la création du stream. Votre source de remplissage rétrospectif doit contenir des données qui s'étendent dans la plage horaire d'ingestion — et non seulement jusqu'à l'heure de création du Stream.

Par exemple, si vous créez un Stream à 3 h 00, le pipeline de remplissage avant commence à lire les messages à partir de 3 h 00 et au-delà. Votre source de remplissage rétroactif doit inclure des Timestamp jusqu'à au moins 16 h 00 (1 heure après le start du remplissage avant) pour satisfaire la vérification de chevauchement. Cela signifie que vous devez mettre à jour votre table de remplissage rétroactif après 16 h 00 pour vous assurer que la table d'ingestion ne présente aucune lacune.

Déduplication

Utilisez deduplication_columns pour spécifier les chemins de colonne afin d'identifier les lignes en double lors de l'ingestion entre les données de Stream de rétro-remplissage et de remplissage anticipé. Utilisez la notation par points pour les champs imbriqués (par exemple, "value.user_id").

Choisissez les colonnes de déduplication en fonction de vos données :

  • Si chaque enregistrement dans votre Stream contient un identifiant unique (par exemple, value.transaction_id), utilisez cette colonne pour la déduplication.
  • Si votre source de remplissage inclut les colonnes kafka_partition et kafka_offset, utilisez-les pour identifier chaque enregistrement de manière unique.
  • Si aucune colonne de déduplication n'est spécifiée, la clé de déduplication default est la combinaison complète de key, value et stream_record_timestamp. Ceci n'est pas recommandé car cette correspondance de critères stricte peut facilement entraîner des doublons.
Python
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
deduplication_columns=["value.transaction_id"],
)

Gérer les Stream

Obtenir un stream

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

Lister les flux

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

Définissez include_schemas=True pour inclure les détails complets du schéma. Les schémas peuvent être volumineux, ce qui peut entraîner une Opération de longue durée. Pour récupérer les schémas individuellement, utilisez plutôt get_stream.

Supprimer un Stream

La suppression d'un stream supprime également son pipeline d'ingestion et sa table d'ingestion.

attention

Les modèles ou fonctionnalités qui référencent le Stream supprimé n'auront plus accès aux données du Stream sous-jacent. Créez une copie de la table d'ingestion avant la suppression si vous avez besoin de ces données mais que vous n'avez plus besoin du Stream.

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

Exemple de Notebook

Pour un exemple de bout en bout qui crée un Stream, définit des fonctionnalités en streaming et déploie vers un Endpoint de service, consultez le Notebook suivant :

Notebook de démarrage rapide des vues de fonctionnalités en streaming