Führen Sie eine integrierte CDC-Pipeline im kontinuierlichen Modus aus

Gilt für: Rotes X-Symbol SaaS-Connectors Grünes Häkchen Datenbankconnectors Rotes X-Symbol Abfragebasierte Connectors

Important

Dieses Feature befindet sich in der Betaversion. Arbeitsbereichsadministratoren können den Zugriff auf dieses Feature über die Vorschauseite steuern. Siehe Manage Azure Databricks Previews.

Im kontinuierlichen Modus wird eine integrierte CDC-Pipeline als kontinuierlich laufender Stream anstelle einer zeitgesteuerten Ausführung betrieben. Standardmäßig läuft eine integrierte CDC-Pipeline im getriggerten Modus, bei dem jedes Update Änderungsdaten extrahiert und anwendet, dann stoppt. Verwenden Sie den kontinuierlichen Modus für:

  • Datenaufnahme mit geringer Latenz. Änderungsdaten werden auf Ziel-Streaming-Tabellen angewendet, sobald sie eintreffen, typischerweise innerhalb weniger Minuten, anstatt auf das nächste geplante Update zu warten.
  • Quellen mit begrenzter Aufbewahrung des Änderungsprotokolls. Einige Datenbanken puffern Änderungen in Transaktionsprotokollen, die stark anwachsen oder zwischen den Updates gelöscht werden können. Durch den kontinuierlichen Betrieb bleibt die Pipeline auf dem Stand der Quelle, wodurch das Risiko verringert wird, hinter das verfügbare Protokollfenster zurückzufallen.

Aktiviere den kontinuierlichen Modus

Um eine integrierte CDC-Pipeline im kontinuierlichen Modus auszuführen, setzen Sie continuous in den Pipeline-Einstellungen auf true. Die Pipeline verwendet standardmäßig den skalierungsoptimierten Modus. Für vollständige Schritte zur Pipeline-Erstellung siehe die Seite zur integrierten Pipeline für Ihren Connector: Erstellen Sie eine integrierte CDC-Pipeline für SQL Server, Erstellen Sie eine integrierte CDC-Pipeline für MySQL oder Erstellen Sie eine integrierte CDC-Pipeline für Oracle.

Deklarative Automatisierungspakete

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

Databricks-Notizbuch

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-CLI

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

Wandle eine triggergesteuerte Pipeline in eine kontinuierliche Pipeline um

Um eine bestehende ausgelöste Pipeline in den kontinuierlichen Modus umzuschalten:

  1. Stoppe das aktuelle Update, falls eines läuft. Siehe Beenden der aktuellen Aktualisierung.
  2. Aktualisieren Sie die Pipeline und setzen Sie continuous auf true.

Note

Die Aktualisierungsoperation ersetzt die gesamte Pipeline-Spezifikation und fügt also die vollständige Pipeline-Definition hinzu, nicht nur das geänderte Feld.

Deklarative Automatisierungspakete

Setze continuous: true in die Pipeline-Ressource und deploye dann das Bundle erneut:

databricks bundle deploy

Databricks-Notizbuch

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-CLI

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

Ausführungsmodi

Kontinuierliche CDC-Pipelines unterstützen zwei Laufzeitmodi:

Modus Description
Skalierungsoptimiert (Standard) Es erfasst bis zu 500 Tabellen, indem die Datenaufnahme intern zwischen ihnen rotiert wird, ohne aggressives automatisches Skalieren. Verwenden Sie diesen Modus, um eine große Anzahl von Tabellen in einer einzigen Pipeline zu importieren.
Geschwindigkeitsoptimiert (Beta) Betreibt den Datenaufnahmestrom für alle Tabellen kontinuierlich, um die Latenz bei der Datenaufnahme zu minimieren, in der Regel auf wenige Minuten. Der geschwindigkeitsoptimierte Modus unterstützt bis zu 50 Tabellen und verwendet aggressive Autoskalierung. Verwenden Sie den geschwindigkeitsoptimierten Modus, wenn Latenz oberste Priorität hat.

Der Laufmodus wird durch die pipelines.managedIngestion.continuous.runMode Spark-Konfiguration in der Pipeline festgelegt. Der skalenoptimierte Modus ist der Standardmodus. Um den geschwindigkeitsoptimierten Modus zu aktivieren, setzen Sie runMode auf SPEED.

Aktivieren Sie den geschwindigkeitsoptimierten Modus

Um den geschwindigkeitsoptimierten Modus zu aktivieren, setzen Sie die pipelines.managedIngestion.continuous.runMode Spark-Konfiguration beim Erstellen der Pipeline zusätzlich zu continuous: true auf SPEED:

Deklarative Automatisierungspakete

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

Databricks-Notizbuch

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-CLI

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

Um den skalenoptimierten Modus zu verwenden, entferne die Konfiguration pipelines.managedIngestion.continuous.runMode .

Wählen Sie den richtigen Modus

Verwenden Sie folgenden Vergleich, um zwischen dem ausgelösten Modus und den beiden kontinuierlichen Laufmodi auszuwählen:

Fähigkeit Ausgelöst Kontinuierlich (skalierungsoptimiert) Kontinuierlich (geschwindigkeitsoptimiert)
Maximale Anzahl unterstützter Tabellen 300 500 50
Benötigte Berechnung Low (läuft nach einem Zeitplan) Mittel (immer an) Hoch (immer aktiviert bei aggressivem Autoscaling)
Last der Quelle Low (läuft nach einem Zeitplan) High Hoch (kontinuierliche Abfragen)
Konsistenz Niedrig (Risiko eines Change-Log-Rollovers) High High

Note

Die Grenzwerte in dieser Tabelle gelten für kontinuierliche CDC-Pipelines. Einzelne Stecker könnten untere Grenzen haben. Siehe die Dokumentation zu deinem Connector.

Stoppt das aktuelle Update

Stoppe das laufende Pipeline-Update, bevor du eine Pipeline in den kontinuierlichen Modus umwandelst oder eine selektive vollständige Aktualisierung durchführst. Ersetze <pipeline-id> sie durch die ID deiner Pipeline, die du in der Pipeline-Benutzeroberfläche findest.

Databricks-Notizbuch

from databricks.sdk import WorkspaceClient

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

Databricks-CLI

databricks pipelines stop <pipeline-id>

REST API

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

Erneuere eine Teilmenge der Tabellen vollständig

In einer kontinuierlichen Pipeline können Sie eine Teilmenge der Tabellen vollständig neu laden, während alle anderen Tabellen im Rahmen derselben Aktualisierung weiterhin Daten aufnehmen. Dies ist nützlich, wenn eine einzelne Tabelle eine vollständige Aktualisierung benötigt (zum Beispiel nach einer inkompatiblen Schemaänderung), ohne den Rest der Pipeline zu stören.

Um eine selektive vollständige Aktualisierung durchzuführen:

  1. Stoppt das aktuelle Update. Siehe Beenden der aktuellen Aktualisierung.

  2. Starten Sie eine neue Aktualisierung, die die Tabellen auflistet, die in full_refresh_selection vollständig aktualisiert werden sollen, und setzen Sie refresh_selection auf den Platzhalter ["*"].

    Databricks-Notizbuch

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

    Databricks-CLI

    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": ["*"]
    }
    

Die Tabellen in full_refresh_selection werden vollständig aktualisiert, während alle anderen Tabellen innerhalb derselben Aktualisierung weiterhin aktualisiert werden. Sobald die vollständige Aktualisierung abgeschlossen ist, setzt die Pipeline die normale kontinuierliche Datenaufnahme für alle Tabellen automatisch fort. Es ist nicht nötig, das Update zu stoppen und ein neues zu starten, ohne full_refresh_selection.

Note

Im kontinuierlichen Modus muss eine der Aktualisierungsauswahlen das Platzhalterzeichen * enthalten, damit alle anderen Tabellen weiterhin Daten erfassen, während die ausgewählten Tabellen vollständig aktualisiert werden. Das Aktualisieren nur einer Teilmenge von Tabellen ohne * (eine teilweise Aktualisierung) wird im kontinuierlichen Modus nicht unterstützt.

Limitations

Der kontinuierliche Modus hat folgende Einschränkungen:

  • Updates werden neu gestartet, um Zustandsänderungen anzuwenden. Die Pipeline verwendet einen Mechanismus zum Abbrechen und Neustarten, um den Graphen der Pipeline neu zu laden oder Zustandsänderungen wie Schemaänderungen zu übernehmen.
  • Eine vollständige Aktualisierung erfordert möglicherweise mehrere Neustarts. Eine vollständige Aktualisierung erfordert möglicherweise mehrere Neustarts der Pipeline, da der Quell-Snapshot asynchron bereitgestellt wird.

Weitere Ressourcen