Ingérer des données de RabbitMQ
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 RabbitMQ 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 CONNECTIONsur 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 CONNECTIONouALL PRIVILEGESsur l'objet de connexion. -
Vous devez disposer de privilèges
USE CATALOGsur le catalogue cible. -
Vous devez disposer des privilèges
USE SCHEMAetCREATE TABLEsur un schéma existant ou des privilègesCREATE SCHEMAsur le catalogue cible.
-
-
Pour ingérer à partir de RabbitMQ, vous devez d'abord suivre les étapes décrites dans Connexion à RabbitMQ.
Créez un pipeline d'ingestion
Chaque file d'attente classique RabbitMQ 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.
La création de pipeline basée sur l'interface utilisateur n'est pas prise en charge pour le connecteur RabbitMQ en version bêta. Utilisez Declarative Automation Bundles ou un notebook Databricks pour créer votre pipeline.
- Declarative Automation Bundles
- Databricks notebook
Utilisez les Declarative Automation Bundles pour gérer les pipelines RabbitMQ en tant de 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?.
-
Créez un nouveau bundle à l'aide de la CLI Databricks :
Bashdatabricks bundle init -
Ajouter un fichier de définition de pipeline au bundle (par exemple,
resources/rabbitmq_pipeline.yml). Voir pipeline.ingestion_definition et Exemples. -
Déployez le bundle à l'aide de la CLI Databricks :
Bashdatabricks bundle deploy
- Importez le notebook suivant dans votre Databricks Workspace :
-
Laissez la cellule un telle quelle. Ne modifiez pas le champ
channel— il doit resterPREVIEW. -
Modifiez la cellule avec les détails de votre configuration de pipeline, y compris le nom de votre file d'attente. Voir pipeline.ingestion_definition et Exemples.
-
Cliquez sur Tout exécuter .
Exemples
Utilisez ces exemples pour configurer votre pipeline.
Pipeline avec file d'attente unique
Cet exemple ingère une seule file d’attente classique RabbitMQ dans une table de streaming :
- Declarative Automation Bundles
variables:
connection_name:
default: my-rabbitmq-connection
dest_catalog:
default: main
dest_schema:
default: rabbitmq_ingest
resources:
pipelines:
rabbitmq_pipeline:
name: rabbitmq-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:
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
destination_table: orders
connector_options:
rabbitmq_options:
queue: orders
Pipeline avec plusieurs files d'attente
Cet exemple ingère deux files d'attente classiques RabbitMQ, chacune dans sa propre table de destination. Définir une entrée de table par file d'attente :
- Declarative Automation Bundles
variables:
connection_name:
default: my-rabbitmq-connection
dest_catalog:
default: main
dest_schema:
default: rabbitmq_ingest
resources:
pipelines:
rabbitmq_pipeline:
name: rabbitmq-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:
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
destination_table: orders
connector_options:
rabbitmq_options:
queue: orders
- table:
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
destination_table: shipments
connector_options:
rabbitmq_options:
queue: shipments
Modèles courants
Pour les configurations de pipeline avancées, consultez Modèles courants pour les pipelines d'ingestion gérés.
Étapes suivantes
Le connecteur RabbitMQ exécute des pipelines continus. Ensuite, start votre pipeline et configurez-y des alertes. Consultez Tâches courantes de maintenance des pipelines.