Créer un pipeline CDC intégré pour Oracle
Bêta
Cette fonctionnalité est en version Bêta. Les administrateurs de workspace peuvent contrôler l'accès à cette fonctionnalité depuis la page Aperçus . Voir Gérer les aperçus Databricks.
Un pipeline CDC intégré ingère les données de modification d'Oracle vers Databricks à l'aide d'un pipeline unique. Le connecteur CDC intégré combine l'extraction et l'application en une seule mise à jour de pipeline.
Le connecteur CDC intégré à Oracle utilise LogMiner en mode de transaction non validée pour lire les modifications des Logs de rétablissement en ligne et des Logs d'archives.
Exigences
-
Votre workspace est activé pour Unity Catalog.
-
Si vous prévoyez de créer une connexion : vous avez
CREATE CONNECTIONprivilèges sur le metastore. Voir Gérer les privilèges dans Unity Catalog.Si votre connecteur prend en charge la création de pipelines basée sur l'interface utilisateur, vous pouvez créer la connexion et le pipeline simultanément en suivant les étapes de cette page. Cependant, si vous utilisez l’édition de pipelines basée sur l’API, vous devez créer la connexion dans l’Explorateur de catalogues avant de suivre les étapes de cette page. Voir Connexion aux sources d'ingestion gérées.
-
Si vous prévoyez d'utiliser une connexion existante : vous disposez des privilèges
USE CONNECTIONouALL PRIVILEGESsur la connexion. -
Vous disposez de
USE CATALOGprivilèges sur le catalogue cible. -
Vous disposez de
USE SCHEMA,CREATE TABLEetCREATE VOLUMEprivilèges sur un schéma existant ou deCREATE SCHEMAprivilèges sur le catalogue cible. -
Votre Workspace doit avoir la fonctionnalité de connecteur CDC intégré activée. Contactez l'équipe de votre compte Databricks.
-
Vous avez terminé la configuration de la base de données source Oracle. Voir Configurer Oracle pour l'ingestion dans Databricks.
-
Vous disposez des autorisations suivantes :
CREATE CONNECTIONsur le metastore (si vous créez une nouvelle connexion Unity Catalog), ouUSE CONNECTIONsur une connexion existante.USE CATALOGsur le catalogue de destination.USE SCHEMAetCREATE TABLEsur le schéma de destination.CREATE VOLUMEsur le schéma de destination, ou sur le schéma spécifié dansdata_staging_options.
Pour les bases de données Oracle multi-tenant, l’utilisateur de connexion doit être un utilisateur commun dans CDB$ROOT. Pour plus de détails, consultez Créer l’utilisateur de réplication.
Exigences en matière de compute
Un pipeline CDC intégré s'exécute sur un compute classique ou serverless :
- Compute classique : le plan de compute s'exécute dans le Virtual Private Cloud (VPC) ou le VNet de votre Workspace Databricks et doit pouvoir atteindre votre instance Oracle via le réseau. Tout chemin réseau qui permet au plan de compute d'atteindre la base de données est pris en charge, y compris le peering Virtual Private Cloud (VPC) ou VNet, les Endpoint publics et, pour Oracle on-premise, AWS Direct Connect, Azure ExpressRoute ou VPN.
- Compute serverless : configurez la connectivité réseau serverless entre le compute serverless de Databricks et votre base de données source. Les sources on-premise nécessitent un chemin réseau via la sortie Serverless configurée (par exemple, une passerelle de transit ou un VNet appairé avec ExpressRoute ou VPN).
Pour le compute classique, vous pouvez utiliser des autorisations de création de cluster illimitées ou une stratégie de cluster personnalisée avec cluster_type défini sur dlt, runtime_engine défini sur STANDARD, et au moins 8 cœurs recommandés pour une extraction efficace.
Créer une connexion Unity Catalog à Oracle
Créer une connexion Unity Catalog à Oracle avant de créer un pipeline. Consultez Créer une connexion Oracle.
Créer un pipeline CDC intégré
Créez un pipeline CDC intégré en utilisant l'interface utilisateur d'ingestion de données, l'API REST, la CLI Databricks, les notebooks ou les Declarative Automation Bundles.
Toute demande de création de pipeline par programmation doit inclure "channel": "PREVIEW". Lorsque vous utilisez l’interface utilisateur, Databricks définit le canal de distribution pour vous.
Pour les pipelines CDC intégrés à Oracle, source_catalog correspond au nom du service Oracle. Pour les bases de données multi-tenant, il doit s'agir du nom de service CDB$ROOT.
- Databricks UI
- Declarative Automation Bundles
- Databricks notebook
- Databricks CLI
- REST API
-
Dans la barre latérale, cliquez sur Data Ingestion , puis sélectionnez Oracle comme type de source.

-
Sélectionnez une connexion à utiliser. Choisissez une connexion Unity Catalog existante ou créez-en une.


-
Indiquez un nom pour le pipeline et un emplacement pour les Logs des événements. L'emplacement du journal des événements est l'endroit où Databricks stocke les données de staging et les métadonnées utilisées pour effectuer la CDC.

-
Cliquez sur Suivant . Databricks provisionne le compute et crée le pipeline. Cette étape peut prendre un certain temps et afficher
Waiting for resources. Une fois l’opération terminée, sélectionnez les tables sources à ingérer.
-
Sélectionnez le schéma de destination dans lequel le pipeline écrit les données capturées à partir de la source. Le pipeline crée automatiquement des tables portant les mêmes noms que la source dans le schéma sélectionné.

-
Cliquez sur Valider et attendez que la validation réussisse.

-
Définissez un calendrier pour le pipeline. Le pipeline s'exécute tant que des données sont disponibles, s'arrête après avoir atteint un état inactif et reprend à partir du même point lors du prochain trigger.

-
Examinez le pipeline. La vue en liste affiche les flux et les statistiques concernant les données répliquées.

-
Pour vérifier ce que fait le pipeline, et notamment pour examiner les avertissements ou les messages d'erreur lorsqu'une mise à jour échoue, ouvrez le panneau Logs des événements sur la droite.

Le pipeline est désormais configuré et en cours d’exécution. Vous pouvez query les tables créées par le pipeline dans le schéma de destination et les traiter comme des tables Bronze dans l’architecture en médaillon.
Définissez la ressource de pipeline dans un fichier de bundle (par exemple, resources/oracle_integrated_cdc_pipeline.yml) :
variables:
pipeline_name:
description: 'Name for the integrated CDC pipeline'
connection_name:
description: 'Unity Catalog connection name'
dest_catalog:
description: 'Destination catalog for ingested data'
dest_schema:
description: 'Destination schema for ingested data'
resources:
pipelines:
oracle_integrated_cdc_pipeline:
name: ${var.pipeline_name}
channel: PREVIEW
catalog: ${var.dest_catalog}
schema: ${var.dest_schema}
ingestion_definition:
connection_name: ${var.connection_name}
connector_type: CDC
objects:
- table:
source_catalog: 'ORCL'
source_schema: 'HR'
source_table: 'EMPLOYEES'
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
destination_table: 'employees'
table_configuration:
scd_type: 'SCD_TYPE_1'
Pour exécuter le pipeline selon un calendrier, définissez un job qui déclenche le pipeline. Étant donné que chaque étape d'extraction dure au moins 10 minutes, un intervalle de 60 minutes ou plus est un bon point de départ :
resources:
jobs:
oracle_integrated_cdc_job:
name: '${var.pipeline_name}-job'
tasks:
- task_key: 'cdc_ingestion'
pipeline_task:
pipeline_id: ${resources.pipelines.oracle_integrated_cdc_pipeline.id}
schedule:
quartz_cron_expression: '0 0 * * * ?'
timezone_id: 'UTC'
Déployez le bundle avec la CLI Databricks :
databricks bundle deploy
databricks bundle run oracle_integrated_cdc_job
Pour plus d'information, consultez Que sont les bundles d'automatisation déclaratifs ?.
from databricks.sdk import WorkspaceClient
from databricks.sdk.service.pipelines import (
ConnectorType,
IngestionConfig,
IngestionPipelineDefinition,
TableSpec,
)
w = WorkspaceClient()
pipeline = w.pipelines.create(
name="<pipeline-name>",
channel="PREVIEW",
catalog="<destination-catalog>",
schema="<destination-schema>",
ingestion_definition=IngestionPipelineDefinition(
connection_name="<oracle-connection-name>",
connector_type=ConnectorType.CDC,
objects=[
IngestionConfig(
table=TableSpec(
source_catalog="<oracle-service-name>",
source_schema="<oracle-schema>",
source_table="<oracle-table>",
destination_catalog="<destination-catalog>",
destination_schema="<destination-schema>",
)
)
],
),
)
print(f"Pipeline created: {pipeline.pipeline_id}")
databricks pipelines create --json '{
"name": "<pipeline-name>",
"channel": "PREVIEW",
"catalog": "<destination-catalog>",
"schema": "<destination-schema>",
"ingestion_definition": {
"connection_name": "<oracle-connection-name>",
"connector_type": "CDC",
"objects": [
{
"table": {
"source_catalog": "<oracle-service-name>",
"source_schema": "<oracle-schema>",
"source_table": "<oracle-table>"
}
}
]
}
}'
POST /api/2.0/pipelines
{
"name": "my-oracle-integrated-pipeline",
"channel": "PREVIEW",
"catalog": "main",
"schema": "ingestion",
"ingestion_definition": {
"connection_name": "my-oracle-connection",
"connector_type": "CDC",
"objects": [
{
"table": {
"source_catalog": "ORCL",
"source_schema": "HR",
"source_table": "EMPLOYEES",
"table_configuration": {
"scd_type": "SCD_TYPE_1"
}
}
}
],
"data_staging_options": {
"catalog_name": "main",
"schema_name": "ingestion_staging"
}
}
}
Pour répliquer chaque table dans un schéma source, utilisez un objet schema au lieu d'objets table individuels :
POST /api/2.0/pipelines
{
"name": "my-oracle-schema-pipeline",
"channel": "PREVIEW",
"catalog": "main",
"schema": "ingestion",
"ingestion_definition": {
"connection_name": "my-oracle-connection",
"connector_type": "CDC",
"objects": [
{
"schema": {
"source_catalog": "ORCL",
"source_schema": "HR",
"destination_catalog": "main",
"destination_schema": "ingestion"
}
}
]
}
}
Pour start une mise à jour de pipeline :
POST /api/2.0/pipelines/<pipeline-id>/updates
{
"full_refresh": false
}
Planifier les mises à jour récurrentes
Les pipelines CDC intégrés s'exécutent uniquement en mode Trigger. Pour ingérer des données selon un calendrier récurrent, créez une tâche Lakeflow Jobs qui exécute le pipeline. Chaque mise à jour dure environ 30 minutes et peut ne pas finir de traiter l'intégralité du backlog de modifications en une seule mise à jour. Planifiez les pipelines suffisamment souvent pour que les mises à jour ultérieures puissent être rattrapées. Un point de départ de 60 minutes convient à la plupart des charges de travail.
Référence de configuration
Paramètres du pipeline
parameter | Type | Description |
|---|---|---|
| chaîne | Un nom pour le pipeline. |
| chaîne | Doit être |
| Booléen | Facultatif. La valeur par défaut est |
| chaîne | Le catalogue de destination par default. |
| chaîne | Le schéma de destination default. |
| chaîne | La connexion Unity Catalog à Oracle. |
| chaîne | Doit être |
| tableau | La liste des tables ou schémas à ingérer. |
| objet | Facultatif. Le catalogue et le schéma où le pipeline crée le volume de préparation. Default au schéma de destination du pipeline. |
Spécification de table
parameter | Obligatoire | Description |
|---|---|---|
| Oui | Le nom du service Oracle. Pour les bases de données multi-tenant, utilisez le nom du service |
| Oui | Le schéma Oracle (généralement le propriétaire de la table). |
| Oui | Le nom de la table Oracle. |
| Non | Le catalogue de destination. Valeur par défaut : |
| Non | Le schéma de destination. Valeur par défaut : |
| Non | Le nom de la table de destination. La valeur default est |
Configuration de la table
parameter | Par défaut | Description |
|---|---|---|
| Détection automatique | Les colonnes qui identifient chaque ligne. Détection automatique à partir de la clé primaire source si non spécifiée. |
|
|
|
| Détection automatique | Les colonnes utilisées pour ordonner les événements CDC. |
Pour les mappages de types de données Oracle, consultez Mappages de types de données.
Sensibilité à la casse pour les identifiants Oracle
Oracle stocke les identifiants non cités en majuscules. Lorsque vous spécifiez source_catalog, source_schema, source_table et primary_keys dans votre configuration de pipeline, la casse doit correspondre à la manière dont Oracle stocke l'identifiant. Pour la plupart des bases de données, cela signifie l'utilisation des majuscules. Si un identifiant a été créé avec des guillemets doubles qui a conservé une casse différente, utilisez cette casse exacte.
Surveiller le pipeline
Après avoir créé et start un pipeline CDC intégré, surveillez son statut à l'aide des éléments suivants :
-
Interface utilisateur Databricks. Ouvrez le pipeline dans la section Pipelines pour afficher l'état de la mise à jour, les métriques d'ingestion par table et la lignée.
-
API REST.
TextGET /api/2.0/pipelines/<pipeline-id> -
API Événements.
TextGET /api/2.0/pipelines/<pipeline-id>/events
La vue en liste sur la page des détails du pipeline affiche le nombre d’enregistrements traités au fur et à mesure de l’ingestion des données. Ces chiffres refresh automatiquement.

La première mise à jour du pipeline effectue un instantané complet de toutes les tables sélectionnées, ce qui peut prendre plus de temps que les mises à jour incrémentielles. Pour les grandes tables, l'instantané initial peut nécessiter plusieurs mises à jour planifiées pour être terminé.
Vous pouvez query les données ingérées dans Unity Catalog.

Pour le full refresh et le comportement d'auto full refresh, voir refresh des tables cibles.
Les pipelines CDC intégrés ont le dimensionnement automatique vertical activé par default. Si une mise à jour de pipeline échoue en raison d'une condition de mémoire insuffisante, la prochaine mise à jour provisionne automatiquement un Driver plus grand.
Limitations
Limitations générales
- Bêta. Le connecteur CDC intégré et le connecteur Oracle nécessitent une activation au niveau du workspace. Contactez l'équipe de votre compte Databricks.
- Mode Trigger uniquement. Les pipelines CDC intégrés ne prennent pas en charge l'exécution continue (toujours active). Planifiez les pipelines à l'aide d'une tâche Lakeflow Jobs.
- Le canal de distribution doit être
PREVIEW. Les spécifications programmatiques du pipeline doivent inclure"channel": "PREVIEW". - Maximum recommandé d’environ 500 tables par pipeline d'ingestion.
- Les pipelines CDC intégrés ne prennent pas encore en charge les modifications de schéma (opérations DDL).
- L'instantané initial pourrait s'étendre sur plusieurs mises à jour pour les grandes tables.
- Chaque mise à jour s'exécute pendant environ 30 minutes. Le pipeline ne traite pas nécessairement l'intégralité du backlog de modifications en une seule mise à jour. Les mises à jour programmées ultérieures reprennent le traitement là où la mise à jour précédente s'était arrêtée. Vous ne pouvez pas configurer ce runtime.
- Le type de connexion et de connecteur est immuable après la création du pipeline.
Limitations spécifiques à Oracle
- **Déploiements Oracle non pris en charge** : Oracle RAC, Exadata en configuration RAC, Physical Standby, bases de données autonomes Oracle et instances de base de données Amazon RDS multi-tenant.
- Types de données non pris en charge :
XML,JSONet types de données spatiales. - Tables ignorées par LogMiner : LogMiner ignore toute table qui contient
BFILE, des tables imbriquées, des colonnes d'identité, des colonnes de validité temporelle, des colonnesPKREFou des colonnesPKOID. Consultez les limitations de LogMiner. - Longueur de l'identificateur : Les noms de table et de colonne ne peuvent pas dépasser 30 caractères.
- Fonctionnalités post-12.2 : Le connecteur ne prend pas en charge les types de données et les fonctionnalités ajoutées après Oracle Database 12c Release 2, y compris
BOOLEAN,VECTORetJSON.
Dépannage
Si une mise à jour de pipeline échoue :
- Examinez le journal des événements du pipeline dans l'interface utilisateur de Databricks ou via
GET /api/2.0/pipelines/<pipeline-id>/events. - Testez la connexion Unity Catalog depuis Catalog Explorer pour confirmer qu'Oracle est accessible.
- Confirmez que le mode archive log et la journalisation supplémentaire sont activés. Voir Étape 1 : Vérifier le mode archive log et la rétention des logs.
- Vérifiez que l'utilisateur de réplication dispose des privilèges accordés par
DBX_ORACLE_SETUP_UTIL.GRANT_PERMISSIONS. Voir les exigences relatives aux utilisateurs de base de données Oracle. - Pour les bases de données multi-tenant, veuillez confirmer que l'utilisateur est un utilisateur commun dans
CDB$ROOTet quesource_catalogest le nom du serviceCDB$ROOT. - Vérifiez que votre spécification de pipeline inclut
"channel": "PREVIEW".
Si Oracle purge les Logs d’archives avant que le pipeline ne puisse les traiter, effectuez une refresh complète sur les tables affectées.