Configurer un Stream
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.
Streams have three-part names (catalog.schema.stream_name). Access to a Stream is governed by its associated ingestion table. See Ingestion and backfill for details.
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
feature-engineering-clientPython version 0.18.0 ou supérieure doit être installé.
Connexion à des sources de Stream
Avant de définir des fonctionnalités de streaming, connectez et testez une connexion de Lakeflow pipeline en streaming à votre broker Kafka. Magasin de fonctionnalités s'appuie sur le SDP Serverless, ce qui signifie que vous aurez besoin d'un mécanisme pour connecter votre compute classique (broker ou Endpoint) au compute Serverless de Databricks. Pour ce faire, vous pouvez utiliser des produits tels que privatelink ou permettre à votre compute classique d'être accessible depuis l'Internet public.
Créer un Stream
Utilisez create_stream() pour créer un nouveau Stream. Un Stream nécessite quatre composants de configuration :
- Source config : Specifies the streaming platform and source-specific details, such as the topic subscription for a Kafka source.
- 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 d’ingestion : spécifie où et comment les données du Stream sont ingérées. Consultez Ingestion et remplissage pour plus de détails.
Pour en savoir plus sur les éléments source_config et la configuration de la connexion propres à la source, ainsi que pour consulter un exemple complet create_stream(), voir Apache Kafka. Les options de schéma et d'ingestion sont partagées entre les sources.
Apache Kafka
Pour effectuer un Stream à partir d'Apache Kafka, utilisez KafkaStreamConfig comme configuration de source et une connexion Unity Catalog pour l'authentification. Consultez Streaming on serverless compute et Connect to Apache Kafka pour en savoir plus sur 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 comment le Stream sélectionne les rubriques Kafka à consommer. Trois modes sont pris en charge :
Mode | Description | Exemple |
|---|---|---|
| Liste de noms de rubriques séparés par des virgules |
|
| Noms de sujets correspondant aux modèles d'expressions régulières Java. |
|
| JSON spécifiant les attributions de partitions de rubrique |
|
Kafka authentication
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.
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 :
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 :
PLAINSCRAM-SHA-256SCRAM-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.
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.
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é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.
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.
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.
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,
uint32versBIGINTetuint64versDECIMAL(20,0)), les champs enum sont décodés vers leur nom de chaîne, et les types wrapper scalaires (par exemple,StringValueetInt32Value) 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 :
- Pris en charge pour les flux Kafka uniquement.
- 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 :
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.
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>-valuepour la valeur et<topic>-keypour la clé. Par exemple, le schéma de valeur 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 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 quetransactions-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é.
Pour les Streams pris en charge par un registre de schémas, le pipeline d’ingestion redémarre automatiquement toutes les quelques heures. À chaque redémarrage, il récupère la dernière version de schéma du sujet, et les champs nouveaux ou modifiés apparaissent dans la table d'ingestion.
Les Stream qui utilisent des schémas directs au lieu d’un registre de schémas font évoluer leur schéma avec update_stream. Consultez Mettre à jour un stream.
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.
Filtrer les enregistrements par type
Un Stream décode chaque enregistrement à l'aide d'un schéma de clé et de valeur unique (le cas échéant), que vous les spécifiiez directement ou que vous utilisiez un registre de schémas. Étant donné qu'un sujet peut transporter plusieurs types d'enregistrements et que les Streams peuvent s'abonner à plusieurs sujets, utilisez record_type_filter pour sélectionner les enregistrements du sujet qui appartiennent à ce Stream.
Fournissez une expression SQL qui fait référence aux champs décodés avec la notation par points, par exemple value.event_type = 'transaction'. Les enregistrements qui ne correspondent pas au filtre sont ignorés. Ils ne sont pas écrits dans la table d'ingestion et ne sont pas utilisés lors de la matérialisation. Pour créer un Stream pour d'autres types d'enregistrements, créez un Stream distinct avec un autre record_type_filter.
stream = client.create_stream(
name="my_catalog.my_schema.my_stream",
# ...source, connection, schema, and ingestion config...
record_type_filter="value.event_type = 'transaction'",
)
Même sans record_type_filter, le décodage n'entraîne jamais l'échec du Stream. Un enregistrement qui ne correspond pas au schéma configuré est décodé de manière permissive. Les enregistrements sont décodés de l'une des manières suivantes :
- Dans une ligne avec
NULLvaleurs pour les champs attendus par le schéma mais omis par l'enregistrement (JSON, Avro et Protobuf). - Dans une ligne contenant des valeurs qui appartiennent à un type d’enregistrement différent (Avro et Protobuf uniquement).
Pour identifier si une ligne de la table d'ingestion appartient au type d'enregistrement attendu, veuillez utiliser l'une des vérifications suivantes :
- Vérifiez qu’un champ correspond à une valeur attendue, par exemple
value.event_type = 'transaction'(recommandé pour Avro et Protobuf). - Vérifiez qu'un champ n'est pas de type
NULL, par exemplevalue.activity_id IS NOT NULL.
Il est recommandé d'utiliser record_type_filter avec des Stream distincts lorsque les schémas diffèrent sensiblement entre les types d'enregistrements sur la rubrique ou si vous souhaitez régir l'accès à chaque type d'enregistrement de manière indépendante. Pour maîtriser les coûts, Databricks vous recommande de conserver un petit nombre de Stream, car chaque Stream possède un pipeline d'ingestion et une table d'ingestion distincts. Chaque Stream utilise également un compute distinct au moment de la matérialisation. Vous pouvez utiliser des filtres spécifiques aux caractéristiques pour la matérialisation.
record_type_filter diffère du filter_condition d'une fonctionnalité. record_type_filter est défini sur le Stream et contrôle quels enregistrements sont ingérés et disponibles pour toutes les fonctionnalités utilisant le Stream comme source, tandis que filter_condition est défini sur une fonctionnalité individuelle et filtre les lignes avant l'agrégation. Consultez Filter conditions on streaming sources pour plus de détails sur filter_condition.
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 :
SELECTsur la table d'ingestion accorde un accès en lecture au Stream.MANAGEsur 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 Stream source et les écrit dans une table Delta (la table d'ingestion). Le pipeline start à partir de la dernière position de la source et s'exécute en continu, ne capturant que 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 le sont également.
Destination d'ingestion
Le ingestion_destination spécifie le nom de table Delta en trois parties où les données de stream sont écrites.
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 de message ainsi que les colonnes de métadonnées. Les colonnes communes sont présentes pour chaque source ; les colonnes kafka_* ne le sont que pour un Kafka Stream.
Colonne | Type | Source | Description |
|---|---|---|---|
| Varie (à partir de | Courantes | La clé du message, structurée selon le schéma que vous avez fourni. |
| Varie (à partir de | Courantes | The message value (payload), structured according to the schema you provided. |
|
| Courantes | Le Timestamp de l'enregistrement. Pour les données de remplissage en continu, il s'agit du Timestamp d'ingestion de la source. Pour les données de remplissage rétroactif, cette valeur est fournie par le client. |
|
| Courantes | Soit |
|
| Kafka | Le sujet Kafka d'où l'enregistrement a été consommé. |
|
| Kafka | La partition Kafka à partir de laquelle l'enregistrement a été consommé. |
|
| Kafka | Le décalage Kafka de l'enregistrement au sein de sa partition. |
Source de remplissage
Puisque le pipeline de remplissage en avant start à partir de la dernière position dans la source, il ne capture pas les messages qui existaient avant la création du Stream. Pour fournir une couverture de données historiques pour l'entraînement, configurez une source de rattrapage 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 doit inclure une colonne stream_record_timestamp de type TIMESTAMP dans le fuseau horaire UTC. Les autres colonnes de métadonnées sont transmises si elles sont présentes sur la source de remplissage, ou définies sur NULL dans le cas contraire. 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 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_timestampmaximum dans la source de remplissage. - Ingestion min : Le minimum
stream_record_timestampde 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 start à partir de la dernière position de la source, il capture uniquement les messages qui arrivent après la création du Stream. Votre source de remplissage doit contenir des données qui s'étendent dans la plage temporelle d'ingestion — et pas 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_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 default est la combinaison complète de
key,valueetstream_record_timestamp. Ceci n'est pas recommandé car cette correspondance de critères stricte peut 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"],
)
Attribution des coûts
Définissez tags et budget_policy_id sur le IngestionConfig pour imputer le coût de l'ingestion gérée du Stream. Databricks les applique au Lakeflow pipeline d'ingestion ainsi qu'à ses Jobs de remplissage automatique et rétroactif lors de la création du Stream.
Pour obtenir un exemple, connaître les limites de tags et savoir comment query les dépenses attribuées, consultez Attribute costs with tags and serverless usage policies.
Exclure les colonnes d’un Stream
Use excluded_columns to drop specific columns from a Stream that you don't want to ingest. Une colonne exclue n'est pas écrite dans la table d'ingestion et ne peut pas être référencée par une fonctionnalité ni être utilisée lors de l'entraînement.
Spécifiez chaque colonne à l'aide d'une notation par points dans la clé ou la valeur du message, telle que value.user.email ou key.account_id. Ces colonnes sont supprimées du key et du value décodeurs lors de l'ingestion, du remplissage et de la matérialisation. Si un chemin pointe vers un struct, tous ses champs imbriqués sont également supprimés (par exemple, value.address supprime également value.address.city et value.address.zip).
stream = client.create_stream(
name="my_catalog.my_schema.my_stream",
# ...source, connection, schema, and ingestion config...
excluded_columns=["value.user.email", "value.user.ssn"],
)
Lors de l’utilisation de schémas directs, la colonne exclue doit déjà exister dans le schéma de clé ou de valeur, ou create_stream échoue. Lors de l’utilisation d’un registre de schémas, vous pouvez exclure une colonne avant qu’elle n’existe. Une colonne exclue ne peut pas non plus être une colonne de déduplication, car les colonnes de déduplication sont requises pour identifier les lignes en double. La création de toute fonctionnalité qui fait référence à une colonne exclue (par exemple, en tant qu’entité, série chronologique ou entrée) échoue.
Vous pouvez modifier les colonnes exclues d’un Stream après sa création avec update_stream, qu'il s'agisse d'un Stream à schéma direct ou basé sur un registre de schémas. Voir Update a stream pour plus de détails.
Gérer les streams
Obtenir un stream
stream = client.get_stream(name="my_catalog.my_schema.my_stream")
Lister les flux
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.
Mettre à jour un stream
Utilisez update_stream pour modifier un Stream après sa création. Transmettez schema_config pour faire évoluer un schéma direct, excluded_columns pour modifier les colonnes supprimées, ou les deux. La mise à jour d'autres champs n'est pas prise en charge. Créez plutôt un nouveau Stream.
La mise à jour d’un Stream redémarre son pipeline d’ingestion pour que la modification prenne effet. L’ingestion reprend généralement en quelques minutes.
Faire évoluer un schéma direct
Pour un Stream qui utilise des schémas directs, transmettez un DirectSchemas à schema_config. Définissez payload_schema, key_schema ou les deux. Un côté que vous ne définissez pas reste inchangé. Les Stream s'appuyant sur un registre de schémas rejettent une mise à jour schema_config et doivent obligatoirement être mis à jour par l'intermédiaire du registre.
from databricks.feature_engineering.entities import DirectSchemas, SchemaConfig
stream = client.update_stream(
name="my_catalog.my_schema.my_stream",
schema_config=DirectSchemas(
payload_schema=SchemaConfig(
json_schema=(
'{'
' "type": "object",'
' "properties": {'
' "user_id": {"type": "string"},'
' "amount": {"type": "number"},'
' "event_time": {"type": "string"},'
' "channel": {"type": "string"}'
' }'
'}'
)
),
),
)
Les mises à jour de schémas doivent être rétrocompatibles afin que le pipeline d’ingestion en cours d’exécution puisse continuer à décoder les enregistrements existants et à écrire dans la table d’ingestion. Toutes les autres modifications sont rejetées.
Ce qui est autorisé dépend du format :
- JSON et Protobuf : ajoutez des champs facultatifs, supprimez des champs et élargissez le type d’un champ (par exemple, de
intàbigint). Protobuf permet également de réordonner les champs. - Avro : permet uniquement d'élargir
intenlonget de supprimer un champ de fin dont aucun champ ultérieur ne lit les octets. Pour faire évoluer un schéma Avro plus librement, utilisez plutôt un stream pris en charge par un registre de schémas.
L'ajout de champs développe les structures key et value décodées de la table d'ingestion. Les lignes écrites avant la mise à jour conservent leur format d'origine, et les champs ajoutés s'affichent en tant que NULL pour ces lignes antérieures. Les suppressions et les modifications de type prennent effet uniquement pour les enregistrements ingérés après la mise à jour.
Modifier les colonnes exclues
Transmettez le nouvel ensemble complet de chemins de colonnes à excluded_columns, qui remplace l’ensemble existant. Transmettez une liste vide ([]) pour effacer toutes les exclusions. Pour plus de détails sur ce comportement, consultez Exclure des colonnes d’un Stream.
stream = client.update_stream(
name="my_catalog.my_schema.my_stream",
excluded_columns=["value.user.email", "value.user.ssn"],
)
La modification des colonnes exclues s'applique uniquement aux données futures. Les colonnes nouvellement exclues cessent d'être écrites (apparaissant à l'emplacement NULL) et les colonnes nouvellement incluses start à être alimentées pour la suite, tandis que les lignes précédemment écrites restent inchangées. Pour empêcher qu'une nouvelle colonne ne soit jamais ingérée :
- Registre de schémas : ajoutez d'abord la colonne à
excluded_columnset attendez que le pipeline d’ingestion redémarre, puis enregistrez la nouvelle version de schéma dans le registre. - Schémas directs : ajoutez la colonne à
schema_configet àexcluded_columnsdans le même appelupdate_stream.
Supprimer un Stream
La suppression d'un stream supprime également son pipeline d'ingestion et sa table d'ingestion.
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.
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 :