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 remblai 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é.

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. Le Feature Store repose sur un SDP serverless, ce qui signifie que vous aurez besoin d’un mécanisme pour connecter votre calcul classique (broker ou terminaison) au calcul serverless de Databrick. Cela se fait via des produits comme privatelink ou en permettant à votre calcul classique d’être accessible depuis l’internet public.

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 : Précise la plateforme de streaming et les détails spécifiques à la source, comme l’abonnement au sujet pour une source 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 remblai pour plus de détails.

Pour la configuration propre à la source source_config et la configuration de la connexion, ainsi qu’un exemple complet de create_stream(), consultez Apache Kafka. Les options de schéma et d’ingestion sont partagées entre les sources.

Apache Kafka

Pour diffuser depuis Apache Kafka, utilisez KafkaStreamConfig comme configuration source et une connexion Unity Catalog pour l’authentification. Voir le streaming sur le calcul sans serveur et Connectez-vous à Apache Kafka pour la connectivité à Kafka.

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

Modes d’abonnement Kafka

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]}')

Authentification Kafka

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 à partir du flux 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 à la fois l’authentification IAM (identifiants de service) et l’authentification SASL.

IAM (identifiant de service)

Authentifier avec une certification de service Unity Catalog, par exemple pour se connecter à Amazon MSK via IAM. Pour créer des informations d’identification de service, consultez Créer des informations d’identification de service. Définir le nom de la certification de service avec l’option credential suivante :

CREATE CONNECTION IF NOT EXISTS `my-kafka-connection`
TYPE KAFKA
OPTIONS (
    bootstrap_servers '<bootstrap_servers>',
    credential '<service_credential>'
)

Outre USE CONNECTION sur la connexion, les identités qui utilisent les informations d’identification du service doivent avoir ACCESS sur celles-ci. Activer l’octroi ACCESS sur la certification de service référencée au créateur du flux et à toute identité qui matérialise des fonctionnalités avec le flux. Consultez Autoriser l'utilisation d'un justificatif de service pour 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 options suivantes :

  • PLAIN
  • SCRAM-SHA-256
  • SCRAM-SHA-512

Fournissez les identifiants à l’aide des options user et password. La connexion stocke ces identifiants de manière sécurisée.

L’exemple suivant utilise SASL/SCRAM. Pour SASL/PLAIN, régler 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>'
)

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

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.
  • Types 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 :

  • Pris en charge uniquement pour les flux Kafka.
  • 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é exécuter en tant que du flux doit disposer de l’autorisation READ sur l’étendue 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>')
)

Définissez à la fois l’option schema_registry_api_secret sur la connexion Kafka et la référence d’étendue de secrets sur le stream sur le même secret.

Créez un flux utilisant un registre de schéma

Transmettez un SchemaRegistryConfig en tant que 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 remblayage

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 lance un pipeline d’ingestion géré qui lit en continu les messages du flux source et les écrit dans une table Delta (la table d’ingestion). Le pipeline part de la position la plus récente dans la source et s’exécute en continu, et ne capture que 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 du message ainsi que les colonnes de métadonnées. Les colonnes communes sont présentes pour chaque source ; les kafka_* colonnes ne sont présentes que pour un ruisseau Kafka.

Column Type Source Description
key Variable (à partir de key_schema) Common La clé du message, structurée selon le schéma que vous avez fourni.
value Variable (à partir de payload_schema) Common La valeur du message (payload), structurée selon le schéma que vous avez fourni.
stream_record_timestamp TIMESTAMP Common L’horodatage de l’enregistrement. Pour les données de remplissage en avant, il s’agit de l’horodatage d’ingestion de la source. Pour les données de remplissage, il s’agit de données fournies par le client.
record_source STRING Common Soit "stream" (remplissage direct depuis le flux en direct), soit "backfill" (depuis la source de recharge).
kafka_topic STRING Kafka Rubrique Kafka à partir duquel l’enregistrement a été consommé.
kafka_partition INT Kafka Partition Kafka à partir delaquelle l’enregistrement a été consommé.
kafka_offset LONG Kafka Offset Kafka de l’enregistrement au sein de sa partition.

Source de remplissage

Comme le pipeline de remplissage à partir de l’avant démarre à partir de la position la plus récente dans la source, 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 sont transmises si elles sont présentes sur la source de remplissage, ou sont définies sur NULL autrement. Pour Kafka, il s’agit de kafka_topic, kafka_partition et kafka_offset.

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.

Comme le pipeline d’ingestion commence à la dernière position dans la source, il ne capture les messages arrivant qu’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