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 de streaming externe, telle qu'Apache Kafka ou Amazon Kinesis. Les Streams stockent les détails de connexion, l'authentification, les schémas et la configuration de l'ingestion. Une fois qu'un stream est créé, vous pouvez le référencer à l'aide des définitions 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 Python feature-engineering-client en version 0.17.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 pipeline Lakeflow en streaming à votre broker Kafka ou à votre endpoint de service Kinesis. 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 autoriser votre compute classique à être accessible depuis l'Internet public.

Pour le streaming géré par AWS (Amazon MSK), consultez Serverless private connectivity to Amazon MSK. Un modèle similaire sera nécessaire pour les autres clouds et pour Kinesis.

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 la configuration de la source source_config et de la connexion, ainsi qu’un exemple complet create_stream(), consultez Apache Kafka ou Amazon Kinesis. 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.

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

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

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

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.

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

Amazon Kinesis​

Pour diffuser à partir d’un Amazon Kinesis data Stream, utilisez KinesisStreamConfig comme configuration de source et une connexion Unity Catalog de type KINESIS pour l’authentification.

Python
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
KinesisStreamConfig,
StreamNameList,
StreamConnectionConfig,
DirectSchemas,
SchemaConfig,
IngestionConfig,
IngestionDestination,
)

client = FeatureEngineeringClient()

stream = client.create_stream(
name="my_catalog.my_schema.my_stream",
source_config=KinesisStreamConfig(
stream_names=StreamNameList(names=["my-kinesis-stream"]),
),
connection_config=StreamConnectionConfig(
uc_connection_name="my-kinesis-connection"
),
schema_config=DirectSchemas(
payload_schema=SchemaConfig(
json_schema=(
'{'
' "type": "object",'
' "properties": {'
' "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"
),
),
)

Identifiants de stream Kinesis​

Identifiez le ou les Kinesis Stream(s) à lire à l'aide d'une seule des options suivantes. Un seul Stream peut lire à partir de plusieurs flux de données Kinesis. Transmettez des options de source supplémentaires via extra_options — par exemple, maxFetchRate pour plafonner le taux de lecture par partition, ou consumerMode="efo" pour effectuer une lecture avec fan-out amélioré (EFO) au lieu du consommateur de scrutin default.

Champ

Description

Exemple

stream_names

Liste des noms de Stream Kinesis

StreamNameList(names=["stream-a", "stream-b"])

stream_arns

Liste des ARNs de Stream Kinesis

StreamArnList(arns=["arn:aws:kinesis:us-west-2:123456789012:stream/stream-a"])

Champ

Description

Exemple

stream_names

Liste des noms de Stream Kinesis

StreamNameList(names=["stream-a", "stream-b"])

stream_arns

Liste des ARNs de Stream Kinesis

StreamArnList(arns=["arn:aws:kinesis:us-west-2:123456789012:stream/stream-a"])

Authentification Kinesis​

Kinesis s'authentifie via une connexion Unity Catalog de type KINESIS qui référence une accréditation de service Unity Catalog (un rôle IAM accordant l'accès en lecture à Kinesis) et la région AWS du stream. Pour créer une connexion, consultez la page S'authentifier avec une connexion Unity Catalog. Le créateur du Stream doit disposer de l'autorisation USE CONNECTION sur la connexion, tout comme tout utilisateur qui matérialise des fonctionnalités en utilisant le Stream comme source.

SQL
CREATE CONNECTION IF NOT EXISTS `my-kinesis-connection`
TYPE KINESIS
OPTIONS (
aws_region '<region>',
credential '<service_credential>'
)

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.

Pour les sources Kinesis, fournissez uniquement payload_schema pour les données d'enregistrement. La clé de partition d'un message Kinesis est une chaîne de routage sans schéma, de sorte que key_schema ne s'applique pas ; la clé de partition est toujours capturée dans la colonne key de la table d'ingestion en tant que chaîne simple.

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 :

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

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

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 du message ainsi que des colonnes de métadonnées. Les colonnes communes sont présentes pour chaque source ; les colonnes kafka_* ne le sont que pour un stream Kafka, et les colonnes kinesis_* uniquement pour un stream Kinesis.

Colonne

Type

Source

Description

key

Varie (à partir de key_schema)

Courantes

La clé du message, structurée selon le schéma que vous avez fourni. Un stream Kinesis ne possède pas de schéma de clé. Il s'agit donc de la clé de partition de l'enregistrement sous la forme d'un objet STRING.

value

Varie (à partir de payload_schema)

Courantes

The message value (payload), structured according to the schema you provided.

stream_record_timestamp

TIMESTAMP

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.

record_source

STRING

Courantes

Soit "stream" (remplissage en continu à partir du Live Stream), soit "backfill" (à partir de la source de remplissage rétroactif).

kafka_topic

STRING

Kafka

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

kafka_partition

INT

Kafka

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

kafka_offset

LONG

Kafka

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

kinesis_stream

STRING

Kinesis

Le data stream Kinesis à partir duquel l'enregistrement a été consommé.

kinesis_shard_id

STRING

Kinesis

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

kinesis_sequence_number

STRING

Kinesis

Le numéro de séquence de l'enregistrement dans sa partition.

Colonne

Type

Source

Description

key

Varie (à partir de key_schema)

Courantes

La clé du message, structurée selon le schéma que vous avez fourni. Un stream Kinesis ne possède pas de schéma de clé. Il s'agit donc de la clé de partition de l'enregistrement sous la forme d'un objet STRING.

value

Varie (à partir de payload_schema)

Courantes

The message value (payload), structured according to the schema you provided.

stream_record_timestamp

TIMESTAMP

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.

record_source

STRING

Courantes

Soit "stream" (remplissage en continu à partir du Live Stream), soit "backfill" (à partir de la source de remplissage rétroactif).

kafka_topic

STRING

Kafka

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

kafka_partition

INT

Kafka

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

kafka_offset

LONG

Kafka

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

kinesis_stream

STRING

Kinesis

Le data stream Kinesis à partir duquel l'enregistrement a été consommé.

kinesis_shard_id

STRING

Kinesis

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

kinesis_sequence_number

STRING

Kinesis

Le numéro de séquence de l'enregistrement dans 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.

Pour Kinesis, les colonnes de métadonnées directes sont kinesis_stream, kinesis_shard_id et kinesis_sequence_number.

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 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_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.

Pour un stream Kinesis, utilisez kinesis_shard_id et kinesis_sequence_number conjointement pour la déduplication — un numéro de séquence n'étant unique qu'au sein d'un shard — ainsi que kinesis_stream lorsque le Stream lit plusieurs streams Kinesis.

Python
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
deduplication_columns=["value.transaction_id"],
)

Gérer les streams​

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