Gestione dei flussi

In questa pagina imparerai a utilizzare l'API Datastream per:

  • Creare stream
  • Ricevere informazioni su flussi e oggetti di flusso
  • Aggiorna i flussi di dati avviandoli, mettendoli in pausa, riprendendoli e modificandoli, nonché avviando e interrompendo il backfill per gli oggetti del flusso di dati
  • Recuperare gli stream con errori permanenti
  • Attiva lo streaming di oggetti di grandi dimensioni per i flussi Oracle
  • Eliminare gli stream

Esistono due modi per utilizzare l'API Datastream. Puoi effettuare chiamate API REST o utilizzare Google Cloud CLI (CLI).

Per informazioni di alto livello sull'utilizzo di Google Cloud CLI per gestire i flussi Datastream, consulta gcloud CLI Datastream streams.

Crea uno stream

In questa sezione imparerai a creare un flusso utilizzato per trasferire i dati dall'origine a una destinazione. Gli esempi che seguono non sono esaustivi, ma mettono in evidenza funzionalità specifiche di Datastream. Per risolvere il tuo caso d'uso specifico, utilizza questi esempi insieme alla documentazione di riferimento dell'API Datastream.

Questa sezione tratta i seguenti casi d'uso:

Esempio 1: trasmetti in streaming oggetti specifici a BigQuery

In questo esempio imparerai a:

  • Flusso di dati da MySQL a BigQuery
  • Includere un insieme di oggetti nello stream
  • Definisci la modalità di scrittura per lo stream come di sola aggiunta
  • Esegui il backfill di tutti gli oggetti inclusi nel flusso

Di seguito è riportata una richiesta per estrarre tutte le tabelle da schema1 e due tabelle specifiche da schema2: tableA e tableC. Gli eventi vengono scritti in un set di dati in BigQuery.

La richiesta non include il parametro customerManagedEncryptionKey, pertanto il sistema di gestione delle chiavi interno Google Cloud viene utilizzato per criptare i dati anziché CMEK.

Il parametro backfillAll associato all'esecuzione del backfill (o dello snapshot) cronologico è impostato su un dizionario vuoto ({}), il che significa che Datastream esegue il backfill dei dati cronologici di tutte le tabelle incluse nel flusso.

REST

POST https://datastream.googleapis.com/v1/projects/myProjectId1/locations/us-central1/streams?streamId=mysqlCdcStream
{
  "displayName": "MySQL CDC to BigQuery",
  "sourceConfig": {
    "sourceConnectionProfileName": "/projects/myProjectId1/locations/us-central1/streams/mysqlCp",
    "mysqlSourceConfig": {
      "includeObjects": {
        "mysqlDatabases": [
          { "database": "schema1" },
          {
            "database": "schema2",
            "mysqlTables": [
              {
                "table": "tableA",
                "table": "tableC"
              }
            ]
          }
        ]
      },
    }
  },
  "destinationConfig": {
    "destinationConnectionProfileName": "BigQueryCp",
    "bigqueryDestinationConfig": {
      "sourceHierarchyDatasets": {
        "datasetTemplate": {
          "location": "us",
          "datasetIdPrefix": "prefix_"
        }
      },
      "dataFreshness": "900s"
    }
  },
  "backfillAll": {}
}

gcloud

Per ulteriori informazioni sull'utilizzo di gcloud per creare uno stream, consulta la documentazione di Google Cloud SDK.

Esempio 2: escludi oggetti specifici da un flusso con un'origine PostgreSQL

In questo esempio imparerai a:

  • Flusso di dati da PostgreSQL a BigQuery
  • Escludere oggetti dallo stream
  • Escludere oggetti dal backfill

Il seguente codice mostra una richiesta per creare uno stream utilizzato per trasferire i dati da un database PostgreSQL di origine a BigQuery. Quando crei uno stream da un database PostgreSQL di origine, devi specificare due campi aggiuntivi specifici di PostgreSQL nella richiesta:

  • replicationSlot: uno slot di replica è un prerequisito per la configurazione di un database PostgreSQL per la replica. Devi creare uno slot di replica per ogni stream.
  • publication: una pubblicazione è un gruppo di tabelle da cui vuoi replicare le modifiche. Il nome della pubblicazione deve esistere nel database prima di avviare un flusso. Come minimo, la pubblicazione deve includere le tabelle specificate nell'elenco includeObjects del flusso.

Il parametro backfillAll associato all'esecuzione del backfill storico (o dello snapshot) è impostato per escludere una tabella.

REST

POST https://datastream.googleapis.com/v1/projects/myProjectId1/locations/
us-central1/streams?streamId=myPostgresStream
{
  "displayName": "PostgreSQL to BigQueryCloud Storage",
  "sourceConfig": {
    "sourceConnectionProfileName": "/projects/myProjectId1/locations/us-central1/connectionProfiles/postgresCp",
    "postgresqlSourceConfig": {
      "replicationSlot": "replicationSlot1",
      "publication": "publicationA",
      "includeObjects": {
        "postgresqlSchemas": {
          "schema": "schema1"
        }
      },
      "excludeObjects": {
        "postgresqlSchemas": [
          { "schema": "schema1",
        "postgresqlTables": [
          {
            "table": "tableA",
            "postgresqlColumns": [
              { "column": "column5" }
              ]
              }
            ]
          }
        ]
      }
    }
  },
  "destinationConfig": {
    "destinationConnectionProfileName": "BigQueryCp",
    "bigqueryDestinationConfig": {