ストリームの管理

このページでは、Datastream API を使用して以下を行う方法を説明します。

  • ストリームを作成する
  • ストリームとストリーム オブジェクトに関する情報を取得する
  • ストリームを開始、一時停止、再開、変更して更新するほか、ストリーム オブジェクトのバックフィルを開始、停止して更新する
  • 恒久的な障害が発生したストリームの復元
  • Oracle ストリームでラージ オブジェクトのストリーミングを有効にする
  • ストリームの削除

Datastream API を使用する方法は 2 つあります。REST API 呼び出しを行うか、Google Cloud CLI(CLI)を使用できます。

Google Cloud CLI を使用して Datastream ストリームを管理する方法の概要については、gcloud CLI Datastream ストリームをご覧ください。

ストリームの作成

このセクションでは、ソースから宛先にデータを転送するために使用するストリームを作成する方法について説明します。次の例は包括的なものではなく、Datastream の特定の機能をハイライトしたものです。特定のユースケースに対応するには、これらの例を Datastream の API リファレンス ドキュメントと併用してください。

このセクションでは、次のユースケースについて説明します。

例 1: 特定のオブジェクトを BigQuery にストリーミングする

この例では、以下の方法について学習します。

  • MySQL から BigQuery へのストリーミング
  • ストリームにオブジェクトのセットを含める
  • ストリームの書き込みモードを追記専用として定義する
  • ストリームに含まれるすべてのオブジェクトをバックフィルする

次のリクエストは、schema1 からすべてのテーブルを、schema2 から 2 つの特定のテーブル(tableAtableC)をそれぞれ pull するリクエストです。イベントは BigQuery のデータセットに書き込まれます。

リクエストに customerManagedEncryptionKey パラメータが含まれていないため、 Google Cloud 内部の鍵管理システムを使用してデータを暗号化し、CMEK を使用します。

履歴バックフィル(またはスナップショット)の実行に関連する backfillAll パラメータが空の辞書({})に設定されます。つまり、Datastream は、ストリームに含まれるすべてのテーブルの履歴データをバックフィルします。

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

gcloud を使用してストリームを作成する方法については、Google Cloud SDK のドキュメントをご覧ください。

例 2: PostgreSQL ソースを含むストリームから特定のオブジェクトを除外する

この例では、以下の方法について学習します。

  • PostgreSQL から BigQuery へのストリーミング
  • ストリームからオブジェクトを除外する
  • バックフィルからオブジェクトを除外する

次のコードは、ソース PostgreSQL データベースから BigQuery にデータを転送するために使用するストリームの作成リクエストを示しています。 ソース PostgreSQL データベースからストリームを作成するときは、リクエストに PostgreSQL 固有の項目を 2 つ指定する必要があります。

  • replicationSlot: レプリケーション スロットは、レプリケーション用に PostgreSQL データベースを構成するための前提条件です。ストリームごとにレプリケーション スロットを作成する必要があります。
  • publication: パブリケーションは、変更を複製するテーブルのグループです。ストリームを開始する前に、パブリケーション名がデータベースに存在している必要があります。少なくとも、パブリケーションにはストリームの includeObjects リストで指定されたテーブルが含まれている必要があります。

履歴バックフィル(またはスナップショット)の実行に関連する backfillAll パラメータが、1 つのテーブルを除外するように設定されます。

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

gcloud を使用してストリームを作成する方法については、Google Cloud SDK のドキュメントをご覧ください。

例 3: ストリームの追記専用の書き込みモードを指定する

BigQuery にストリーミングするときに、書き込みモード(merge または appendOnly)を定義できます。詳細については、書き込みモードを構成するをご覧ください。

ストリームを作成するリクエストで書き込みモードを指定しない場合、デフォルトの merge モードが使用されます。

次のリクエストは、MySQL から BigQuery へのストリームを作成するときに appendOnly モードを定義する方法を示しています。

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

gcloud を使用してストリームを作成する方法については、Google Cloud SDK のドキュメントをご覧ください。

例 4: BigQuery の別のプロジェクトにストリーミングする

Datastream リソースを 1 つのプロジェクトで作成したが、BigQuery の別のプロジェクトにストリーミングする場合は、次のようなリクエストを使用してストリーミングできます。

宛先データセットに sourceHierarchyDatasets を指定する場合は、projectId フィールドに入力する必要があります。

宛先データセットに singleTargetDataset を指定する場合は、projectId:datasetId 形式で datasetId フィールドに入力します。

REST

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": {}
}

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

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
  

source_hierarchy_cross_project_config.json 構成ファイルの内容:

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

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
  

single_target_cross_project_config.json 構成ファイルの内容:

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

gcloud を使用してストリームを作成する方法については、Google Cloud SDK のドキュメントをご覧ください。

例 5: Cloud Storage の宛先にストリーミングする

この例では、以下の方法について学習します。

  • Oracle から Cloud Storage へのストリーミング
  • ストリームに含めるオブジェクトのセットを定義する
  • 保存データを暗号化するための CMEK を定義する

次のリクエストは、Cloud Storage のバケットにイベントを書き込むストリームを作成する方法を示しています。

このリクエスト例では、イベントは JSON 出力形式で書き込まれ、100 MB または 30 秒ごとに新しいファイルが作成されます(デフォルト値の 50 MB と 60 秒をオーバーライドしています)。

JSON 形式では、次のことが可能です。

  • パスに統合型スキーマ ファイルを含めるその結果、データストリームによって JSON データファイルと Avro スキーマ ファイルの 2 つのファイルが Cloud Storage に書き込まれます。スキーマ ファイルは、データファイルと同じ名前で、拡張子は .schema です。

  • gzip 圧縮を有効にする。これによって、Cloud Storage に書き込まれたファイルをデータストリームが圧縮するようにします。

backfillNone パラメータを使用すると、このリクエストでは、バックフィルなしに、進行中の変更のみが宛先にストリーミングされることが指定されます。

リクエストでは、顧客管理の暗号鍵パラメータを指定します。これにより、 Google Cloud プロジェクト内の保存データの暗号化に使用する鍵を制御できます。このパラメータは、ソースから宛先にストリーミングされるデータの暗号化に Datastream が使用する CMEK を指します。また、CMEK のキーリングも指定します。

キーリングの詳細については、Cloud KMS リソースをご覧ください。暗号鍵を使用したデータの保護の詳細については、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"
          }
        ]
      }
    }
  },
  "destinationConfig": {
    "destinationConnectionProfileName": "GcsBucketCp",
    "gcsDestinationConfig": {
      "path": "/folder1",
      "jsonFileFormat": {
        "schemaFileFormat": "AVRO_SCHEMA_FILE"
      },
      "fileRotationMb": 100,
      "fileRotationInterval": 30
    }
  },
  "customerManagedEncryptionKey": "projects/myProjectId1/locations/us-central1/
  keyRings/myRing/cryptoKeys/myEncryptionKey",
  "backfillNone": {}
}

gcloud

gcloud を使用してストリームを作成する方法については、Google Cloud SDK のドキュメントをご覧ください。

例 6: Apache Iceberg テーブルにストリーミングする

この例では、append-only モードで MySQL データベースから Apache Iceberg テーブルにデータを複製するようにストリームを構成する方法について説明します。リクエストを作成する前に、次の手順を完了していることを確認してください。

  • データを保存する Cloud Storage バケットがある
  • Cloud リソース接続を作成する
  • Cloud Storage バケットへのアクセス権を Cloud リソース接続に付与する

次のリクエストを使用してストリームを作成できます。

REST

POST https://datastream.googleapis.com/v1/projects/myProjectId1/locations/us-central1/streams?streamId=mysqlIcebergStream
{
  "displayName": "MySQL to Apache Iceberg stream",
  "sourceConfig": {
    "sourceConnectionProfileName": "/projects/myProjectId1/locations/us-central1/streams/mysqlIcebergCp",
    "mysqlSourceConfig": {
      "includeObjects": {
        "mysqlDatabases": [
          {
            "database": "my-mysql-database"
          }
        ]
      }
    }
  },
  "destinationConfig": {
    "destinationConnectionProfileName": "projects/myProjectId1/locations/us-central1/connectionProfiles/my-bq-cp-id",
    "bigqueryDestinationConfig": {
      "blmtConfig": {
        "bucket": "my-gcs-bucket-name",
        "rootPath": "my/folder",
        "connectionName": "my-project-id.us-central1.my-bigquery-connection-name",
        "fileFormat": "PARQUET",
        "tableFormat": "ICEBERG"
        },
      "singleTargetDataset": {
        "datasetId": "my-project-id:my-bigquery-dataset-id"
      },
      "appendOnly": {}
    }
  },
  "backfillAll": {}
}

gcloud

datastream streams create mysqlIcebergStream --location=us-central1
--display-name=mysql-to-bl-stream --source=source --mysql-source-config=mysql_source_config.json
--destination=destination --bigquery-destination-config=bl_config.json
--backfill-none

mysql_source_config.json ソース構成ファイルの内容:

{"excludeObjects": {}, "includeObjects": {"mysqlDatabases":[{"database":"my-mysql-database"}]}}

bl_config.json 構成ファイルの内容:

{ "blmtConfig": { "bucket": "my-gcs-bucket-name", "rootPath": "my/folder", "connectionName": "my-project-id.us-central1.my-bigquery-connection-name", "fileFormat": "PARQUET", "tableFormat": "ICEBERG" }, "singleTargetDataset": {"datasetId": "my-project-id:my-bigquery-dataset-id"}, "appendOnly": {} }

Terraform

resource "google_datastream_stream" "stream" {
  stream_id    = "mysqlBlStream"
  location     = "us-central1"
  display_name = "MySQL to Apache Iceberg stream"

  source_config {
    source_connection_profile = "/projects/myProjectId1/locations/us-central1/streams/mysqlBlCp"
    mysql_source_config {
      include_objects {
        mysql_databases {
          database = "my-mysql-database"
        }
      }
    }
  }

  destination_config {
    destination_connection_profile = "projects/myProjectId1/locations/us-central1/connectionProfiles/my-bq-cp-id"
    bigquery_destination_config {
      single_target_dataset {
        dataset_id = "my-project-id:my-bigquery-dataset-id"
      }
      blmt_config {
        bucket          = "my-gcs-bucket-name"
        table_format    = "ICEBERG"
        file_format     = "PARQUET"
        connection_name = "my-project-id.us-central1.my-bigquery-connection-name"
        root_path       = "my/folder"
      }
      append_only {}
    }
  }

  backfill_none {}
}
    

ストリームの定義を検証する

ストリームは、作成する前にその定義を検証できます。これにより、すべての検証チェックに合格し、作成時にストリームが正常に実行されることが確実になります。

ストリームの検証では、次のことを確認します。

  • データストリームがソースからデータをストリーミングできるようにソースが適切に構成されているかどうか
  • ストリームがソースと宛先の両方に接続できるかどうか
  • ストリームのエンドツーエンド構成

ストリームを検証するには、リクエストの本文の前の URL に &validate_only=true を追加します。

POST "https://datastream.googleapis.com/v1/projects/PROJECT_ID/locations/LOCATION/streams?streamId=STREAM_ID&validate_only=true"

このリクエストを行うと、Datastream がソースと宛先に対して実行する検証チェックと、チェックの合否が表示されます。検証で不合格となった場合は、失敗の理由と問題を解決する方法に関する情報が表示されます。

たとえば、ソースから宛先にストリーミングされるデータの暗号化に Datastream で使用する顧客管理の暗号鍵(CMEK)があるとします。ストリームの検証の一環として、Datastream はキーが存在することと、Datastream にキーを使用する権限があることを確認します。いずれかの条件が満たされていない場合、ストリームを検証すると次のエラー メッセージが返されます。

CMEK_DOES_NOT_EXIST_OR_MISSING_PERMISSIONS

この問題を解決するには、指定した鍵が存在し、Datastream サービス アカウントにその鍵に対する cloudkms.cryptoKeys.get 権限があることを確認します。

適切な修正を行った後、リクエストを再度実行して、すべての検証チェックに合格することを確認します。上記の例では、CMEK_VALIDATE_PERMISSIONS チェックでエラー メッセージが返されなくなり、ステータスが PASSED になります。

ストリームに関する情報を取得する

次のコマンドでは、ストリームに関する情報を取得するリクエストを示します。これには以下の情報が含まれます。

  • ストリームの名前(固有識別子)
  • ストリームのユーザー フレンドリーな名前(表示名)
  • ストリームの作成日時と最終更新日時のタイムスタンプ
  • ストリームに関連付けられたソースと宛先の接続プロファイルに関する情報
  • ストリームの状態

REST

GET https://datastream.googleapis.com/v1/projects/PROJECT_ID/locations/LOCATION/streams/STREAM_ID

レスポンスは次のように表示されます。

{
  "name"