Remarque
L’accès à cette page requiert une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page requiert 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 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-clientpaquet 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
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 à 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 :
PLAINSCRAM-SHA-256SCRAM-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
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.
- Types 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 :
- 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>-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 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 :
-
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 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_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.
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_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 :