Remarque
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de modifier des répertoires.
Important
Cette fonctionnalité est disponible en préversion publique. Les administrateurs d’espace de travail peuvent contrôler l’accès à cette fonctionnalité à partir de la page Aperçus . Consultez Gérer les préversions d’Azure Databricks.
Un flux représente une source de données de streaming externe, telle qu’Apache Kafka. Les flux stockent les détails de connexion, l’authentification, les schémas et la configuration d’ingestion. Une fois qu’un flux est créé, vous pouvez le référencer à l’aide de définitions d’affichage des fonctionnalités pour créer des fonctionnalités de streaming en temps réel.
Les flux ont des noms en trois parties (catalog.schema.stream_name). L’accès à un flux est régi par sa table d’ingestion associée. Voir Ingestion et rétroremplissage pour plus de détails.
Exigences
- Pour exécuter des commandes de bloc-notes : sans serveur ou un cluster de calcul classique exécutant Databricks Runtime 17.0 ML ou une version ultérieure.
- Le
feature-engineering-clientpaquet Python version 0.17.0 ou supérieure doit être installé.
Créer un flux
Utilisez create_stream() pour créer un nouveau Stream. Un flux nécessite quatre composants de configuration :
- Configuration de la source : spécifie la plateforme de diffusion en continu (par exemple, Kafka) et les détails spécifiques à la source (tels que l’abonnement à la rubrique pour Kafka).
- Configuration de la connexion : spécifie comment se connecter et s’authentifier auprès de la plateforme de diffusion en continu, y compris les serveurs de démarrage et les informations d’identification.
- Configuration du schéma : définit la structure des clés et des valeurs de message.
- Configuration d’ingestion : spécifie où et comment les données de flux sont ingérées. Voir Ingestion et rétroremplissage 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 flux
Avant de définir des fonctionnalités de diffusion en continu, connectez-vous et testez une connexion de pipeline Lakeflow en streaming à votre répartiteur Kafka. Consultez Streaming sur un calcul sans serveur et Connexion à Apache Kafka.
Pour le streaming géré d’AWS (Amazon MSK), consultez la connectivité privée sans serveur vers Amazon MSK. Pour plus d’informations sur les options d’authentification Kafka, consultez Authentification.
Authentication
Connexion du catalogue Unity (recommandée)
Utilisez une connexion de catalogue Unity 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, consultez Créer une connexion. Le créateur du Stream doit avoir USE CONNECTION sur la connexion. Tout utilisateur créant des fonctionnalités matérialisées à partir du Stream comme source doit également disposer de USE CONNECTION sur la connexion.
connection_config = StreamConnectionConfig(
uc_connection_name="my-kafka-connection"
)
Fournisseur de solutions cloud (mTLS) direct
Pour l’authentification mTLS directe, fournissez des fichiers keystore et truststore stockés sur un volume Unity Catalog, avec des mots de passe référencés à l’aide de scopes de secrets Databricks. Pour plus d’informations sur l’authentification SSL avec Kafka, consultez Utiliser SSL pour se connecter Azure 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 (SASL/SCRAM et SASL/PLAIN) n’est pas prise en charge pendant la préversion.
Modes d’abonnement
Le mode d’abonnement spécifie la façon dont le flux sélectionne les rubriques Kafka à utiliser. Trois modes sont pris en charge :
| Mode | Description | Example |
|---|---|---|
subscribe |
Liste séparée par des virgules des noms de rubriques | KafkaSubscriptionMode(subscribe="topic1,topic2") |
subscribe_pattern |
Modèles regex Java correspondant aux noms de rubriques correspondantes | KafkaSubscriptionMode(subscribe_pattern="events-.*") |
assign |
JSON spécifiant les affectations sujet-partition | KafkaSubscriptionMode(assign='{"my-topic": [0, 1, 2]}') |
Configuration du schéma
Définir la structure des clés et des valeurs de message afin que les définitions d’ingestion et de caractéristiques puissent lire les champs individuels. Pour les sources Kafka, payload_schema correspond à la valeur du message Kafka (le value modèle clé-valeur de Kafka) et key_schema correspond à la clé de message Kafka. Au moins un des éléments payload_schema ou key_schema doit être fourni.
Chacun SchemaConfig accepte l’un des trois formats, correspondant à la façon 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, il est traité comme une chaîne simple.
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, voir Registre de schéma pour plus de détails.
Schéma JSON
Fournir une chaîne de 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
Fournir une chaîne de schémas Avro à avro_schema. Les types logiques Avro sont pris en charge, y compris 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 ainsi que le nom du message de charge utile. Importez ProtoSchemaSpec depuis databricks.feature_engineering.entities.
message_name doit être le nom du message pleinement qualifié, incluant le package déclaré dans le .proto texte (par exemple, com.example.Event, non Event). Les syntaxes proto2 et proto3 sont toutes deux prises en charge.
google.protobuf.Timestamp et les types d’enveloppes 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 ces valeurs comme un scalaire ou un message supporté à la place. Les types scalaires fixed32 et fixed64, ainsi que map avec des clés qui ne sont pas des chaînes, 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 de from_jsonSpark, from_avro, et from_protobuf . 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éma :
- Enregistrements mal formés. Le décodage utilise ce
PERMISSIVEmode, donc un enregistrement qui ne correspond pas à son schéma se décode vers une valeur nulle au lieu de faire échouer le flux. - Les syndicats Avro. Une union de plusieurs types d’enregistrements décode en une structure avec un champ par type d’enregistrement, chacun nommé d’après son enregistrement Avro.
- Du type Protobuf. Les entiers non signés décodent vers un type signé plus large (par exemple,
uint32versBIGINTetuint64versDECIMAL(20,0)), les champs enum décodent vers leur nom de chaîne, et les types d’enveloppe scalaire (par exemple,StringValueetInt32Value) décodent vers une colonne annulable du type enroulé.
Registre de schémas
Les registres de schémas stockent et versionnent les schémas utilisés par les producteurs et consommateurs de streaming, en appliquant les règles de compatibilité au fur et à mesure que ces schémas évoluent. Lorsqu’un registre de schéma externe est configuré, le Feature Store lit le schéma du registre et l’utilise pour décoder le message en flux. Vous ne déclarez pas le schéma en ligne sur le Stream lorsque vous utilisez un registre de schéma.
La prise en charge du registre de schémas présente les limitations suivantes :
- Seul le Registre de Schéma Confluent est pris en charge
- Seuls les formats Avro et Protobuf sont pris en charge. Pour lire les messages JSON, déclarez le schéma en ligne à la place. Voir le schéma JSON.
- Chaque flux est connecté à exactement un sujet Confluent pour la valeur du message, et un autre pour la clé du message (si disponible). Les sujets de flux contenant plusieurs enregistrements de schéma ne sont 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 nulls.
Connectez-vous à un registre de schéma
Indiquer les paramètres de connexion au registre comme options de la connexion Kafka Unity Catalog, et stocker le secret d’API du registre dans un Databricks secret scope. L’identité d’exécution du Stream doit disposer de l’autorisation READ sur la portée du secret, car le pipeline d’ingestion lit le secret au moment de l’exécution. Pour savoir comment créer et configurer une connexion, voir 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 courtier avec un identifiant 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>')
)
Réglez à la fois l’option schema_registry_api_secret sur la connexion Kafka et la référence secrète de portée sur le Stream au même secret.
Créez un flux utilisant un registre de schéma
Passez a SchemaRegistryConfig comme le schema_config. Référez le secret de l’API du registre avec api_secret_ref, et identifiez l’objet 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 les différences ici par rapport aux exemples directs de schémas dans la section Configuration des schémas . Lorsque vous utilisez un registre de schémas, vous ne fournissez pas directement le schéma dans le flux à 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 champ 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 nommage du sujet :
-
TopicNameStrategy (par défaut, dérive le sujet du nom du sujet) :
<topic>-valuepour la valeur et<topic>-keypour la clé. Par exemple, le schéma de valeurs 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 pleinement qualifié, tel que
com.example.Payment. Il s’agit de l’espace de noms et du nom d’un enregistrement pour Avro, ou du package et du nom du message pour Protobuf. -
TopicRecordNameStrategy (combine le nom du topic et le nom de l’enregistrement) :
<topic>-<fully-qualified-record-name>, commetransactions-com.example.Payment.
format est obligatoire. Réglez-le sur SchemaLocatorFormat.FORMAT_AVRO ou SchemaLocatorFormat.FORMAT_PROTOBUF pour qu’il corresponde à la façon dont le sujet est sérialisé.
Évolution du schéma
Le pipeline d’ingestion résout le schéma courant du sujet au démarrage. Lorsque vous enregistrez une nouvelle version de schéma rétrocompatible sur le sujet dans le registre de schéma, le pipeline en cours continue d’utiliser la version avec laquelle il a commencé.
Parce que Databricks gère le pipeline d’ingestion comme un pipeline Lakeflow sans serveur, le pipeline redémarre périodiquement. Lors de son redémarrage suivant, il détecte la nouvelle version du schéma. Il peut falloir jusqu’à une semaine pour que de nouveaux champs ou des champs modifiés apparaissent dans la table d’ingestion.
Pour la manière dont le pipeline gère les enregistrements qui ne correspondent pas au schéma qu’il utilise actuellement, voir Décodage des données à l’aide de schémas.
Ingestion et rattrapage
Le ingestion_config paramètre configure la façon dont les données de flux sont capturées et stockées pour l’entraînement et le service.
L’accès à un flux est régi par la table d’ingestion :
-
SELECTsur la table d’ingestion accorde l’accès en lecture au flux. -
MANAGEsur la table d’ingestion accorde l’accès à la suppression.
Pour plus d’informations sur les privilèges de table, consultez Table et la référence des privilèges Unity Catalog.
Pipeline d’ingestion
Lorsqu’un flux est créé, Databricks démarre un pipeline d’ingestion managé qui lit en continu les messages de la rubrique Kafka et les écrit dans une table Delta (la table d’ingestion). Le pipeline démarre à partir du dernier offset Kafka et s’exécute en continu, en capturant uniquement les nouveaux messages qui arrivent après la création du flux. Cette table d’ingestion est utilisée pour l’entraînement avec les fonctionnalités de streaming. Lorsqu’un flux est supprimé, son pipeline d’ingestion et sa table d’ingestion sont également supprimés.
Destination d’ingestion
Le paramètre ingestion_destination spécifie le nom en trois parties de la table Delta dans laquelle les données de flux sont écrites.
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
)
Schéma de la table d’ingestion
La table d’ingestion contient les données de message ainsi que les colonnes de métadonnées :
| Column | Type | Description |
|---|---|---|
key |
Variable (à partir de key_schema) |
Clé de message Kafka, structurée en fonction du schéma que vous avez fourni. |
value |
Variable (à partir de payload_schema) |
Valeur du message Kafka (charge utile), structurée en fonction du schéma que vous avez fourni. |
stream_record_timestamp |
TIMESTAMP |
L’horodatage de l’enregistrement. Pour les données de forward-fill, il s’agit de l’horodatage d’ingestion du broker Kafka. Pour les données de remplissage, il s’agit de données fournies par le client. |
kafka_topic |
STRING |
Rubrique Kafka à partir duquel l’enregistrement a été consommé. |
kafka_partition |
INT |
Partition Kafka à partir delaquelle l’enregistrement a été consommé. |
kafka_offset |
LONG |
Offset Kafka de l’enregistrement au sein de sa partition. |
record_source |
STRING |
Soit "stream" (préremplissage à partir du flux Kafka en direct), soit "backfill" (à partir de la source de reprise). |
Source de remplissage
Comme le pipeline de remplissage progressif démarre à partir du dernier offset Kafka, il ne capture pas les messages qui existaient avant la création du flux. Pour fournir une couverture des données historiques pour l’apprentissage, configurez une source de remplissage facultative.
Lorsqu’une source de renvoi est configurée, Databricks exécute une tâche ponctuelle MERGE INTO qui copie les lignes de renvoi vers la table d’ingestion avec record_source="backfill". La fusion s’exécute uniquement après que le vérificateur de chevauchement confirme que la source de remplissage et le flux de remplissage avant ont des horodatages qui se chevauchent (voir Chevauchement entre le remplissage et les données de flux en direct). Si la condition de chevauchement n’est pas remplie au bout de 2 jours, la FUSION s’exécute de toute façon pour éviter un blocage indéfini.
La table de remplissage doit inclure une stream_record_timestamp colonne 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 ou définies 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 et les données de flux en direct
Avant d’exécuter une FUSION entre le renvoi et la table d’ingestion, une vérification de chevauchement compare les horodatages dans les deux tables :
-
Remplissage maximal : maximum
stream_record_timestampdans la source de remplissage. -
Ingestion min: nombre minimal
stream_record_timestampde lignes (record_source="stream") dans la table d’ingestion.
La FUSION a lieu lorsque le plus récent horodatage du renvoi est supérieur d’au moins 1 heure au plus ancien horodatage de la table d’ingestion. Ce chevauchement garantit qu’il n’y a pas d’écarts dans la table d’ingestion. Si la condition de chevauchement n’est pas remplie au bout de 2 jours, la FUSION s’exécute de toute façon 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 qui arrivent après la création du flux. Votre source de rétroremplissage doit contenir des données qui couvrent la période d’ingestion, et non pas seulement jusqu’à l’heure de création du flux.
Par exemple, si vous créez un flux à 15 h 00, le pipeline de propagation commence à lire les messages à partir de 15 h 00. Votre source de renvoi doit inclure des données horodatées jusqu’à au moins 16 h 00 (1 heure après le début du remplissage vers l’avant) afin de satisfaire à la vérification du chevauchement. Cela signifie que vous devez mettre à jour votre table de renvoi après 17 h 00 afin de vous assurer que la table d’ingestion ne présente aucune lacune.
Deduplication
Utilisez deduplication_columns pour spécifier des chemins d’accès aux colonnes afin d’identifier les lignes en double lors de l’ingestion entre les données de flux de renvoi et de transfert. Utilisez la notation par points pour les champs imbriqués (par exemple, "value.user_id").
Choisissez des colonnes de déduplication en fonction de vos données :
- Si chaque enregistrement de votre flux contient un identificateur unique (par exemple,
value.transaction_id), utilisez cette colonne pour la déduplication. - Si votre source de backfill contient 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 par défaut est la combinaison complète de
key,valueetstream_record_timestamp. Cela n’est pas recommandé, car ces critères stricts peuvent 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 flux
Obtenir un flux
stream = client.get_stream(name="my_catalog.my_schema.my_stream")
Répertorier les flux
streams = client.list_streams(
catalog_name="my_catalog",
schema_name="my_schema",
max_results=50,
include_schemas=False,
)
Définissez cette option include_schemas=True pour inclure les détails complets du schéma. Les schémas peuvent être volumineux et cela peut entraîner une opération de longue durée. Pour récupérer des schémas individuellement à la place, utilisez get_stream.
Supprimer un flux
La suppression d’un flux supprime également son pipeline d’ingestion et sa table d’ingestion.
Warning
Tous les modèles ou fonctionnalités qui référencent le flux supprimé n’ont plus accès aux données de flux sous-jacents. Créez une copie de la table d’ingestion avant la suppression si vous avez besoin de ces données, mais n’avez plus besoin du flux.
client.delete_stream(name="my_catalog.my_schema.my_stream")
Exemple de bloc-notes
Pour obtenir un exemple de bout en bout qui crée un flux, définit des fonctionnalités de diffusion en continu et se déploie sur un point de terminaison de service, consultez le notebook suivant :