Esegui una pipeline CDC integrata in modalità continua

Si applica a: icona X rossa Connettori SaaS icona con segno di spunta verde Connettori di database icona X rossa Connettori basati su query

Importante

Questa funzionalità è in versione beta. Gli amministratori dell'area di lavoro possono controllare l'accesso a questa funzionalità dalla pagina Anteprime . Vedere Gestire le anteprime di Azure Databricks.

La modalità continua esegue una pipeline CDC integrata come flusso sempre attivo anziché in base a una pianificazione. Di default, una pipeline CDC integrata funziona in modalità triggered, dove ogni aggiornamento estrae e applica i dati di cambiamento, per poi interrompere. Usa la modalità continua per:

  • Ingestione a bassa latenza. I dati di modifica vengono applicati alle tabelle di streaming di destinazione man mano che arrivano, tipicamente entro pochi minuti, invece di attendere il prossimo aggiornamento programmato.
  • Fonti con limitata ritenzione del journal delle modifiche. Alcuni database memorizzano temporaneamente le modifiche nei log delle transazioni, che possono diventare molto grandi o venire eliminati tra un aggiornamento e l'altro. Funzionare continuamente mantiene la pipeline in linea con la sorgente, riducendo il rischio di rimanere indietro rispetto alla finestra di log disponibile.

Abilita modalità continua

Per eseguire una pipeline CDC integrata in modalità continua, imposta continuous su true nelle impostazioni della pipeline. La pipeline utilizza di default la modalità ottimizzata per la scala. Per i passaggi completi di creazione della pipeline, consulta la pagina della pipeline integrata per il tuo connettore: Crea una pipeline CDC integrata per SQL Server, crea una pipeline CDC integrata per MySQL, oppure Crea una pipeline CDC integrata per Oracle.

Pacchetti di automazione dichiarativa

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

Notebook di Databricks

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

Interfaccia a riga di comando di Databricks

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

REST API

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"
        }
      }
    ]
  }
}

Converti una pipeline attivata in continua

Per passare una pipeline esistente attivata alla modalità continua:

  1. Interrompi l'aggiornamento corrente, se è in corso. Vedi Ferma l'aggiornamento attuale.
  2. Aggiorna la pipeline e imposta continuous a true.

Note

L'operazione di aggiornamento sostituisce l'intera specifica della pipeline, quindi include la definizione completa della pipeline, non solo il campo modificato.

Pacchetti di automazione dichiarativa

Imposta continuous: true nella risorsa pipeline, poi ridistribuisci il bundle:

databricks bundle deploy

Notebook di Databricks

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

Interfaccia a riga di comando di Databricks

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

REST API

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"
        }
      }
    ]
  }
}

Modalità di esecuzione

Le pipeline CDC continue supportano due modalità di esecuzione:

Modalità Descrzione
Ottimizzato per la scala (predefinito) Ingerisce fino a 500 tabelle ruotando internamente l'ingestione tra di esse, senza una scala automatica aggressiva. Usa questa modalità per assorbire un gran numero di tabelle in una singola pipeline.
Ottimizzato per la velocità (Beta) Esegue il flusso di ingestione di tutte le tabelle in modo continuo per minimizzare la latenza di ingestione, tipicamente fino a pochi minuti. La modalità ottimizzata per la velocità supporta fino a 50 tabelle e utilizza una scala automatica aggressiva. Usa la modalità ottimizzata per la velocità quando la latenza è la priorità.

La modalità di esecuzione è impostata dalla configurazione pipelines.managedIngestion.continuous.runMode Spark nella pipeline. La modalità ottimizzata per la scala è la predefinita. Per abilitare la modalità ottimizzata per la velocità, imposta runMode su SPEED.

Abilita modalità ottimizzata per velocità

Per abilitare la modalità ottimizzata per la velocità, imposta la configurazione Spark pipelines.managedIngestion.continuous.runMode su SPEED quando crei la pipeline, oltre a continuous: true:

Pacchetti di automazione dichiarativa

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

Notebook di Databricks

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

Interfaccia a riga di comando di Databricks

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

REST API

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"
        }
      }
    ]
  }
}

Note

Per usare la modalità ottimizzata per la scala, rimuovi la pipelines.managedIngestion.continuous.runMode configurazione.

Seleziona la modalità giusta

Usa il seguente confronto per scegliere tra modalità innescata e le due modalità di esecuzione continua:

Capability Attivato Continuo (ottimizzato per la scalabilità) Continuo (ottimizzato per velocità)
Numero massimo di tabelle supportate 300 500 50
Calcolo richiesto Basso (viene eseguito in base a una pianificazione) Medio (sempre attivo) Alta (sempre attivo con scalabilità automatica aggressiva)
Carico di origine Basso (viene eseguito in base a una pianificazione) High Alta (query continue)
Coerenza Basso (rischio di rollover nel registro delle modifiche) High High

Note

I limiti in questa tabella si applicano alle pipeline CDC continue. I singoli connettori potrebbero avere limiti inferiori. Consulta la documentazione del tuo connettore.

Interrompi l'aggiornamento attuale

Interrompi l'aggiornamento della pipeline in esecuzione prima di convertire una pipeline in modalità continua o di eseguire un refresh selettivo completo. Sostituisci <pipeline-id> con l'ID della pipeline, che puoi trovare nell'interfaccia della pipeline.

Notebook di Databricks

from databricks.sdk import WorkspaceClient

w = WorkspaceClient()
w.pipelines.stop(pipeline_id="<pipeline-id>")

Interfaccia a riga di comando di Databricks

databricks pipelines stop <pipeline-id>

REST API

POST /api/2.0/pipelines/<pipeline-id>/stop

Aggiorna completamente un sottoinsieme di tabelle

In una pipeline continua, puoi aggiornare integralmente un sottoinsieme di tabelle mentre tutte le altre tabelle continuano ad acquisire dati nello stesso aggiornamento. Questo è utile quando una singola tabella necessita di un aggiornamento completo (ad esempio, dopo un cambiamento di schema incompatibile) senza disturbare il resto della pipeline.

Per eseguire un aggiornamento selettivo completo:

  1. Interrompi l'aggiornamento attuale. Vedi Ferma l'aggiornamento attuale.

  2. Avvia un nuovo aggiornamento in cui sono elencate le tabelle da aggiornare completamente in full_refresh_selection e imposta refresh_selection come carattere jolly ["*"].

    Notebook di Databricks

    from databricks.sdk import WorkspaceClient
    
    w = WorkspaceClient()
    w.pipelines.start_update(
        pipeline_id="<pipeline-id>",
        full_refresh_selection=["customers", "orders"],
        refresh_selection=["*"],
    )
    

    Interfaccia a riga di comando di Databricks

    databricks pipelines start-update <pipeline-id> --json '{
      "full_refresh_selection": ["customers", "orders"],
      "refresh_selection": ["*"]
    }'
    

    REST API

    POST /api/2.0/pipelines/<pipeline-id>/updates
    
    {
      "full_refresh_selection": ["customers", "orders"],
      "refresh_selection": ["*"]
    }
    

Le tabelle in full_refresh_selection vengono aggiornate completamente, mentre tutte le altre tabelle continuano ad aggiornarsi nell’ambito dello stesso aggiornamento. Dopo il completamento del completo refresh, la pipeline riprende automaticamente l'ingestione continua normale di tutte le tabelle. Non c'è bisogno di interrompere l'aggiornamento e avviarne uno nuovo senza full_refresh_selection.

Note

In modalità continua, una delle selezioni di aggiornamento deve includere il carattere jolly * in modo che tutte le altre tabelle continuino a essere acquisite mentre quelle selezionate vengono aggiornate completamente. L'aggiornamento solo di un sottoinsieme di tabelle senza * (un aggiornamento parziale) non è supportato in modalità continua.

Limitations

La modalità continua presenta le seguenti limitazioni:

  • Gli aggiornamenti si riavviano per applicare modifiche allo stato. La pipeline utilizza un meccanismo di cancellazione e riavvio per ricaricare il grafo della pipeline o applicare cambiamenti di stato, come modifiche allo schema.
  • Il refresh completo potrebbe richiedere più riavvii. Un aggiornamento completo potrebbe richiedere più riavvii della pipeline per essere completato, perché lo snapshot sorgente è gestito in modo asincrono.

Risorse aggiuntive