Configurer un flux

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-client paquet 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

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 PERMISSIVE mode, 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, uint32 vers BIGINT et uint64 vers DECIMAL(20,0)), les champs enum décodent vers leur nom de chaîne, et les types d’enveloppe scalaire (par exemple, StringValue et Int32Value) 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>-value pour la valeur et <topic>-key pour la clé. Par exemple, le schéma de valeurs pour le sujet transactions utilise le sujet transactions-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>, comme transactions-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 :

  • SELECT sur la table d’ingestion accorde l’accès en lecture au flux.
  • MANAGE sur 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_timestamp dans la source de remplissage.
  • Ingestion min: nombre minimal stream_record_timestamp de 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_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 par défaut est la combinaison complète de key, valueet stream_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 :

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

Obtenir un ordinateur portable