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-clienten version 0.17.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 gérée. Pour créer une connexion, consultez Créer une connexion. Le créateur du Stream doit disposer de USE CONNECTION sur la connexion. Tout utilisateur matérialisant des fonctionnalités avec le Stream comme source doit également disposer de USE CONNECTION sur la connexion.
connection_config = StreamConnectionConfig(
uc_connection_name="my-kafka-connection"
)
La connexion prend en charge l’authentification IAM (identifiant de service) et SASL.
IAM (identifiant de service)
Authentifiez-vous avec un identifiant de service Unity Catalog, par exemple pour vous connecter à Amazon MSK avec IAM. Pour créer un identifiant de service, voir Créer des identifiants de service. Définissez le nom de l’identifiant de service avec l’option credential :
CREATE CONNECTION IF NOT EXISTS `my-kafka-connection`
TYPE KAFKA
OPTIONS (
bootstrap_servers '<bootstrap_servers>',
credential '<service_credential>'
)
En plus de USE CONNECTION sur la connexion, les identités qui utilisent l’identifiant de service ont besoin de ACCESS sur celle-ci. Accordez ACCESS sur l'identifiant de service référencé au créateur du Stream et à toute identité qui matérialise des features avec le Stream. Voir Accorder des autorisations pour utiliser un identifiant de service afin d'accéder à un service cloud externe.
SASL
L'authentification SASL utilise un nom d'utilisateur et un mot de passe. Définissez sasl_mechanism sur l'une des valeurs suivantes :
PLAINSCRAM-SHA-256SCRAM-SHA-512
Fournissez les identifiants avec les options user et password. La connexion stocke ces identifiants en toute sécurité.
L’exemple suivant utilise SASL/SCRAM. Pour SASL/PLAIN, définissez sasl_mechanism sur 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>'
)
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"
),
),
)
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 afin que les définitions d'ingestion et de fonctionnalités puissent lire les champs individuels. Pour les sources Kafka, payload_schema correspond à la valeur du message Kafka (le value dans le modèle clé-valeur de Kafka) et key_schema correspond à la clé du message Kafka. Au moins l'un des éléments payload_schema ou key_schema doit être fourni.
Chaque SchemaConfig accepte l'un des trois formats, correspondant à la manière dont la source sérialise ses messages : json_schema, avro_schema ou proto_schema. Si aucun schéma n'est fourni pour une clé ou une charge utile, celle-ci est traitée comme une simple chaîne de caractères.
Les exemples de code de cette section utilisent des schémas déclarés en ligne avec DirectSchemas, où le schéma est fourni sous forme de chaîne. Pour gérer les schémas à l’aide d’un registre de schémas externe, consultez Schema registry pour plus de détails.
Schéma JSON
Fournissez une chaîne Schéma JSON à 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"}'
),
)
Schéma Avro
Fournissez une chaîne de schéma Avro à avro_schema. Les types logiques Avro sont pris en charge, notamment timestamp-millis, date et 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"}}'
' ]'
'}'
)
),
)
Schéma Protobuf
Fournissez un ProtoSchemaSpec à proto_schema avec le texte source Protocol Buffers .proto et le nom du message de charge utile. Importer ProtoSchemaSpec depuis databricks.feature_engineering.entities.
message_name doit être le nom de message entièrement qualifié, incluant le package déclaré dans le texte .proto (par exemple, com.example.Event, et non Event). Les syntaxes proto2 et proto3 sont toutes deux prises en charge.
google.protobuf.Timestamp et les types de wrapper scalaires (StringValue, Int32Value, etc.) sont pris en charge, et leurs importations sont résolues automatiquement. D’autres types bien connus, tels que Duration, Struct et Any, sont rejetés ; encodez plutôt ces valeurs sous forme de scalaire ou de message pris en charge. Les types scalaires fixed32 et fixed64 ainsi que map avec des clés non-chaîne ne sont pas non plus pris en charge.
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",
)
),
)
Décodage des données à l’aide de schémas
Databricks décode chaque message avec les fonctions from_json, from_avro et from_protobuf de Spark. Les comportements suivants s’appliquent, que vous déclariez le schéma en ligne ou que vous le résolviez à partir d’un registre de schémas :
- Enregistrements mal formés. Le décodage utilise le mode
PERMISSIVE; ainsi, un enregistrement qui ne correspond pas à son schéma est décodé en une valeur nulle au lieu de faire échouer le Stream. - Unions Avro. Une union de plusieurs types d’enregistrement est décodée en une structure avec un champ par type d’enregistrement, chacun nommé d’après son enregistrement Avro.
- Types Protobuf. Les entiers non signés sont décodés vers un type signé plus large (par exemple,
uint32versBIGINTetuint64versDECIMAL(20,0)), les champs enum sont décodés vers leur nom de chaîne, et les types wrapper scalaires (par exemple,StringValueetInt32Value) sont décodés vers une colonne nullable du type encapsulé.
Registre de schémas
Les registres de schémas stockent et versionnent les schémas utilisés par les producteurs et les consommateurs de streaming, en appliquant des règles de compatibilité à mesure que ces schémas évoluent. Lorsqu'un registre de schémas externe est configuré, le Magasin de fonctionnalités lit le schéma à partir du registre et l'utilise pour décoder le message en streaming. Vous ne déclarez pas le schéma en ligne sur le Stream lorsque vous utilisez un registre de schémas.
La prise en charge du registre de schémas présente les limitations suivantes :
- Seul Confluent Schema Registry est pris en charge
- Seuls les formats Avro et Protobuf sont pris en charge. Pour lire des messages JSON, déclarez plutôt le schéma en ligne. Voir Schéma JSON.
- Chaque Stream est connecté à exactement un sujet Confluent pour la valeur du message, et un pour la clé du message (si fournie). Stream de sujets contenant plusieurs enregistrements de schéma n’est pas une configuration prise en charge. Si votre Stream se connecte à des sujets contenant plusieurs schémas, les enregistrements qui ne correspondent pas au schéma du sujet spécifié sont décodés comme null.
Se connecter à un registre de schémas
Fournissez les détails de connexion au registre en tant qu'options sur la connexion Kafka Unity Catalog, et stockez le secret de l'API du registre dans un secret scope Databricks. L'identité d'exécution du Stream doit disposer de l'autorisation READ sur le secret scope, car le pipeline d'ingestion lit le secret au moment de l'exécution. Pour savoir comment créer et configurer une connexion, consultez Créer une connexion.
Ajoutez les options schema_registry_url, schema_registry_api_key et schema_registry_api_secret à la connexion utilisée pour l’authentification. L’exemple suivant crée une connexion Kafka qui s’authentifie auprès du broker avec une information d’identification de service Unity Catalog et auprès du registre avec une clé API :
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>')
)
Définissez à la fois l’option schema_registry_api_secret sur la connexion Kafka et la référence Secret Scope sur le Stream vers le même secret.
Créer un stream qui utilise un registre de schémas
Transmettez un SchemaRegistryConfig en tant que schema_config. Référencez le secret de l'API de registre avec api_secret_ref, et identifiez le sujet et le format avec payload_schema_locator pour la valeur du message, ou key_schema_locator pour la clé du message. Au moins un localisateur doit être fourni.
Notez ici les différences par rapport aux exemples de schéma directs dans la section Configuration du schéma. Lorsque vous utilisez un registre de schémas, vous ne fournissez pas le schéma en ligne sur le Stream vers schema_config. À la place, vous spécifiez un SchemaRegistryConfig qui identifie le schéma dans le registre.
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"
),
),
)
Un sujet Confluent est le scope nommé sous lequel l'historique des versions d'un schéma est enregistré et la compatibilité est appliquée. Définissez subject sur le nom du scope concerné, qui est généralement déterminé à partir de la stratégie de nom de sujet:
- TopicNameStrategy (default, dérive le sujet du nom du sujet) :
<topic>-valuepour la valeur et<topic>-keypour la clé. Par exemple, le schéma de valeur pour le sujettransactionsutilise le sujettransactions-value. - RecordNameStrategy (dérive le sujet du nom d’enregistrement du schéma, indépendamment du sujet) : le nom d’enregistrement entièrement qualifié, tel que
com.example.Payment. Il s’agit de l’espace de noms et du nom de l’enregistrement pour Avro, ou du package et du nom du message pour Protobuf. - TopicRecordNameStrategy (combine les noms de sujet et d’enregistrement) :
<topic>-<fully-qualified-record-name>, tel quetransactions-com.example.Payment.
format est obligatoire. Définissez-le sur SchemaLocatorFormat.FORMAT_AVRO ou SchemaLocatorFormat.FORMAT_PROTOBUF pour correspondre à la manière dont le sujet est sérialisé.
évolution des schémas
Le pipeline d’ingestion résout le schéma actuel du sujet lorsqu’il start. Lorsque vous enregistrez une nouvelle version de schéma compatible avec les versions antérieures pour le sujet dans le registre de schémas, le pipeline en cours d’exécution continue d’utiliser la version avec laquelle il a démarré.
Comme Databricks gère le pipeline d’ingestion en tant que pipeline Lakeflow serverless, le pipeline redémarre périodiquement. Lors de son prochain redémarrage, il prend en compte la nouvelle version du schéma. Il peut s’écouler jusqu’à une semaine avant que les champs nouveaux ou modifiés n’apparaissent dans la table d’ingestion.
Pour savoir comment le pipeline gère les enregistrements qui ne correspondent pas au schéma qu’il utilise actuellement, consultez Décodage des données à l’aide de schémas.
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 :