Exécuter un pipeline CDC intégré en mode continu
S’applique à : Connecteurs SaaS
Connecteurs de base de données
Connecteurs basés sur la query
Bêta
Cette fonctionnalité est en version bêta. Les administrateurs du Workspace peuvent contrôler l’accès à cette fonctionnalité depuis la page Previews . Consultez Gérer les aperçus Databricks.
Le mode continu exécute un pipeline CDC intégré sous forme de Stream toujours actif au lieu d’une planification. By default, an integrated CDC pipeline runs in triggered mode, where each update extracts and applies change data, then stops. Utilisez le mode continu pour :
- Ingestion à faible latence. Les données de modification sont appliquées aux tables de streaming de destination au fur et à mesure de leur arrivée, généralement en quelques minutes, au lieu d’attendre la prochaine mise à jour planifiée.
- Sources avec une rétention limitée du log de modifications. Certaines bases de données mettent en mémoire tampon les modifications dans des Logs de transactions qui peuvent devenir volumineux ou être purgés entre les mises à jour. L'exécution en continu permet au pipeline de rester synchronisé avec la source, réduisant ainsi le risque de retard par rapport à la fenêtre de log disponible.
Activer le mode continu
Pour exécuter un pipeline CDC intégré en mode continu, définissez continuous sur true dans les paramètres du pipeline. Le pipeline utilise le mode optimisé pour la montée en charge par default. Pour connaître les étapes complètes de création de pipeline, consultez la page du pipeline intégré pour votre connecteur : Créer un pipeline CDC intégré pour SQL Server, Créer un pipeline CDC intégré pour MySQL ou Créer un pipeline CDC intégré pour Oracle.
- Declarative Automation Bundles
- Databricks notebook
- Databricks CLI
- REST API
resources:
pipelines:
continuous_cdc_pipeline:
name: my-continuous-cdc-pipeline
channel: PREVIEW
continuous: true
catalog: main
schema: ingestion
ingestion_definition:
connection_name: my-sqlserver-connection
connector_type: CDC
objects:
- table:
source_catalog: my_database
source_schema: dbo
source_table: customers
from databricks.sdk import WorkspaceClient
from databricks.sdk.service.pipelines import (
ConnectorType,
IngestionConfig,
IngestionPipelineDefinition,
TableSpec,
)
w = WorkspaceClient()
pipeline = w.pipelines.create(
name="my-continuous-cdc-pipeline",
channel="PREVIEW",
continuous=True,
catalog="main",
schema="ingestion",
ingestion_definition=IngestionPipelineDefinition(
connection_name="my-sqlserver-connection",
connector_type=ConnectorType.CDC,
objects=[
IngestionConfig(
table=TableSpec(
source_catalog="my_database",
source_schema="dbo",
source_table="customers",
)
)
],
),
)
print(f"Pipeline created: {pipeline.pipeline_id}")
databricks pipelines create --json '{
"name": "my-continuous-cdc-pipeline",
"channel": "PREVIEW",
"continuous": true,
"catalog": "main",
"schema": "ingestion",
"ingestion_definition": {
"connection_name": "my-sqlserver-connection",
"connector_type": "CDC",
"objects": [
{
"table": {
"source_catalog": "my_database",
"source_schema": "dbo",
"source_table": "customers"
}
}
]
}
}'
POST /api/2.0/pipelines
{
"name": "my-continuous-cdc-pipeline",
"channel": "PREVIEW",
"continuous": true,
"catalog": "main",
"schema": "ingestion",
"ingestion_definition": {
"connection_name": "my-sqlserver-connection",
"connector_type": "CDC",
"objects": [
{
"table": {
"source_catalog": "my_database",
"source_schema": "dbo",
"source_table": "customers"
}
}
]
}
}
Convertir un Trigger pipeline en pipeline continu
Pour faire passer un pipeline Trigger existant en mode continu :
- Arrêtez la mise à jour en cours, si une est en cours d’exécution. Voir Arrêter la mise à jour en cours.
- Mettez à jour le pipeline et définissez
continuoussurtrue.
L’opération de mise à jour remplace l’intégralité de la spécification du pipeline ; incluez donc la définition complète du pipeline, et pas seulement le champ modifié.
- Declarative Automation Bundles
- Databricks notebook
- Databricks CLI
- REST API
Définissez continuous: true dans la ressource de pipeline, puis redéployez le bundle :
databricks bundle deploy
from databricks.sdk import WorkspaceClient
w = WorkspaceClient()
w.pipelines.update(
pipeline_id="<pipeline-id>",
name="my-continuous-cdc-pipeline",
channel="PREVIEW",
continuous=True,
catalog="main",
schema="ingestion",
ingestion_definition=existing_ingestion_definition,
)
databricks pipelines update --json '{
"pipeline_id": "<pipeline-id>",
"name": "my-continuous-cdc-pipeline",
"channel": "PREVIEW",
"continuous": true,
"catalog": "main",
"schema": "ingestion",
"ingestion_definition": {
"connection_name": "my-sqlserver-connection",
"connector_type": "CDC",
"objects": [
{
"table": {
"source_catalog": "my_database",
"source_schema": "dbo",
"source_table": "customers"
}
}
]
}
}'
PUT /api/2.0/pipelines/<pipeline-id>
{
"pipeline_id": "<pipeline-id>",
"name": "my-continuous-cdc-pipeline",
"channel": "PREVIEW",
"continuous": true,
"catalog": "main",
"schema": "ingestion",
"ingestion_definition": {
"connection_name": "my-sqlserver-connection",
"connector_type": "CDC",
"objects": [
{
"table": {
"source_catalog": "my_database",
"source_schema": "dbo",
"source_table": "customers"
}
}
]
}
}
Modes d’exécution
Les pipelines CDC continus prennent en charge deux modes d'exécution :
Mode | Description |
|---|---|
Optimisé pour la montée en charge (default) | Ingère jusqu'à 500 tables en faisant pivoter l'ingestion en interne entre elles, sans montée en charge automatique agressive. Utilisez ce mode pour ingérer un grand nombre de tables dans un seul pipeline. |
Optimisé pour la vitesse (Bêta) | Exécute le Stream d'ingestion pour toutes les tables en continu afin de minimiser la latence d'ingestion, généralement à quelques minutes près. Le mode optimisé pour la vitesse prend en charge jusqu'à 50 tables et utilise un autoscaling agressif. Utilisez le mode optimisé pour la vitesse lorsque la latence est la priorité absolue. |
Le mode d’exécution est défini par la configuration Spark pipelines.managedIngestion.continuous.runMode sur le pipeline. Le mode optimisé pour la montée en charge est le mode default. Pour activer le mode optimisé pour la vitesse, définissez runMode sur SPEED.
Activer le mode optimisé pour la vitesse
Pour activer le mode optimisé pour la vitesse, définissez la configuration Spark pipelines.managedIngestion.continuous.runMode sur SPEED lors de la création du pipeline, en plus de continuous: true:
- Declarative Automation Bundles
- Databricks notebook
- Databricks CLI
- REST API
resources:
pipelines:
continuous_cdc_pipeline:
name: my-continuous-cdc-pipeline
channel: PREVIEW
continuous: true
catalog: main
schema: ingestion
configuration:
pipelines.managedIngestion.continuous.runMode: SPEED
ingestion_definition:
connection_name: my-sqlserver-connection
connector_type: CDC
objects:
- table:
source_catalog: my_database
source_schema: dbo
source_table: customers
from databricks.sdk import WorkspaceClient
from databricks.sdk.service.pipelines import (
ConnectorType,
IngestionConfig,
IngestionPipelineDefinition,
TableSpec,
)
w = WorkspaceClient()
pipeline = w.pipelines.create(
name="my-continuous-cdc-pipeline",
channel="PREVIEW",
continuous=True,
catalog="main",
schema="ingestion",
configuration={"pipelines.managedIngestion.continuous.runMode": "SPEED"},
ingestion_definition=IngestionPipelineDefinition(
connection_name="my-sqlserver-connection",
connector_type=ConnectorType.CDC,
objects=[
IngestionConfig(
table=TableSpec(
source_catalog="my_database",
source_schema="dbo",
source_table="customers",
)
)
],
),
)
print(f"Pipeline created: {pipeline.pipeline_id}")
databricks pipelines create --json '{
"name": "my-continuous-cdc-pipeline",
"channel": "PREVIEW",
"continuous": true,
"catalog": "main",
"schema": "ingestion",
"configuration": {
"pipelines.managedIngestion.continuous.runMode": "SPEED"
},
"ingestion_definition": {
"connection_name": "my-sqlserver-connection",
"connector_type": "CDC",
"objects": [
{
"table": {
"source_catalog": "my_database",
"source_schema": "dbo",
"source_table": "customers"
}
}
]
}
}'
POST /api/2.0/pipelines
{
"name": "my-continuous-cdc-pipeline",
"channel": "PREVIEW",
"continuous": true,
"catalog": "main",
"schema": "ingestion",
"configuration": {
"pipelines.managedIngestion.continuous.runMode": "SPEED"
},
"ingestion_definition": {
"connection_name": "my-sqlserver-connection",
"connector_type": "CDC",
"objects": [
{
"table": {
"source_catalog": "my_database",
"source_schema": "dbo",
"source_table": "customers"
}
}
]
}
}
Pour utiliser le mode optimisé pour Monter en charge, supprimez la configuration pipelines.managedIngestion.continuous.runMode.
Sélectionnez le mode approprié
Utilisez la comparaison suivante pour choisir entre le mode Trigger et les deux modes d'exécution continus :
Compétence | Déclenché | Continu (optimisé pour Monter en charge) | Continu (optimisé pour la vitesse) |
|---|---|---|---|
Nombre maximum de tables prises en charge | 300 | 500 | 50 |
Compute requis | Faible (s’exécute selon un calendrier) | Moyen (toujours activé) | Élevé (toujours activé avec une montée en charge automatique agressive) |
Chargement de la source | Faible (s’exécute selon un calendrier) | Haute | Élevé (requêtes continues) |
Cohérence | Faible (risque de basculement du log de modifications) | Haute | Haute |
Les limites de ce tableau s'appliquent aux pipelines CDC continus. Les connecteurs individuels peuvent avoir des limites inférieures. Consultez la documentation de votre connecteur.
Arrêter la mise à jour en cours
Arrêtez la mise à jour du pipeline en cours avant de convertir un pipeline en mode continu ou d'exécuter un refresh complet sélectif. Remplacez <pipeline-id> par l'ID de votre pipeline, que vous pouvez trouver dans l'interface utilisateur des pipelines.
- Databricks notebook
- Databricks CLI
- REST API
from databricks.sdk import WorkspaceClient
w = WorkspaceClient()
w.pipelines.stop(pipeline_id="<pipeline-id>")
databricks pipelines stop <pipeline-id>
POST /api/2.0/pipelines/<pipeline-id>/stop
Fully refresh un sous-ensemble de tables
Dans un pipeline continu, vous pouvez entièrement refresh un sous-ensemble de tables tandis que toutes les autres tables continuent de s’ingérer lors de la même mise à jour. Ceci est utile lorsqu’une seule table nécessite un full refresh (par exemple, après un changement de schéma incompatible) sans perturber le reste du pipeline.
Pour effectuer une refresh complète sélective :
- Arrêter la mise à jour en cours. Voir Arrêter la mise à jour en cours.
- Start une nouvelle mise à jour qui liste les tables à refresh entièrement dans
full_refresh_selectionet définitrefresh_selectionsur le caractère générique["*"].
- Databricks notebook
- Databricks CLI
- REST API
from databricks.sdk import WorkspaceClient
w = WorkspaceClient()
w.pipelines.start_update(
pipeline_id="<pipeline-id>",
full_refresh_selection=["customers", "orders"],
refresh_selection=["*"],
)
databricks pipelines start-update <pipeline-id> --json '{
"full_refresh_selection": ["customers", "orders"],
"refresh_selection": ["*"]
}'
POST /api/2.0/pipelines/<pipeline-id>/updates
{
"full_refresh_selection": ["customers", "orders"],
"refresh_selection": ["*"]
}
Les tables dans full_refresh_selection sont entièrement actualisées, tandis que toutes les autres tables continuent de s'actualiser au cours de la même mise à jour. Une fois le full refresh terminé, le pipeline reprend automatiquement l'ingestion continue normale pour toutes les tables. Il n'est pas nécessaire d'arrêter la mise à jour et d'en start une nouvelle sans full_refresh_selection.
En mode continu, l’une des sélections de refresh doit inclure le caractère générique * afin que toutes les autres tables continuent l’ingestion pendant que les tables sélectionnées sont entièrement refreshed. Le refreshing d’un sous-ensemble de tables uniquement sans * (un refresh partiel) n’est pas prise en charge en mode continu.
Limitations
Le mode continu présente les limitations suivantes :
- Les mises à jour redémarrent pour appliquer les changements d'état. Le pipeline utilise un mécanisme d'annulation et de redémarrage pour recharger le graphe du pipeline ou appliquer des changements d'état, tels que des changements de schéma.
- Une refresh complète peut nécessiter plusieurs redémarrages. Un refresh complet peut nécessiter plusieurs redémarrages du pipeline pour se terminer, car l'instantané source est mis en zone de transit de manière asynchrone.