Aller au contenu principal

Configurer un Stream

info

Aperçu

Cette fonctionnalité est en Aperçu public. Les administrateurs du Workspace peuvent contrôler l'accès à cette fonctionnalité à partir de la page Previews . Consultez Gérer les aperçus Databricks.

Un Stream représente une source de données streaming externe, telle qu’Apache Kafka. Les Streams stockent les détails de connexion, l'authentification, les schémas et la configuration d'ingestion. Une fois qu'un Stream est créé, vous pouvez le référencer à l'aide des définitions de Feature View pour créer des fonctionnalités de streaming en temps réel.

Les Stream ont des noms en trois parties (catalog.schema.stream_name). L'accès à un Stream est régi par sa table d'ingestion associée. Consultez Ingestion et remplissage pour plus de détails.

Exigences

  • Pour exécuter des commandes de notebook : serverless ou un cluster de compute classique exécutant Databricks Runtime 17.0 ML ou une version ultérieure.
  • Le package Python feature-engineering-client en version 0.17.0 ou supérieure doit être installé.

Créer un Stream

Utilisez create_stream() pour créer un nouveau Stream. Un Stream nécessite quatre composants de configuration :

  • Configuration de la source : Spécifie la plateforme de streaming (par exemple, Kafka) et les détails spécifiques à la source (tels que l'abonnement aux rubriques pour Kafka).
  • Configuration de la connexion : spécifie comment se connecter et s’authentifier à la plateforme de streaming, y compris les serveurs d’amorçage et les informations d’identification.
  • Configuration du schéma : Définit la structure des clés et des valeurs de message.
  • Configuration de l'ingestion : Spécifie où et comment les données Stream sont ingérées. Consultez Ingestion et remplissage pour plus de détails.
Python
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
KafkaStreamConfig,
KafkaSubscriptionMode,
StreamConnectionConfig,
DirectSchemas,
SchemaConfig,
IngestionConfig,
IngestionDestination,
StreamBackfillSource,
)

client = FeatureEngineeringClient()

stream = client.create_stream(
name="my_catalog.my_schema.my_stream",
source_config=KafkaStreamConfig(
subscription_mode=KafkaSubscriptionMode(subscribe="events-topic"),
),
connection_config=StreamConnectionConfig(
uc_connection_name="my-kafka-connection"
),
schema_config=DirectSchemas(
payload_schema=SchemaConfig(
json_schema=(
'{'
' "type": "object",'
' "properties": {'
' "transaction_id": {"type": "string"},'
' "user_id": {"type": "string"},'
' "amount": {"type": "number"},'
' "event_time": {"type": "string", "format": "date-time"}'
' }'
'}'
)
),
),
ingestion_config=IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
),
)

Connexion à des sources de Stream

Avant de définir les fonctionnalités de streaming, connectez et testez une connexion de LakeFlow Pipelines de streaming à votre courtier Kafka. Consultez Streaming sur compute serverless et Se connecter à Apache Kafka.

Pour le streaming géré par AWS (Amazon MSK), consultez Connectivité privée Serverless à Amazon MSK. Pour plus de détails sur les options d'authentification Kafka, consultez Authentification.

Authentification

Connexion Unity Catalog (recommandée)

Utilisez une connexion Unity Catalog pour vous authentifier auprès de votre cluster Kafka. Il s'agit de l'approche recommandée pour l'authentification gérée. Pour créer une connexion, consultez Créer une connexion. Le créateur du Stream doit disposer de USE CONNECTION sur la connexion. Tout utilisateur matérialisant des fonctionnalités avec le Stream comme source doit également disposer de USE CONNECTION sur la connexion.

Python
connection_config = StreamConnectionConfig(
uc_connection_name="my-kafka-connection"
)

La connexion prend en charge l’authentification IAM (identifiant de service) et SASL.

IAM (identifiant de service)

Authentifiez-vous avec un identifiant de service Unity Catalog, par exemple pour vous connecter à Amazon MSK avec IAM. Pour créer un identifiant de service, voir Créer des identifiants de service. Définissez le nom de l’identifiant de service avec l’option credential :

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

En plus de USE CONNECTION sur la connexion, les identités qui utilisent l’identifiant de service ont besoin de ACCESS sur celle-ci. Accordez ACCESS sur l'identifiant de service référencé au créateur du Stream et à toute identité qui matérialise des features avec le Stream. Voir Accorder des autorisations pour utiliser un identifiant de service afin d'accéder à un service cloud externe.

SASL

L'authentification SASL utilise un nom d'utilisateur et un mot de passe. Définissez sasl_mechanism sur l'une des valeurs suivantes :

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

Fournissez les identifiants avec les options user et password. La connexion stocke ces identifiants en toute sécurité.

L’exemple suivant utilise SASL/SCRAM. Pour SASL/PLAIN, définissez sasl_mechanism sur PLAIN.

SQL
CREATE CONNECTION IF NOT EXISTS `my-kafka-connection`
TYPE KAFKA
OPTIONS (
bootstrap_servers '<bootstrap_servers>',
sasl_mechanism 'SCRAM-SHA-512',
user '<username>',
password '<password>'
)

mTLS direct

Pour l'authentification mTLS directe, fournissez les fichiers de keystore et de truststore stockés sur un volume Unity Catalog, avec des mots de passe référencés via les Secret Scopes Databricks. Pour plus d'informations sur l'authentification SSL avec Kafka, voir Utiliser SSL pour connecter Databricks à Kafka.

Python
from databricks.feature_engineering.entities import (
DirectMtlsConfig,
MtlsConfig,
SecretScopeReference,
)

connection_config = DirectMtlsConfig(
bootstrap_servers="broker1:9092,broker2:9092",
mtls_config=MtlsConfig(
keystore_location="/Volumes/my_catalog/my_schema/my_volume/keystore.jks",
keystore_password_ref=SecretScopeReference(
scope="my_scope", key="keystore_password"
),
key_password_ref=SecretScopeReference(
scope="my_scope", key="key_password"
),
truststore_location="/Volumes/my_catalog/my_schema/my_volume/truststore.jks",
truststore_password_ref=SecretScopeReference(
scope="my_scope", key="truststore_password"
),
),
)

Modes d'abonnement

Le mode d'abonnement spécifie comment le Stream sélectionne les rubriques Kafka à consommer. Trois modes sont pris en charge :

Mode

Description

Exemple

subscribe

Liste de noms de rubriques séparés par des virgules

KafkaSubscriptionMode(subscribe="topic1,topic2")

subscribe_pattern

Noms de sujets correspondant aux modèles d'expressions régulières Java.

KafkaSubscriptionMode(subscribe_pattern="events-.*")

assign

JSON spécifiant les attributions de partitions de rubrique

KafkaSubscriptionMode(assign='{"my-topic": [0, 1, 2]}')

Mode

Description

Exemple

subscribe

Liste de noms de rubriques séparés par des virgules

KafkaSubscriptionMode(subscribe="topic1,topic2")

subscribe_pattern

Noms de sujets correspondant aux modèles d'expressions régulières Java.

KafkaSubscriptionMode(subscribe_pattern="events-.*")

assign

JSON spécifiant les attributions de partitions de rubrique

KafkaSubscriptionMode(assign='{"my-topic": [0, 1, 2]}')

Configuration du schéma

Définissez la structure des clés et des valeurs de message afin que les définitions d'ingestion et de fonctionnalités puissent lire les champs individuels. Pour les sources Kafka, payload_schema correspond à la valeur du message Kafka (le value dans le modèle clé-valeur de Kafka) et key_schema correspond à la clé du message Kafka. Au moins l'un des éléments payload_schema ou key_schema doit être fourni.

Chaque SchemaConfig accepte l'un des trois formats, correspondant à la manière dont la source sérialise ses messages : json_schema, avro_schema ou proto_schema. Si aucun schéma n'est fourni pour une clé ou une charge utile, celle-ci est traitée comme une simple chaîne de caractères.

Les exemples de code de cette section utilisent des schémas déclarés en ligne avec DirectSchemas, où le schéma est fourni sous forme de chaîne. Pour gérer les schémas à l’aide d’un registre de schémas externe, consultez Schema registry pour plus de détails.

Schéma JSON

Fournissez une chaîne Schéma JSON à json_schema.

Python
schema_config = DirectSchemas(
payload_schema=SchemaConfig(
json_schema=(
'{'
' "type": "object",'
' "properties": {'
' "user_id": {"type": "string"},'
' "amount": {"type": "number"},'
' "event_time": {"type": "string"}'
' }'
'}'
)
),
key_schema=SchemaConfig(
json_schema='{"type": "string"}'
),
)

Schéma Avro

Fournissez une chaîne de schéma Avro à avro_schema. Les types logiques Avro sont pris en charge, notamment timestamp-millis, date et decimal.

Python
schema_config = DirectSchemas(
payload_schema=SchemaConfig(
avro_schema=(
'{'
' "type": "record",'
' "name": "Event",'
' "fields": ['
' {"name": "user_id", "type": "string"},'
' {"name": "amount", "type": "double"},'
' {"name": "event_time",'
' "type": {"type": "long", "logicalType": "timestamp-millis"}}'
' ]'
'}'
)
),
)

Schéma Protobuf

Fournissez un ProtoSchemaSpec à proto_schema avec le texte source Protocol Buffers .proto et le nom du message de charge utile. Importer ProtoSchemaSpec depuis databricks.feature_engineering.entities.

message_name doit être le nom de message entièrement qualifié, incluant le package déclaré dans le texte .proto (par exemple, com.example.Event, et non Event). Les syntaxes proto2 et proto3 sont toutes deux prises en charge.

google.protobuf.Timestamp et les types de wrapper scalaires (StringValue, Int32Value, etc.) sont pris en charge, et leurs importations sont résolues automatiquement. D’autres types bien connus, tels que Duration, Struct et Any, sont rejetés ; encodez plutôt ces valeurs sous forme de scalaire ou de message pris en charge. Les types scalaires fixed32 et fixed64 ainsi que map avec des clés non-chaîne ne sont pas non plus pris en charge.

Python
from databricks.feature_engineering.entities import ProtoSchemaSpec

schema_config = DirectSchemas(
payload_schema=SchemaConfig(
proto_schema=ProtoSchemaSpec(
schema_text=(
'syntax = "proto3";\n'
'package com.example;\n'
'import "google/protobuf/timestamp.proto";\n'
'message Event {\n'
' string user_id = 1;\n'
' double amount = 2;\n'
' google.protobuf.Timestamp event_time = 3;\n'
'}'
),
message_name="com.example.Event",
)
),
)

Décodage des données à l’aide de schémas

Databricks décode chaque message avec les fonctions from_json, from_avro et from_protobuf de Spark. Les comportements suivants s’appliquent, que vous déclariez le schéma en ligne ou que vous le résolviez à partir d’un registre de schémas :

  • Enregistrements mal formés. Le décodage utilise le mode PERMISSIVE ; ainsi, un enregistrement qui ne correspond pas à son schéma est décodé en une valeur nulle au lieu de faire échouer le Stream.
  • Unions Avro. Une union de plusieurs types d’enregistrement est décodée en une structure avec un champ par type d’enregistrement, chacun nommé d’après son enregistrement Avro.
  • Types Protobuf. Les entiers non signés sont décodés vers un type signé plus large (par exemple, uint32 vers BIGINT et uint64 vers DECIMAL(20,0)), les champs enum sont décodés vers leur nom de chaîne, et les types wrapper scalaires (par exemple, StringValue et Int32Value) sont décodés vers une colonne nullable du type encapsulé.

Registre de schémas

Les registres de schémas stockent et versionnent les schémas utilisés par les producteurs et les consommateurs de streaming, en appliquant des règles de compatibilité à mesure que ces schémas évoluent. Lorsqu'un registre de schémas externe est configuré, le Magasin de fonctionnalités lit le schéma à partir du registre et l'utilise pour décoder le message en streaming. Vous ne déclarez pas le schéma en ligne sur le Stream lorsque vous utilisez un registre de schémas.

La prise en charge du registre de schémas présente les limitations suivantes :

  • Seul Confluent Schema Registry est pris en charge
  • Seuls les formats Avro et Protobuf sont pris en charge. Pour lire des messages JSON, déclarez plutôt le schéma en ligne. Voir Schéma JSON.
  • Chaque Stream est connecté à exactement un sujet Confluent pour la valeur du message, et un pour la clé du message (si fournie). Stream de sujets contenant plusieurs enregistrements de schéma n’est pas une configuration prise en charge. Si votre Stream se connecte à des sujets contenant plusieurs schémas, les enregistrements qui ne correspondent pas au schéma du sujet spécifié sont décodés comme null.

Se connecter à un registre de schémas

Fournissez les détails de connexion au registre en tant qu'options sur la connexion Kafka Unity Catalog, et stockez le secret de l'API du registre dans un secret scope Databricks. L'identité d'exécution du Stream doit disposer de l'autorisation READ sur le secret scope, car le pipeline d'ingestion lit le secret au moment de l'exécution. Pour savoir comment créer et configurer une connexion, consultez Créer une connexion.

Ajoutez les options schema_registry_url, schema_registry_api_key et schema_registry_api_secret à la connexion utilisée pour l’authentification. L’exemple suivant crée une connexion Kafka qui s’authentifie auprès du broker avec une information d’identification de service Unity Catalog et auprès du registre avec une clé API :

SQL
CREATE CONNECTION IF NOT EXISTS `my-kafka-connection`
TYPE KAFKA
OPTIONS (
bootstrap_servers '<bootstrap_servers>',
credential '<service_credential>',
schema_registry_url 'https://<registry-host>',
schema_registry_api_key '<registry_api_key>',
schema_registry_api_secret secret('<scope>', '<key>')
)

Définissez à la fois l’option schema_registry_api_secret sur la connexion Kafka et la référence Secret Scope sur le Stream vers le même secret.

Créer un stream qui utilise un registre de schémas

Transmettez un SchemaRegistryConfig en tant que schema_config. Référencez le secret de l'API de registre avec api_secret_ref, et identifiez le sujet et le format avec payload_schema_locator pour la valeur du message, ou key_schema_locator pour la clé du message. Au moins un localisateur doit être fourni.

Notez ici les différences par rapport aux exemples de schéma directs dans la section Configuration du schéma. Lorsque vous utilisez un registre de schémas, vous ne fournissez pas le schéma en ligne sur le Stream vers schema_config. À la place, vous spécifiez un SchemaRegistryConfig qui identifie le schéma dans le registre.

Python
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
KafkaStreamConfig,
KafkaSubscriptionMode,
StreamConnectionConfig,
SchemaRegistryConfig,
SchemaLocator,
SchemaLocatorConfluentSchema,
SchemaLocatorFormat,
SecretScopeReference,
IngestionConfig,
IngestionDestination,
)

client = FeatureEngineeringClient()

stream = client.create_stream(
name="my_catalog.my_schema.my_stream",
source_config=KafkaStreamConfig(
subscription_mode=KafkaSubscriptionMode(subscribe="transactions"),
),
connection_config=StreamConnectionConfig(
uc_connection_name="my-kafka-connection"
),
schema_config=SchemaRegistryConfig(
api_secret_ref=SecretScopeReference(
scope="my_scope", key="sr_api_secret"
),
payload_schema_locator=SchemaLocator(
confluent_schema=SchemaLocatorConfluentSchema(
subject="transactions-value"
),
format=SchemaLocatorFormat.FORMAT_AVRO,
),
),
ingestion_config=IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.transactions_ingestion"
),
),
)

Un sujet Confluent est le scope nommé sous lequel l'historique des versions d'un schéma est enregistré et la compatibilité est appliquée. Définissez subject sur le nom du scope concerné, qui est généralement déterminé à partir de la stratégie de nom de sujet:

  • TopicNameStrategy (default, dérive le sujet du nom du sujet) : <topic>-value pour la valeur et <topic>-key pour la clé. Par exemple, le schéma de valeur 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 entièrement qualifié, tel que com.example.Payment. Il s’agit de l’espace de noms et du nom de l’enregistrement pour Avro, ou du package et du nom du message pour Protobuf.
  • TopicRecordNameStrategy (combine les noms de sujet et d’enregistrement) : <topic>-<fully-qualified-record-name>, tel que transactions-com.example.Payment.

format est obligatoire. Définissez-le sur SchemaLocatorFormat.FORMAT_AVRO ou SchemaLocatorFormat.FORMAT_PROTOBUF pour correspondre à la manière dont le sujet est sérialisé.

évolution des schémas

Le pipeline d’ingestion résout le schéma actuel du sujet lorsqu’il start. Lorsque vous enregistrez une nouvelle version de schéma compatible avec les versions antérieures pour le sujet dans le registre de schémas, le pipeline en cours d’exécution continue d’utiliser la version avec laquelle il a démarré.

Comme Databricks gère le pipeline d’ingestion en tant que pipeline Lakeflow serverless, le pipeline redémarre périodiquement. Lors de son prochain redémarrage, il prend en compte la nouvelle version du schéma. Il peut s’écouler jusqu’à une semaine avant que les champs nouveaux ou modifiés n’apparaissent dans la table d’ingestion.

Pour savoir comment le pipeline gère les enregistrements qui ne correspondent pas au schéma qu’il utilise actuellement, consultez Décodage des données à l’aide de schémas.

Ingestion et remplissage

Le parameter ingestion_config configure la manière dont les données de Stream sont capturées et stockées pour l'entraînement et le service.

L'accès à un Stream est régi par la table d'ingestion :

  • SELECT sur la table d'ingestion accorde un accès en lecture au Stream.
  • MANAGE sur la table d'ingestion accorde l'accès en suppression.

Pour plus d'informations sur les privilèges de table, consultez Table et Référence des privilèges Unity Catalog.

Pipeline d'ingestion

Lorsqu'un Stream est créé, Databricks start un pipeline d'ingestion géré qui lit en continu les messages du sujet Kafka et les écrit dans une table Delta (la table d'ingestion). Le pipeline start à partir du dernier offset Kafka et s'exécute en continu, capturant uniquement les nouveaux messages qui arrivent après la création du Stream. Cette table d'ingestion est utilisée pour l'entraînement avec des fonctionnalités de streaming. Lorsqu'un Stream est supprimé, son pipeline d'ingestion et sa table d'ingestion sont également supprimés.

Destination d'ingestion

Le ingestion_destination spécifie le nom de table Delta en trois parties où les données de stream sont écrites.

Python
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
)

Schéma de table d’ingestion

La table d'ingestion contient les données des messages ainsi que les colonnes de métadonnées :

Colonne

Type

Description

key

Varie (à partir de key_schema)

La clé de message Kafka, structurée selon le schéma que vous avez fourni.

value

Varie (à partir de payload_schema)

La valeur du message Kafka (payload), structurée selon le schéma que vous avez fourni.

stream_record_timestamp

TIMESTAMP

L'Timestamp du record. Pour les données de remplissage avant, il s'agit de l'ingestion du broker Kafka Timestamp. Pour les données de remplissage rétroactif, celles-ci sont fournies par le client.

kafka_topic

STRING

Le sujet Kafka d'où l'enregistrement a été consommé.

kafka_partition

INT

La partition Kafka à partir de laquelle l'enregistrement a été consommé.

kafka_offset

LONG

Le décalage Kafka de l'enregistrement au sein de sa partition.

record_source

STRING

Soit "stream" (remplissage vers l'avant à partir du Stream Kafka en direct) ou "backfill" (à partir de la source de remplissage rétroactif).

Colonne

Type

Description

key

Varie (à partir de key_schema)

La clé de message Kafka, structurée selon le schéma que vous avez fourni.

value

Varie (à partir de payload_schema)

La valeur du message Kafka (payload), structurée selon le schéma que vous avez fourni.

stream_record_timestamp

TIMESTAMP

L'Timestamp du record. Pour les données de remplissage avant, il s'agit de l'ingestion du broker Kafka Timestamp. Pour les données de remplissage rétroactif, celles-ci sont fournies par le client.

kafka_topic

STRING

Le sujet Kafka d'où l'enregistrement a été consommé.

kafka_partition

INT

La partition Kafka à partir de laquelle l'enregistrement a été consommé.

kafka_offset

LONG

Le décalage Kafka de l'enregistrement au sein de sa partition.

record_source

STRING

Soit "stream" (remplissage vers l'avant à partir du Stream Kafka en direct) ou "backfill" (à partir de la source de remplissage rétroactif).

Source de remplissage

Étant donné que le pipeline de remplissage avant start du dernier décalage Kafka, il ne capture pas les messages qui existaient avant la création du Stream. Pour fournir une couverture des données historiques pour l'entraînement, configurez une source de remplissage historique facultative.

Lorsqu’une source de remplissage est configurée, Databricks exécute un job MERGE INTO unique qui copie les lignes de remplissage dans la table d’ingestion avec record_source="backfill". Le Merge ne s'exécute qu'après que le vérificateur de chevauchement confirme que la source de remplissage rétrospectif et le Stream de remplissage prospectif ont des Timestamp qui se chevauchent (voir Chevauchement entre les données de remplissage rétrospectif et de Stream en direct). Si la condition de chevauchement n'est pas remplie dans les 2 jours, le Merge s'exécute quand même pour éviter un blocage indéfini.

La table de remplissage rétroactif doit inclure une colonne stream_record_timestamp de type TIMESTAMP dans le fuseau horaire UTC. D'autres colonnes de métadonnées Kafka (kafka_topic, kafka_partition, kafka_offset) sont transmises si elles sont présentes sur la source de remplissage rétroactif, ou définies sur NULL dans le cas contraire.

Python
from databricks.feature_engineering.entities import StreamBackfillSource

ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
backfill_source=StreamBackfillSource(
delta_table_name="my_catalog.my_schema.historical_events"
),
)

Chevauchement entre le remplissage historique et les données de Live Stream en direct

Avant d'exécuter un MERGE entre le remplissage rétrospectif et la table d'ingestion, une vérification de chevauchement compare les Timestamp sur les deux tables :

  • Remplissage max : le stream_record_timestamp maximum dans la source de remplissage.
  • Ingestion min : Le minimum stream_record_timestamp de lignes (record_source="stream") dans la table d'ingestion.

Le MERGE se poursuit lorsque le dernier Timestamp du remplissage excède le Timestamp le plus ancien de la table d'ingestion d'au moins 1 heure. Ce chevauchement garantit l'absence de lacunes dans la table d'ingestion. Si la condition de chevauchement n'est pas remplie dans les 2 jours, le MERGE s'exécute quand même pour éviter un blocage indéfini.

Étant donné que le pipeline d'ingestion démarre à partir du dernier offset Kafka, il ne capture que les messages arrivant après la création du stream. Votre source de remplissage rétrospectif doit contenir des données qui s'étendent dans la plage horaire d'ingestion — et non seulement jusqu'à l'heure de création du Stream.

Par exemple, si vous créez un Stream à 3 h 00, le pipeline de remplissage avant commence à lire les messages à partir de 3 h 00 et au-delà. Votre source de remplissage rétroactif doit inclure des Timestamp jusqu'à au moins 16 h 00 (1 heure après le start du remplissage avant) pour satisfaire la vérification de chevauchement. Cela signifie que vous devez mettre à jour votre table de remplissage rétroactif après 16 h 00 pour vous assurer que la table d'ingestion ne présente aucune lacune.

Déduplication

Utilisez deduplication_columns pour spécifier les chemins de colonne afin d'identifier les lignes en double lors de l'ingestion entre les données de Stream de rétro-remplissage et de remplissage anticipé. Utilisez la notation par points pour les champs imbriqués (par exemple, "value.user_id").

Choisissez les colonnes de déduplication en fonction de vos données :

  • Si chaque enregistrement dans votre Stream contient un identifiant unique (par exemple, value.transaction_id), utilisez cette colonne pour la déduplication.
  • Si votre source de remplissage inclut les colonnes kafka_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 default est la combinaison complète de key, value et stream_record_timestamp. Ceci n'est pas recommandé car cette correspondance de critères stricte peut facilement entraîner des doublons.
Python
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
deduplication_columns=["value.transaction_id"],
)

Gérer les Stream

Obtenir un stream

Python
stream = client.get_stream(name="my_catalog.my_schema.my_stream")

Lister les flux

Python
streams = client.list_streams(
catalog_name="my_catalog",
schema_name="my_schema",
max_results=50,
include_schemas=False,
)

Définissez include_schemas=True pour inclure les détails complets du schéma. Les schémas peuvent être volumineux, ce qui peut entraîner une Opération de longue durée. Pour récupérer les schémas individuellement, utilisez plutôt get_stream.

Supprimer un Stream

La suppression d'un stream supprime également son pipeline d'ingestion et sa table d'ingestion.

attention

Les modèles ou fonctionnalités qui référencent le Stream supprimé n'auront plus accès aux données du Stream sous-jacent. Créez une copie de la table d'ingestion avant la suppression si vous avez besoin de ces données mais que vous n'avez plus besoin du Stream.

Python
client.delete_stream(name="my_catalog.my_schema.my_stream")

Exemple de Notebook

Pour un exemple de bout en bout qui crée un Stream, définit des fonctionnalités en streaming et déploie vers un Endpoint de service, consultez le Notebook suivant :

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