Il modello di acquisizione dei dati di modifica da MySQL a BigQuery tramite Debezium e Pub/Sub è una pipeline di flussi di dati che legge i messaggi Pub/Sub con dati di modifica da un database MySQL e scrive i record su BigQuery. Un connettore Debezium acquisisce le modifiche al database MySQL e pubblica i dati modificati su Pub/Sub. Il modello legge quindi i messaggi Pub/Sub e li scrive su BigQuery.
Puoi utilizzare questo modello per sincronizzare i database MySQL e le tabelle BigQuery. La pipeline scrive i dati modificati in una tabella di gestione temporanea BigQuery e aggiorna periodicamente una tabella BigQuery che replica il database MySQL.
Requisiti della pipeline
- Il connettore Debezium deve essere sottoposto a deployment.
- I messaggi Pub/Sub devono essere serializzati in una riga Beam.
Parametri del modello
Parametri obbligatori
- inputSubscriptions: l'elenco separato da virgole delle sottoscrizioni di input Pub/Sub da cui leggere, nel formato
<SUBSCRIPTION_NAME>,<SUBSCRIPTION_NAME>, .... - changeLogDataset: il set di dati BigQuery in cui archiviare le tabelle di gestione temporanea, nel formato <DATASET_NAME>.
- replicaDataset: la posizione del set di dati BigQuery in cui archiviare le tabelle di replica, nel formato <DATASET_NAME>.
Parametri facoltativi
- inputTopics: elenco separato da virgole degli argomenti Pub/Sub in cui vengono inseriti i dati CDC.
- updateFrequencySecs: l'intervallo in cui la pipeline aggiorna la tabella BigQuery che replica il database MySQL.
- useSingleTopic: imposta questo valore su
truese configuri il connettore Debezium in modo che pubblichi tutti gli aggiornamenti delle tabelle in un singolo argomento. Il valore predefinito è false. - useStorageWriteApi: se è true, la pipeline utilizza l'API BigQuery Storage Write (https://cloud.google.com/bigquery/docs/write-api). Il valore predefinito è
false. Per ulteriori informazioni, vedi Utilizzo dell'API Storage Write (https://beam.apache.org/documentation/io/built-in/google-bigquery/#storage-write-api). - useStorageWriteApiAtLeastOnce: quando utilizzi l'API Storage Write, specifica la semantica di scrittura. Per utilizzare la semantica almeno una volta (https://beam.apache.org/documentation/io/built-in/google-bigquery/#at-least-once-semantics), imposta questo parametro su
true. Per utilizzare la semantica esattamente una volta, imposta il parametro sufalse. Questo parametro si applica solo quandouseStorageWriteApiètrue. Il valore predefinito èfalse. - numStorageWriteApiStreams: quando utilizzi l'API Storage Write, specifica il numero di flussi di scrittura. Se
useStorageWriteApiètrueeuseStorageWriteApiAtLeastOnceèfalse, devi impostare questo parametro. Il valore predefinito è 0. - storageWriteApiTriggeringFrequencySec: quando utilizzi l'API Storage Write, specifica la frequenza di attivazione, in secondi. Se
useStorageWriteApiètrueeuseStorageWriteApiAtLeastOnceèfalse, devi impostare questo parametro.
Esegui il modello
Per eseguire questo modello:
- Clona il repository DataflowTemplates sulla tua macchina locale.
- Passa alla directory
v2/cdc-parent. - Assicurati che il connettore Debezium sia sottoposto a deployment.
- Utilizza Maven per eseguire il modello Dataflow:
mvn exec:java -pl cdc-change-applier -Dexec.args="--runner=DataflowRunner \ --inputSubscriptions=SUBSCRIPTIONS \ --updateFrequencySecs=300 \ --changeLogDataset=CHANGELOG_DATASET \ --replicaDataset=REPLICA_DATASET \ --project=PROJECT_ID \ --region=REGION_NAME"
Sostituisci quanto segue:
PROJECT_ID: l'ID progetto Google Cloud in cui vuoi eseguire il job DataflowSUBSCRIPTIONS: l'elenco separato da virgole dei nomi delle sottoscrizioni Pub/SubCHANGELOG_DATASET: il set di dati BigQuery per i dati del log delle modificheREPLICA_DATASET: il set di dati BigQuery per le tabelle di replica
Passaggi successivi
- Scopri di più sui modelli Dataflow.
- Consulta l'elenco di modelli forniti da Google.