Nota
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare ad accedere o modificare le directory.
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare a modificare le directory.
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.
Una pipeline CDC integrata inserisce i dati delle modifiche da SQL Server a Azure Databricks usando una singola pipeline. A differenza dell'architettura standard basata su gateway, che richiede un gateway di inserimento e una pipeline di inserimento separati, una pipeline CDC integrata esegue le fasi di estrazione e applicazione in un unico aggiornamento della pipeline.
Quando usare il connettore CDC integrato
La tabella seguente confronta le pipeline CDC integrate con l'architettura standard basata su gateway:
| Feature | CDC standard (basato su gateway) | CDC integrato |
|---|---|---|
| Numero di pipeline | Due (gateway di inserimento e pipeline di inserimento) | Uno (pipeline unificata) |
| Configurazione | Creare un gateway, quindi creare una pipeline di inserimento che faccia riferimento all'ID gateway | Creare una singola pipeline che faccia riferimento a una connessione al catalogo Unity |
| La modalità gateway | Il gateway è in esecuzione continua | La pipeline integra l'estrazione in ogni aggiornamento |
| Informazioni di riferimento sulla connessione | ingestion_gateway_id |
connection_name (una connessione al catalogo Unity) |
| Tipo di connettore | Implicito | Esplicito: connector_type: CDC |
| Volume di gestione temporanea | Il gateway gestisce internamente il volume di staging | Si configura il volume di staging tramite data_staging_options. Se non specificato, la pipeline ne crea automaticamente uno. |
Per la configurazione del database di origine, vedere Configurare Microsoft SQL Server per l'inserimento in Azure Databricks. La stessa configurazione di origine si applica a entrambe le architetture.
Modalità di esecuzione di una pipeline CDC integrata
Ogni aggiornamento della pipeline esegue due fasi in sequenza:
- Estrazione. La pipeline si connette al database di origine usando la connessione al catalogo Unity. Durante la prima esecuzione o un aggiornamento completo, acquisisce uno snapshot iniziale. Nelle esecuzioni successive acquisisce le modifiche incrementali (inserimenti, aggiornamenti ed eliminazioni) usando il meccanismo predefinito di rilevamento delle modifiche del database. La pipeline scrive i dati estratti in un volume di staging di Unity Catalog.
- Applicazione. La pipeline legge dal volume di staging e applica le modifiche alle tabelle di streaming di destinazione in Unity Catalog. Le operazioni di merge utilizzano le chiavi primarie configurate e il tipo SCD. La pipeline garantisce la semantica di tipo exactly-once.
Ogni aggiornamento della pipeline estrae le modifiche e poi si arresta automaticamente una volta allineato alla sorgente, entro un tempo massimo di esecuzione. Per informazioni dettagliate, vedere Chiusura intelligente per pipeline CDC integrate. Per acquisire dati con cadenza regolare, pianifica la pipeline usando un'attività di Lakeflow Jobs.
Requirements
L'area di lavoro è configurata per Unity Catalog.
Se prevedi di creare una connessione: hai privilegi
CREATE CONNECTIONnel metastore. Consulta Gestione dei privilegi in Unity Catalog.Se il connettore supporta la creazione di pipeline basate sull'interfaccia utente, è possibile creare la connessione e la pipeline contemporaneamente completando i passaggi in questa pagina. Tuttavia, se si usa la creazione di pipeline basate su API, è necessario creare la connessione in Esplora cataloghi prima di completare i passaggi in questa pagina. Vedere Connettersi alle origini di inserimento gestite.
Se si prevede di usare una connessione esistente: si dispone di
USE CONNECTIONprivilegi oALL PRIVILEGESsulla connessione.Hai i privilegi
USE CATALOGsul catalogo di destinazione.Si dispone di
USE SCHEMAprivilegi,CREATE TABLE, eCREATE VOLUMEper uno schema esistente oCREATE SCHEMAprivilegi sul catalogo di destinazione.
- L'area di lavoro deve avere la funzionalità integrata del connettore CDC abilitata. Contatta il team responsabile del tuo account Azure Databricks.
- È possibile accedere all'istanza di primary SQL Server. Il connettore CDC integrato non supporta repliche in lettura, istanze di standby o istanze secondarie.
- La configurazione dell'origine SQL Server è stata completata. Vedere Configurare Microsoft SQL Server per l'inserimento in Azure Databricks.
- Sono disponibili le autorizzazioni seguenti:
-
CREATE CONNECTIONnel metastore (se si crea una nuova connessione al catalogo Unity) oUSE CONNECTIONin una connessione esistente. -
USE CATALOGnel catalogo di destinazione. -
USE SCHEMAeCREATE TABLEnello schema di destinazione. -
CREATE VOLUMEnello schema di destinazione o nello schema specificato indata_staging_options. È necessario un volume di staging anche sedata_staging_optionsnon è impostato, perché la pipeline ne crea automaticamente una nello schema di destinazione.
-
Requisiti di calcolo
Una pipeline CDC integrata viene eseguita in un ambiente di calcolo classico o serverless:
- Calcolo classico: l'ambiente di calcolo classico viene eseguito nella VPC o nella VNet dell'area di lavoro di Azure Databricks e deve poter raggiungere l'istanza di SQL Server tramite la rete. Qualsiasi percorso di rete che consente al piano di calcolo di raggiungere il database è supportato, tra cui peering VPC o rete virtuale, endpoint pubblici e, per SQL Server locali, AWS Direct Connect, Azure ExpressRoute o VPN.
- Elaborazione serverless: Configurare la connettività di rete serverless tra l'elaborazione serverless di Azure Databricks e il database di origine. Le origini locali richiedono un percorso di rete tramite l'uscita serverless configurata, ad esempio un gateway di transito o una rete virtuale con peering con ExpressRoute o VPN.
Per il calcolo classico, è possibile usare autorizzazioni di creazione di cluster senza restrizioni o criteri di cluster personalizzati con cluster_type fisso su dlt, runtime_engine fisso su STANDARDe almeno 8 core consigliati per un'estrazione efficiente.
Creare una connessione del catalogo Unity a SQL Server
Creare una connessione al catalogo Unity per SQL Server prima di creare una pipeline. Vedere Creare una connessione SQL Server.
Creare una pipeline CDC integrata
Crea pipeline CDC integrate usando l'API, la CLI di Databricks, i notebook o i bundle di automazione dichiarativa. La creazione dell'interfaccia utente non è ancora disponibile.
Importante
Tutte le richieste di creazione della pipeline devono includere "channel": "PREVIEW".
Pacchetti di automazione dichiarativa
Definire la risorsa della pipeline in un file di bundle , ad esempio resources/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:
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: 'my_database'
source_schema: 'dbo'
source_table: 'customers'
destination_catalog: ${var.dest_catalog}
destination_schema: ${var.dest_schema}
destination_table: 'customers'
table_configuration:
scd_type: 'SCD_TYPE_1'
Per eseguire la pipeline in base a una pianificazione, definisci un job (ad esempio, resources/integrated_cdc_job.yml) che attivi la pipeline. Poiché ogni fase di estrazione viene eseguita per almeno 10 minuti, un intervallo di 60 minuti o più è un buon punto di partenza:
resources:
jobs:
integrated_cdc_job:
name: '${var.pipeline_name}-job'
tasks:
- task_key: 'cdc_ingestion'
pipeline_task:
pipeline_id: ${resources.pipelines.integrated_cdc_pipeline.id}
schedule:
quartz_cron_expression: '0 0 * * * ?'
timezone_id: 'UTC'
Distribuisci il bundle con Databricks CLI:
databricks bundle deploy
databricks bundle run integrated_cdc_job
Per altre informazioni, vedere Che cosa sono i bundle di automazione dichiarativa?.
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="<pipeline-name>",
channel="PREVIEW",
catalog="<destination-catalog>",
schema="<destination-schema>",
ingestion_definition=IngestionPipelineDefinition(
connection_name="<unity-catalog-connection-name>",
connector_type=ConnectorType.CDC,
objects=[
IngestionConfig(
table=TableSpec(
source_catalog="<source-database>",
source_schema="<source-schema>",
source_table="<source-table>",
destination_catalog="<destination-catalog>",
destination_schema="<destination-schema>",
)
)
],
),
)
print(f"Pipeline created: {pipeline.pipeline_id}")
Interfaccia a riga di comando di Databricks
databricks pipelines create --json '{
"name": "<pipeline-name>",
"channel": "PREVIEW",
"catalog": "<destination-catalog>",
"schema": "<destination-schema>",
"ingestion_definition": {
"connection_name": "<unity-catalog-connection-name>",
"connector_type": "CDC",
"objects": [
{
"table": {
"source_catalog": "<source-database>",
"source_schema": "<source-schema>",
"source_table": "<source-table>"
}
}
]
}
}'
REST API
Nell'esempio seguente vengono replicate due tabelle da un database SQL Server. La tabella customers usa scD Type 1 e la tabella orders usa scD Type 2 (che richiede SQL Server CDC nell'origine). Entrambi ereditano la destinazione main.ingestiondi primo livello. Nell'esempio viene omesso serverless, che per impostazione predefinita è false (calcolo classico). Aggiungere "serverless": true per l'esecuzione in un ambiente di calcolo serverless.
POST /api/2.0/pipelines
{
"name": "my-integrated-cdc-pipeline",
"channel": "PREVIEW",
"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",
"table_configuration": {
"scd_type": "SCD_TYPE_1"
}
}
},
{
"table": {
"source_catalog": "my_database",
"source_schema": "dbo",
"source_table": "orders",
"table_configuration": {
"scd_type": "SCD_TYPE_2"
}
}
}
],
"data_staging_options": {
"catalog_name": "main",
"schema_name": "ingestion_staging"
}
}
}
Per replicare ogni tabella in uno schema di origine, usare un schema oggetto anziché singoli table oggetti. La pipeline ignora le tabelle senza CDC o il rilevamento delle modifiche abilitato nell'origine.
POST /api/2.0/pipelines
{
"name": "my-integrated-cdc-schema-pipeline",
"channel": "PREVIEW",
"catalog": "main",
"schema": "ingestion",
"ingestion_definition": {
"connection_name": "my-sqlserver-connection",
"connector_type": "CDC",
"objects": [
{
"schema": {
"source_catalog": "my_database",
"source_schema": "dbo",
"destination_catalog": "main",
"destination_schema": "ingestion"
}
}
]
}
}
Per avviare un aggiornamento della pipeline:
POST /api/2.0/pipelines/<pipeline-id>/updates
{
"full_refresh": false
}
Pianificare gli aggiornamenti ricorrenti
Di default, le pipeline CDC integrate funzionano in modalità triggered. Per l'esecuzione continua, vedi Eseguire una pipeline CDC integrata in modalità continua. Per acquisire dati secondo una pianificazione ricorrente, crea un'attività di Lakeflow Jobs che esegue la pipeline. La durata dell'aggiornamento varia in base alla quantità di dati di modifica presenti nell'origine e un backlog elevato potrebbe non essere smaltito in un singolo aggiornamento (vedi Chiusura intelligente per pipeline CDC integrate). Pianifica le pipeline abbastanza frequentemente in modo che gli aggiornamenti successivi riescano a recuperare. Un punto di partenza di 60 minuti funziona bene per la maggior parte dei carichi di lavoro. Se un attivatore si attiva mentre è ancora in esecuzione un aggiornamento precedente, il nuovo aggiornamento viene messo in coda.
Informazioni di riferimento sulla configurazione
Parametri della pipeline
| Parametro | Tipo | Description |
|---|---|---|
name |
string | Un nome per la pipeline. |
channel |
string | Deve essere PREVIEW. |
serverless |
Boolean | Optional. Di default è false. Impostare su true per il calcolo serverless o false per il calcolo classico. Il calcolo serverless richiede una connettività serverless verso il tuo database di origine. |
catalog |
string | Catalogo di destinazione predefinito. Utilizzato quando non viene specificato un destination_catalog per tabella. |
schema |
string | Schema di destinazione predefinito. Utilizzato quando non viene specificato un destination_schema per tabella. |
ingestion_definition.connection_name |
string | Connessione del catalogo Unity al database di origine. |
ingestion_definition.connector_type |
string | Deve essere CDC. |
ingestion_definition.objects |
array | Elenco di tabelle o schemi da inserire. |
ingestion_definition.data_staging_options |
object | Optional. Il catalogo e lo schema in cui la pipeline crea il volume di staging. Il valore predefinito è lo schema di destinazione della pipeline. |
Specifica della tabella
| Parametro | Obbligatorio | Description |
|---|---|---|
source_catalog |
Yes | Nome del database di origine. |
source_schema |
Yes | Nome dello schema di origine. |
source_table |
Yes | Nome della tabella di origine. |
destination_catalog |
No | Catalogo di destinazione. Per impostazione predefinita, usa quello della pipeline catalog. |
destination_schema |
No | Schema di destinazione. Per impostazione predefinita, usa quello della pipeline schema. |
destination_table |
No | Nome della tabella di destinazione. Di default è source_table. |
Configurazione delle tabelle
| Parametro | Predefinito | Description |
|---|---|---|
primary_keys |
Rilevamento automatico | Colonne che identificano ogni riga. Rilevato automaticamente dalla chiave primaria di origine, se non specificato. |
scd_type |
SCD_TYPE_1 |
SCD_TYPE_1 mantiene solo la versione più recente.
SCD_TYPE_2 mantiene la cronologia integrale e richiede SQL Server CDC nell'origine dati. Il tipo SCD 2 non è supportato con il rilevamento delle modifiche. |
sequence_by |
Rilevamento automatico | Colonne utilizzate per ordinare gli eventi CDC. Rilevato automaticamente in base al meccanismo CDC di origine, se non specificato. |
Per il mapping dei tipi di dati di SQL Server, vedere SQL Server connector reference. Le pipeline CDC integrate supportano l'estensione automatica dei tipi: quando un tipo di colonna di origine viene esteso (ad esempio, INT a BIGINT), la tabella di destinazione si adatta automaticamente.
Monitorare la pipeline
Dopo aver creato e avviato una pipeline CDC integrata, monitorarne lo stato usando quanto segue:
Azure Databricks'interfaccia utente. Aprire la pipeline nella sezione Pipelines per visualizzare lo stato di aggiornamento, le metriche di inserimento per tabella e la derivazione.
API REST.
GET /api/2.0/pipelines/<pipeline-id>API degli eventi.
GET /api/2.0/pipelines/<pipeline-id>/events
Il primo aggiornamento della pipeline esegue uno snapshot completo di tutte le tabelle selezionate, che può richiedere più tempo rispetto agli aggiornamenti incrementali. Per le tabelle di grandi dimensioni, lo snapshot iniziale potrebbe richiedere il completamento di più aggiornamenti pianificati. Ogni aggiornamento successivo riprende dove è stato interrotto il precedente.
Per verificare l'acquisizione:
-- Check row counts in the destination table
SELECT COUNT(*) FROM <destination_catalog>.<destination_schema>.<destination_table>;
-- View recent changes (SCD Type 2 tables)
SELECT * FROM <destination_catalog>.<destination_schema>.<destination_table>
ORDER BY __START_AT DESC
LIMIT 10;
Per l'aggiornamento completo e il comportamento di aggiornamento automatico, vedere Tabelle di destinazione di aggiornamento completo.
Le pipeline CDC integrate hanno la scalabilità automatica verticale abilitata per impostazione predefinita. Se un aggiornamento della pipeline non riesce a causa di una condizione di memoria insufficiente, l'aggiornamento successivo effettua automaticamente il provisioning di un driver più grande. Per ignorare questo comportamento, usare un criterio personalizzato per il cluster.
Limitations
- Beta. Il connettore CDC integrato richiede l'abilitazione a livello di area di lavoro. Contatta il team responsabile del tuo account Azure Databricks.
- Attivato di default. Per impostazione predefinita, le pipeline CDC integrate vengono eseguite in modalità attivata; pianificarle utilizzando un'attività di Lakeflow Jobs. La modalità continua è disponibile in Beta. Vedere Eseguire una pipeline CDC integrata in modalità continua.
- Creazione esclusivamente tramite API. La creazione della pipeline può essere eseguita tramite l'API REST, la Databricks CLI, i notebook e i bundle di automazione dichiarativi. La creazione dell'interfaccia utente non è ancora supportata.
- Il canale deve essere
PREVIEW. Le specifiche della pipeline devono includere"channel": "PREVIEW". - La connessione e il tipo di connettore non sono modificabili.
connection_nameeconnector_typenon possono essere modificati dopo la creazione della pipeline. Per modificare l'origine, creare una nuova pipeline. - Massimo consigliato: 300 tabelle per pipeline.
- Solo le istanze primarie. Il connettore CDC integrato non supporta repliche in lettura, istanze di standby o istanze secondarie.
- Tabelle senza chiavi primarie. La pipeline considera tutte le colonne non LOB come una chiave composita. Le righe duplicate potrebbero essere unite in una singola riga, a meno che non si attivi SCD Type 2.
- Lo snapshot iniziale può estendersi su più aggiornamenti. Per le tabelle di grandi dimensioni, lo snapshot iniziale potrebbe non terminare in un singolo aggiornamento. Gli aggiornamenti pianificati successivi riprendono la posizione in cui è stato interrotto l'aggiornamento precedente.
- Il runtime di aggiornamento viene gestito automaticamente: La chiusura intelligente determina quando ogni aggiornamento si arresta. Un aggiornamento viene completato dopo aver raggiunto la sorgente, entro un tempo massimo di esecuzione. Vedere Chiusura intelligente per pipeline CDC integrate. Non è possibile configurare il runtime minimo o massimo. Un backlog di modifiche di grandi dimensioni potrebbe estendersi su più aggiornamenti. Gli aggiornamenti pianificati successivi riprendono la posizione in cui è stato interrotto l'aggiornamento precedente.
- L'eliminazione dei log richiede l'aggiornamento completo. Se SQL Server elimina i log di rilevamento delle modifiche o i log CDC prima dell'elaborazione della pipeline, eseguire un aggiornamento completo nelle tabelle interessate. La pipeline rileva questa condizione e rileva un errore nel registro eventi.
Risoluzione dei problemi
Annotazioni
Alcuni codici di errore usano il INGESTION_GATEWAY_ prefisso . Si tratta di una convenzione di denominazione storica e non indica che sia necessario un gateway di acquisizione separato.
| Error | Motivo | Resolution |
|---|---|---|
NOT_IN_DEFAULT_PUBLISHING_MODE |
La pipeline non è in modalità pubblicazione diretta. | La modalità di pubblicazione diretta viene impostata automaticamente per le pipeline CDC integrate. Se viene visualizzato questo errore, ricrea la pipeline. |
INGESTION_GATEWAY_CDC_NOT_ENABLED |
CDC o il rilevamento delle modifiche non è abilitato in una o più tabelle di origine. | Abilitare CDC o il rilevamento delle modifiche nelle tabelle interessate. Vedere Configurare Microsoft SQL Server per l'inserimento in Azure Databricks. |
INGESTION_GATEWAY_MISSING_TABLE_IN_SOURCE |
La tabella di origine specificata non esiste o è stata eliminata. | Verificare che la tabella esista e che l'utente di connessione abbia accesso. |
INGESTION_GATEWAY_SOURCE_SCHEMA_MISSING_ENTITY |
Lo schema di origine non esiste. | Verificare che lo schema esista nel database di origine. |
UNSUPPORTED_SOURCE_TYPE_FOR_CDC_CONNECTOR |
Il tipo di database di origine non è supportato. | Il connettore CDC integrato supporta SQL Server e Oracle. |
SOURCE_TABLE_REQUIRED |
Nella specifica della tabella manca source_table. |
Aggiungi source_table a ogni specifica della tabella nell'array objects. |
Integrated CDC connector is disabled |
Il feature flag dell'area di lavoro non è abilitato. | Contatta il team del tuo account Azure Databricks per abilitare il connettore CDC integrato nella tua area di lavoro. |
Se si verifica un problema non trattato qui:
- Esaminare il log eventi della pipeline nell'interfaccia utente di Azure Databricks o tramite
GET /api/2.0/pipelines/<pipeline-id>/events. - Testa la connessione a Unity Catalog da Catalog Explorer per confermare che la sorgente sia raggiungibile.
- Verificare che il rilevamento delle modifiche o CDC sia abilitato nel database e nelle tabelle di origine.
- Verificare che l'utente del database disponga delle autorizzazioni SQL Server elencate in database di Microsoft SQL Server requisiti utente.
- Verifica che la specifica della pipeline includa
"channel": "PREVIEW".