Configurer un Stream
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-clientversion 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.
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.
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.
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 |
|---|---|---|
| Liste de noms de rubriques séparés par des virgules |
|
| Noms de sujets correspondant aux modèles d'expressions régulières Java. |
|
| JSON spécifiant les attributions de partitions de rubrique |
|
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.
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 :
SELECTsur la table d'ingestion accorde un accès en lecture au Stream.MANAGEsur 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.
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 |
|---|---|---|
| Varie (à partir de | La clé de message Kafka, structurée selon le schéma que vous avez fourni. |
| Varie (à partir de | La valeur du message Kafka (payload), structurée selon le schéma que vous avez fourni. |
|
| 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. |
|
| Le sujet Kafka d'où l'enregistrement a été consommé. |
|
| La partition Kafka à partir de laquelle l'enregistrement a été consommé. |
|
| Le décalage Kafka de l'enregistrement au sein de sa partition. |
|
| Soit |
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.
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_timestampmaximum dans la source de remplissage. - Ingestion min : Le minimum
stream_record_timestampde 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_partitionetkafka_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,valueetstream_record_timestamp. Ceci n'est pas recommandé car cette correspondance de critères stricte peut facilement entraîner des doublons.
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
stream = client.get_stream(name="my_catalog.my_schema.my_stream")
Lister les flux
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.
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.
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 :