Faça a gestão das streams

Nesta página, vai saber como usar a API Datastream para:

  • Crie streams
  • Obtenha informações sobre streams e objetos de stream
  • Atualizar streams iniciando-as, pausando-as, retomando-as e modificando-as, bem como iniciando e parando o preenchimento para objetos de stream
  • Recupere streams com falhas definitivas
  • Ative o streaming de objetos grandes para streams Oracle
  • Elimine streams

Existem duas formas de usar a API Datastream. Pode fazer chamadas à API REST ou usar a CLI (interface de linha de comandos) do Google Cloud.

Para ver informações de alto nível sobre a utilização da Google Cloud CLI para gerir streams do Datastream, consulte o artigo Streams do Datastream da CLI gcloud.

Crie uma stream

Nesta secção, vai aprender a criar uma stream que é usada para transferir dados da sua origem para um destino. Os exemplos que se seguem não são exaustivos, mas realçam funcionalidades específicas do fluxo de dados. Para abordar o seu exemplo de utilização específico, use estes exemplos juntamente com a documentação de referência da API do fluxo de dados.

Esta secção aborda os seguintes exemplos de utilização:

Exemplo 1: transmita objetos específicos para o BigQuery

Neste exemplo, vai aprender a:

  • Transmita do MySQL para o BigQuery
  • Incluir um conjunto de objetos na stream
  • Defina o modo de escrita para a stream como apenas anexar
  • Preencha todos os objetos incluídos na stream

Segue-se um pedido para extrair todas as tabelas de schema1 e duas tabelas específicas de schema2: tableA e tableC. Os eventos são escritos num conjunto de dados no BigQuery.

O pedido não inclui o parâmetro customerManagedEncryptionKey. Por isso, o sistema de gestão de chaves interno é usado para encriptar os seus dados em vez da CMEK. Google Cloud

O parâmetro backfillAll associado à realização do preenchimento histórico (ou instantâneo) está definido como um dicionário vazio ({}), o que significa que o Datastream preenche os dados históricos de todas as tabelas incluídas na stream.

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

Para mais informações sobre como usar gcloud para criar uma stream, consulte a documentação do SDK Google Cloud.

Exemplo 2: exclua objetos específicos de uma stream com uma origem PostgreSQL

Neste exemplo, vai aprender a:

  • Faça streaming do PostgreSQL para o BigQuery
  • Exclua objetos da stream
  • Exclua objetos do preenchimento

O código seguinte mostra um pedido para criar uma stream que é usada para transferir dados de uma base de dados PostgreSQL de origem para o BigQuery. Quando cria uma stream a partir de uma base de dados PostgreSQL de origem, tem de especificar dois campos adicionais específicos do PostgreSQL no seu pedido:

  • replicationSlot: um espaço de replicação é um pré-requisito para configurar uma base de dados PostgreSQL para replicação. Tem de criar um espaço de replicação para cada stream.
  • publication: uma publicação é um grupo de tabelas a partir do qual quer replicar alterações. O nome da publicação tem de existir na base de dados antes de iniciar uma stream. No mínimo, a publicação tem de incluir as tabelas especificadas na lista includeObjects da stream.

O parâmetro backfillAll associado à realização do preenchimento histórico (ou da captura instantânea) está definido para excluir uma tabela.

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": {
      "dataFreshness": "900s",
      "sourceHierarchyDatasets": {
        "datasetTemplate": {
           "location": "us",
           "datasetIdPrefix": "prefix_"
        }
      }
    }
  },
  "backfillAll": {
    "postgresqlExcludedObjects": {
        "postgresqlSchemas": [
          { "schema": "schema1",
            "postgresqlTables": [
              { "table": "tableA" }
            ]
          }
        ]
      }
    }
  }

gcloud

Para mais informações sobre como usar gcloud para criar uma stream, consulte a documentação do SDK Google Cloud.

Exemplo 3: especifique o modo de escrita apenas de anexação para uma stream

Quando faz streaming para o BigQuery, pode definir o modo de escrita: merge ou appendOnly. Para mais informações, consulte o artigo Configure o modo de escrita.

Se não especificar o modo de escrita no seu pedido para criar uma stream, é usado o modo merge predefinido.

O pedido seguinte mostra como definir o modo appendOnly quando cria uma stream do MySQL para o BigQuery.

REST

POST https://datastream.googleapis.com/v1/projects/myProjectId1/locations/us-central1/streams?streamId=appendOnlyStream
{
  "displayName": "My append-only stream",
  "sourceConfig": {
    "sourceConnectionProfileName": "/projects/myProjectId1/locations/us-central1/streams/mysqlCp",
    "mysqlSourceConfig": {
      "includeObjects": {
        "mysqlDatabases": [
          { "database": "myMySqlDb"
          }
        ]
      }
    }
  },
  "destinationConfig": {
    "destinationConnectionProfileName": "BigQueryCp",
    "bigqueryDestinationConfig": {
      "sourceHierarchyDatasets": {
        "datasetTemplate": {
          "location": "us",
          "datasetIdPrefix": "prefix_"
        }
      },
      "appendOnly": {}
    }
  },
  "backfillAll": {}
}

gcloud

Para mais informações sobre como usar gcloud para criar uma stream, consulte a documentação do SDK Google Cloud.

Exemplo 4: faça streaming para um projeto diferente no BigQuery

Se criou os recursos da visualização de propriedade num projeto, mas quer transmitir para um projeto diferente no BigQuery, pode fazê-lo através de um pedido semelhante ao que se segue.

Se especificar sourceHierarchyDatasets para o conjunto de dados de destino, tem de preencher o campo projectId.

Se especificar singleTargetDataset para o conjunto de dados de destino, preencha o campo datasetId no formato projectId:datasetId.

REST

Para sourceHierarchyDatasets:

POST https://datastream.googleapis.com/v1/projects/myProjectId1/locations/us-central1/streams?streamId=crossProjectBqStream1
{
  "displayName": "My cross-project stream",
  "sourceConfig": {
    "sourceConnectionProfileName": "/projects/myProjectId1/locations/us-central1/streams/mysqlCp",
    "mysqlSourceConfig": {
      "includeObjects": {
        "mysqlDatabases": [
          { "database": "myMySqlDb"
          }
        ]
      }
    }
  },
  "destinationConfig": {
    "destinationConnectionProfileName": "BigQueryCp",
    "bigqueryDestinationConfig": {
      "sourceHierarchyDatasets": {
        "datasetTemplate": {
          "location": "us",
          "datasetIdPrefix": "prefix_"
        },
        "projectId": "myProjectId2"
      }
    }
  },
  "backfillAll": {}
}

Para singleTargetDataset:

POST https://datastream.googleapis.com/v1/projects/myProjectId1/locations/us-central1/streams?streamId=crossProjectBqStream2
{
  "displayName": "My cross-project stream",
  "sourceConfig": {
    "sourceConnectionProfileName": "/projects/myProjectId1/locations/us-central1/streams/mysqlCp",
    "mysqlSourceConfig": {
      "includeObjects": {
        "mysqlDatabases": [
          { "database": "myMySqlDb"
          }
        ]
      }
    }
  },
  "destinationConfig": {
    "destinationConnectionProfileName": "BigQueryCp",
    "bigqueryDestinationConfig": {
      "singleTargetDataset": {
        "datasetId": "myProjectId2:myDatasetId"
      },
    }
  },
  "backfillAll": {}
}

gcloud

Para sourceHierarchyDatasets:

  datastream streams create crossProjectBqStream1 --location=us-central1
  --display-name=my-cross-project-stream --source=source-cp --mysql-source-config=mysql_source_config.json
  --destination=destination-cp --bigquery-destination-config=source_hierarchy_cross_project_config.json
  --backfill-none
  

O conteúdo do ficheiro de configuração source_hierarchy_cross_project_config.json:

  {"sourceHierarchyDatasets": {"datasetTemplate": {"location": "us-central1", "datasetIdPrefix": "prefix_"}, "projectId": "myProjectId2"}}
  

Para singleTargetDataset:

  datastream streams create crossProjectBqStream --location=us-central1
  --display-name=my-cross-project-stream --source=source-cp --mysql-source-config=mysql_source_config.json
  --destination=destination-cp --bigquery-destination-config=single_target_cross_project_config.json
  --backfill-none
  

O conteúdo do ficheiro de configuração single_target_cross_project_config.json:

  {"singleTargetDataset": {"datasetId": "myProjectId2:myDatastetId"}}
  

Para mais informações sobre como usar gcloud para criar uma stream, consulte a documentação do SDK Google Cloud.

Exemplo 5: faça streaming para um destino do Cloud Storage

Neste exemplo, vai aprender a:

  • Faça streaming da Oracle para o Cloud Storage
  • Defina um conjunto de objetos a incluir na stream
  • Defina a CMEK para encriptar dados em repouso

O pedido seguinte mostra como criar uma stream que escreve os eventos num contentor no Cloud Storage.

Neste pedido de exemplo, os eventos são escritos no formato de saída JSON e é criado um novo ficheiro a cada 100 MB ou 30 segundos (substituindo os valores predefinidos de 50 MB e 60 segundos).

Para o formato JSON, pode:

  • Inclua um ficheiro de esquema de tipos unificados no caminho. Como resultado, o fluxo de dados escreve dois ficheiros no Cloud Storage: um ficheiro de dados JSON e um ficheiro de esquema Avro. O ficheiro de esquema tem o mesmo nome que o ficheiro de dados, com uma extensão .schema.

  • Ative a compressão gzip para que o Datastream comprima os ficheiros escritos no Cloud Storage.

Ao usar o parâmetro backfillNone, o pedido especifica que apenas as alterações em curso são transmitidas para o destino, sem preenchimento.

O pedido especifica o parâmetro da chave de encriptação gerida pelo cliente, que lhe permite controlar as chaves usadas para encriptar os dados em repouso num Google Cloud projeto. O parâmetro refere-se à CMEK que o Datastream usa para encriptar os dados transmitidos em fluxo contínuo da origem para o destino. Também especifica o conjunto de chaves para a CMEK.

Para mais informações acerca dos conjuntos de chaves, consulte o artigo Recursos do Cloud KMS. Para mais informações sobre a proteção dos seus dados através de chaves de encriptação, consulte o Cloud Key Management Service (KMS).

REST

POST https://datastream.googleapis.com/v1/projects/myProjectId1/locations/
us-central1/streams?streamId=myOracleCdcStream
{
  "displayName": "Oracle CDC to Cloud Storage",
  "sourceConfig": {
    "sourceConnectionProfileName": "/projects/myProjectId1/locations/us-central1/
    connectionProfiles/OracleCp",
    "oracleSourceConfig": {
      "includeObjects": {
        "oracleSchemas": [
          {
            "schema": "schema1"
          }