Aller au contenu principal

Ingérer des données d'Apache Kafka

info

Bêta

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

Cette page explique comment créer un pipeline d'ingestion Kafka géré à l'aide de Databricks Lakeflow Connect.

Exigences

  • Pour créer un pipeline d'ingestion, vous devez d'abord satisfaire aux exigences suivantes :

    • Votre workspace doit être activé pour Unity Catalog.

    • Le compute Serverless doit être activé pour votre Workspace. Consultez les exigences du compute Serverless.

    • Pour créer une nouvelle connexion, vous devez disposer des privilèges CREATE CONNECTION sur le métastore. Voir Gérer les privilèges dans Unity Catalog.

      Si le connecteur prend en charge la création de pipelines basée sur l'interface utilisateur, un administrateur peut créer la connexion et le pipeline en même temps en suivant les étapes décrites sur cette page. Cependant, si les utilisateurs qui créent des pipelines utilisent la création de pipelines basée sur l'API ou ne sont pas des utilisateurs administrateurs, un administrateur doit d'abord créer la connexion dans l'Explorateur de catalogues. Voir Connexion aux sources d'ingestion gérées.

    • Pour utiliser une connexion existante, vous devez avoir les privilèges USE CONNECTION ou ALL PRIVILEGES sur l'objet de connexion.

    • Vous devez disposer de privilèges USE CATALOG sur le catalogue cible.

    • Vous devez disposer des privilèges USE SCHEMA et CREATE TABLE sur un schéma existant ou des privilèges CREATE SCHEMA sur le catalogue cible.

  • Pour ingérer depuis Kafka, vous devez d'abord suivre les étapes de la section Se connecter à Apache Kafka pour une ingestion gérée.

Créez un pipeline d'ingestion

Chaque rubrique Kafka est ingérée dans une table de streaming. Pour une liste des données prises en charge et des limitations, voir Données prises en charge.

remarque

La création de pipelines basée sur l'interface utilisateur n'est pas prise en charge pour le connecteur Kafka en version bêta. Utilisez Declarative Automation Bundles ou un notebook Databricks pour créer votre pipeline.

Utilisez les Declarative Automation Bundles pour gérer les pipelines Kafka en tant que code. Les bundles peuvent contenir des définitions YAML de jobs et de tâches, sont gérés à l’aide de la Databricks CLI, et peuvent être partagés et exécutés dans différents Workspace cibles (comme le développement, la pré-production et la production). Pour plus d'informations, consultez What are Declarative Automation Bundles?.

  1. Créez un nouveau bundle à l'aide de la CLI Databricks :

    Bash
    databricks bundle init
  2. Ajouter un fichier de définition de pipeline au bundle (par exemple, resources/kafka_pipeline.yml). Voir pipeline.ingestion_definition et Exemples.

  3. Déployez le bundle à l'aide de la CLI Databricks :

    Bash
    databricks bundle deploy

Exemples

Utilisez ces exemples pour configurer votre pipeline.

Pipeline minimal — clé et valeur binaires brutes

Cet exemple ingère un ou plusieurs sujets Kafka avec les colonnes de clé et de valeur conservées comme BINARY:

YAML
variables:
connection_name:
default: my-kafka-connection
dest_catalog:
default: main
dest_schema:
default: kafka_ingest
resources:
pipelines:
kafka_pipeline:
name: kafka-ingestion-pipeline
serverless: true
continuous: true
channel: PREVIEW
catalog: ${var.dest_catalog}
target: ${var.dest_schema}
ingestion_definition:
connection_name: ${var.connection_name}
objects:
- table:
source_table: N/A
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
destination_table: user_events
connector_options:
kafka_options:
topics: [user-events, power-user-events]

Pipeline avec transformateurs — valeur JSON et clé de chaîne

Cet exemple désérialise les clés de message en STRING et les valeurs en JSON, avec l'évolution des schémas activée :

YAML
variables:
connection_name:
default: my-kafka-connection
dest_catalog:
default: main
dest_schema:
default: kafka_ingest
resources:
pipelines:
kafka_pipeline:
name: kafka-ingestion-pipeline
serverless: true
continuous: true
channel: PREVIEW
catalog: ${var.dest_catalog}
target: ${var.dest_schema}
ingestion_definition:
connection_name: ${var.connection_name}
objects:
- table:
source_table: N/A
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
destination_table: user_events
connector_options:
kafka_options:
topics:
- user-events
starting_offset: latest
key_transformer:
format: STRING
value_transformer:
format: JSON
json_options:
schema_evolution_mode: rescue

Acheminement des enregistrements vers plusieurs tables (diffusion)

info

Aperçu

Cette fonctionnalité est en préversion privée. Pour l'essayer, veuillez contacter votre interlocuteur Databricks.

Fanout achemine chaque enregistrement d'une seule source Kafka vers l'une des nombreuses tables de destination. Vous configurez fanout sur un objet de schéma au lieu d'un objet de table. Une clé de routage dérivée de chaque enregistrement détermine le nom de la table de destination : {destination_catalog}.{destination_schema}.{key_value}.

Définissez fanout_options sur l'objet de schéma. Le champ fanout_by est une expression SQL évaluée par rapport à l'enregistrement source brut (avec les colonnes Kafka key et value), et son résultat devient le segment final (le nom de la table) de la table de destination {destination_catalog}.{destination_schema}.{result}. Étant donné que fanout_by s'exécute sur l'enregistrement brut, convertissez le binaire value en chaîne de caractères avant d'en extraire un champ, comme dans l'exemple suivant. Sa valeur résolue est utilisée telle quelle comme segment de nom de table, sans guillemets ni assainissement, elle doit donc être un identifiant de table valide sans guillemets. Les valeurs avec des espaces, des points ou d'autres caractères non valides dans un identifiant non entre guillemets échouent l'écriture, tout comme une valeur entièrement numérique (par exemple, 123). Un chiffre en début est autorisé lorsque la valeur contient également une lettre ou un trait de soulignement, tel que 2024_events.

Vous pouvez appliquer facultativement une seule transformation JSON à chaque route. Cette transformation s'exécute sur chaque enregistrement acheminé après l'acheminement, elle n'affecte donc pas la valeur que fanout_by reçoit. Pour la liste complète des options et les contraintes de la v1, consultez Options d'éclatement et Limites d'éclatement.

L'exemple suivant lit les enregistrements des rubriques correspondant à un modèle et achemine chaque enregistrement vers une table nommée d'après son champ event_type. La transformation JSON facultative analyse la colonne du message value sur place dans chaque table de destination :

YAML
variables:
connection_name:
default: my-kafka-connection
dest_catalog:
default: main
dest_schema:
default: kafka_ingest
resources:
pipelines:
kafka_pipeline:
name: kafka-ingestion-pipeline
serverless: true
continuous: true
channel: PREVIEW
catalog: ${var.dest_catalog}
target: ${var.dest_schema}
ingestion_definition:
connection_name: ${var.connection_name}
objects:
- schema:
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
connector_options:
kafka_options:
topic_pattern: 'events-.*'
starting_offset: earliest
fanout_options:
fanout_by: 'cast(value as string):event_type::string'
transforms:
- format: JSON
input_column: value

Modèles courants

Pour les configurations de pipeline avancées, consultez Modèles courants pour les pipelines d'ingestion gérés.

Étapes suivantes

Start et configurer des alertes sur votre pipeline. Consultez les Tâches de maintenance courantes du pipeline.

Ressources supplémentaires