串流資料產生器範本

串流資料產生器範本會產生合成記錄或訊息,並傳送至目的地接收器。您可以設定記錄結構定義和記錄產生速率。

範本支援下列接收器:

  • Apache Kafka 主題
  • BigQuery 資料表
  • Cloud Storage bucket
  • Java Database Connectivity (JDBC) 端點
  • Pub/Sub 主題
  • Spanner 資料表

以下列舉幾種可能的用途:

  • 模擬大規模即時事件發布至 Pub/Sub 主題,以評估及判斷處理發布事件所需的消費者數量和大小。
  • 生成合成資料,評估效能基準或做為概念驗證。
  • 驗證端對端管道。舉例來說,您可以將記錄傳送至 Kafka 主題,然後由下游消費者讀取。

管道相關規定

定義記錄結構定義

範本會為產生的資料提供預先定義的結構。如要使用這個結構定義,請將 schemaTemplate 範本參數設為 GAME_EVENT

或者,您也可以提供自己的資料結構定義,方法如下:

  1. 建立結構定義檔案,其中包含所產生資料的 JSON 範本。這個範本使用 JSON 資料產生器程式庫,支援多種隨機產生資料的函式。例如:

    {
      "id": {{integer(0,1000)}},
      "name": "{{uuid()}}",
      "isInStock": {{bool()}}
    }

    詳情請參閱 json-data-generator 說明文件

  2. 將結構定義檔案上傳至 Cloud Storage bucket。
  3. schemaLocation 範本參數設為範本檔案的 Cloud Storage URI。

指定輸出格式

根據預設,範本會產生 JSON 資料。對於部分目的地,範本也支援 Avro 或 Parquet 格式:

  • Avro:支援 Cloud Storage、Apache Kafka 和 Pub/Sub
  • Parquet:支援 Cloud Storage。

如要輸出 Avro 或 Parquet 格式,請按照下列步驟操作:

  1. outputType 範本參數設為 AVRO (適用於 Avro 格式) 或 PARQUET (適用於 Parquet 格式)。
  2. 建立 Avro 結構定義檔案。
  3. 將結構定義檔案上傳至 Cloud Storage。
  4. avroSchemaLocation 範本參數設為結構定義檔案的 Cloud Storage URI。

指定目的地接收器

以下各節說明如何為各類接收器設定範本。

Apache Kafka 主題

如要寫入 Kafka 主題,請設定下列範本參數:

  • sinkType: KAFKA.
  • bootstrapServer:Kafka 叢集的啟動位址。
  • kafkaTopic:要寫入的 Kafka 主題。

如果您要寫入 Google Cloud Managed Service for Apache Kafka 叢集,請將「代管 Kafka 用戶端」(roles/managedkafka.client) 角色授予工作人員服務帳戶。

BigQuery 資料表

如要寫入 BigQuery 資料表,請設定下列範本參數:

  • sinkType: BIGQUERY.
  • outputTableSpec:要寫入的 BigQuery 資料表。請按照下列格式設定這個參數: PROJECT_ID:DATASET.TABLE

以下為選用參數:

  • outputDeadletterTable:管道寫入失敗記錄的資料表名稱。如果未指定,管道會建立名為 OUTPUT_TABLE_error_records 的資料表,其中 OUTPUT_TABLE 是輸出資料表的名稱。
  • writeDisposition:指定如何寫入現有資料表。支援的值如下:

    • WRITE_APPEND。將資料列附加至現有資料表。
    • WRITE_TRUNCATE。截斷現有資料列。
    • WRITE_EMPTY. 僅在資料表空白時寫入。如果資料表已有資料,工作就會失敗。

    預設值為 WRITE_APPEND

將 BigQuery 資料編輯者 (roles/bigquery.dataEditor) 角色授予 工作站服務帳戶

Cloud Storage

如要寫入 Cloud Storage bucket,請設定下列範本參數:

  • sinkType: GCS.
  • outputDirectory:要寫入的 Cloud Storage 資料夾路徑。

以下為選用參數:

  • numShards:分片數量上限。值越高,處理量就越高,但資料匯總費用可能會更高。如果值為 0,Dataflow 會選取分片數量。 預設值為 0。
  • outputFilenamePrefix:檔案名稱前置字串。預設值為 output-
  • windowDuration:管道將檔案寫入 Cloud Storage 的間隔。允許的格式為 Ns (秒)、Nm (分鐘) 和 Nh (小時)。預設值為 1m (1 分鐘)。

將 Storage 物件管理員 (roles/storage.objectAdmin) 角色授予 工作人員服務帳戶

JDBC 端點

如要寫入 JDBC 端點,請設定下列範本參數:

  • sinkType: JDBC.
  • driverClassName:要使用的 JDBC 驅動程式類別。示例: com.mysql.jdbc.Driver
  • connectionUrl:用於連線至 JDBC 來源的連線字串。
  • statement:用於寫入資料庫的 INSERT INTO SQL 陳述式。陳述式必須指定要寫入的資料表資料欄,並為 VALUES 子句使用 '?' 字元做為預留位置。管道會將預留位置替換為 JSON 資料中的對應欄位值。

    示例: INSERT INTO tableName (column1, column2) VALUES (?,?)

以下為選用參數:

  • username:JDBC 連線的使用者名稱。
  • password:JDBC 連線的密碼。
  • connectionProperties:JDBC 連線的屬性字串。範例:unicode=true;characterEncoding=UTF-8

Pub/Sub 主題

如要寫入 Pub/Sub 主題,請設定下列範本參數:

  • sinkType: PUBSUB.
  • topic:要寫入的 Pub/Sub 主題。

將 Pub/Sub 發布者 (roles/pubsub.publisher) 角色授予 工作人員服務帳戶

Spanner 資料表

如要寫入 Spanner 資料表,請設定下列範本參數:

  • sinkType: SPANNER.
  • projectId:包含 Spanner 表格的專案 ID。
  • spannerInstanceName:Spanner 執行個體的名稱。
  • spannerDatabaseName:Spanner 資料庫的名稱。
  • spannerTableName:Spanner 資料表名稱。

以下為選用參數:

  • maxNumMutations:每個批次中突變儲存格的數量上限。
  • maxNumRows:每個批次中變動的資料列數量上限。
  • batchSizeBytes:每個批次變動的位元組數上限。
  • commitDeadlineSeconds:提交 API 呼叫的期限,以秒為單位。

將 Cloud Spanner 資料庫使用者 (roles/spanner.databaseUser) 角色授予 工作人員服務帳戶

範本參數

必要參數

  • qps:表示每秒發布至 Pub/Sub 的訊息速率。

選用參數

  • schemaTemplate:要使用的現有結構定義範本。值必須是 [GAME_EVENT] 其中之一。
  • schemaLocation:結構定義位置的 Cloud Storage 路徑。例如:gs://<bucket-name>/prefix
  • topic:管道應發布資料的主題名稱。例如:projects/<project-id>/topics/<topic-name>
  • messagesLimit:表示要生成的輸出訊息數量上限。0 代表無限制。預設為:0。
  • outputType:訊息輸出類型。預設值為 JSON。
  • avroSchemaLocation:Avro 結構定義位置的 Cloud Storage 路徑。如果輸出類型為 AVRO 或 PARQUET,則為必填欄位。例如:gs://your-bucket/your-path/schema.avsc
  • sinkType:訊息接收器類型。預設值為 PUBSUB。
  • outputTableSpec:輸出 BigQuery 資料表。當 sinkType 為 BIGQUERY 時,這是必填欄位。例如 <project>:<dataset>.<table_name>
  • writeDisposition:BigQuery WriteDisposition。例如 WRITE_APPEND、WRITE_EMPTY 或 WRITE_TRUNCATE。預設值為:WRITE_APPEND。
  • outputDeadletterTable:訊息無法到達輸出資料表的所有原因 (例如結構定義不相符、JSON 格式錯誤) 會寫入此資料表。如果該資料表不存在,將會於管道執行期間建立。例如:your-project-id:your-dataset.your-table-name
  • windowDuration:資料寫入 Cloud Storage 的時間間隔/大小。允許的格式為 Ns (以秒為單位例如 5s)、Nm (以分鐘為單位例如 12m)、Nh (以小時為單位 2h)。例如,1m。預設值為 1 分鐘。
  • outputDirectory:寫入輸出檔案的路徑和檔案名稱前置字串,結尾必須為斜線。日期時間格式設定用於剖析目錄路徑,以取得日期和時間格式設定器。例如:gs://your-bucket/your-path/
  • outputFilenamePrefix:加在每個固定時段檔案的前置字元,例如,output-。預設值為:output-。
  • numShards:寫入時產生的輸出資料分割數量上限;資料分割數量越多,寫入 Cloud Storage 的處理量就越高,但處理輸出 Cloud Storage 檔案時,資料分割間的資料匯總費用可能會更高。預設值由 Dataflow 決定。
  • driverClassName:要使用的 JDBC 驅動程式類別名稱。例如:com.mysql.jdbc.Driver
  • connectionUrl:用於連線至 JDBC 來源的網址連線字串。例如:jdbc:mysql://some-host:3306/sampledb
  • username:JDBC 連線要使用的使用者名稱。
  • password:JDBC 連線使用的密碼。
  • connectionProperties:JDBC 連線的屬性字串,字串格式必須為 [propertyName=property;]*。例如:unicode=true;characterEncoding=UTF-8
  • 陳述式:將執行的 SQL 陳述式,用於寫入資料庫。陳述式必須指定資料表的資料欄名稱,順序不限。系統只會從 JSON 讀取指定資料欄名稱的值,並加入陳述式。例如:INSERT INTO tableName (column1, column2) VALUES (?,?)
  • projectId:Spanner 表格所在的 GCP 專案 ID。
  • spannerInstanceName:Cloud Spanner 執行個體名稱。
  • spannerDatabaseName:Cloud Spanner 資料庫名稱。
  • spannerTableName:Cloud Spanner 資料表名稱。
  • maxNumMutations:指定儲存格修改作業限制 (每個批次修改作業的儲存格數量上限)。預設值為 5000。
  • maxNumRows:指定資料列修改作業限制 (每個批次修改作業的資料列數量上限)。預設值為 1000。
  • batchSizeBytes:指定批量大小限制 (每個批量變動的位元組數量上限)。預設值為 1MB。
  • commitDeadlineSeconds:指定 Commit API 呼叫的截止時間 (以秒為單位)。
  • bootstrapServer:Kafka Bootstrap 伺服器,例如 localhost:9092
  • kafkaTopic:要寫入的 Kafka 主題。例如:topic

執行範本

控制台

  1. 前往 Dataflow 的「Create job from template」(依據範本建立工作) 頁面。
  2. 前往「依範本建立工作」
  3. 在「Job name」(工作名稱) 欄位中,輸入專屬工作名稱。
  4. 選用:如要使用區域端點,請從下拉式選單中選取值。預設區域為 us-central1

    如需可執行 Dataflow 工作的區域清單,請參閱「Dataflow 位置」。

  5. 從「Dataflow template」(Dataflow 範本) 下拉式選單中,選取「Streaming Data Generator」(串流資料產生器) 範本。
  6. 在提供的參數欄位中輸入參數值。
  7. 按一下「Run Job」(執行工作)

gcloud

在殼層或終端機中執行範本:

gcloud dataflow flex-template run JOB_NAME \
    --project=PROJECT_ID \
    --region=REGION_NAME \
    --template-file-gcs-location=gs://dataflow-templates-REGION_NAME/VERSION/flex/ \
    --parameters \
schemaLocation=SCHEMA_LOCATION,\
qps=QPS,\
topic=PUBSUB_TOPIC
  

更改下列內容:

  • PROJECT_ID:您要執行 Dataflow 工作的 Google Cloud 專案 ID
  • REGION_NAME:您要部署 Dataflow 工作的區域,例如 us-central1
  • JOB_NAME:您選擇的不重複工作名稱
  • VERSION:您要使用的範本版本

    您可以使用下列值:

  • SCHEMA_LOCATION:Cloud Storage 中結構定義檔案的路徑。例如:gs://mybucket/filename.json
  • QPS:每秒發布的訊息數量
  • PUBSUB_TOPIC:輸出 Pub/Sub 主題。例如:projects/my-project-id/topics/my-topic-id

API

如要使用 REST API 執行範本,請傳送 HTTP POST 要求。如要進一步瞭解 API 和授權範圍,請參閱 projects.templates.launch

POST https://dataflow.googleapis.com/v1b3/projects/PROJECT_ID/locations/LOCATION/flexTemplates:launch
{
   "launch_parameter": {
      "jobName": "JOB_NAME",
      "parameters": {
          "schemaLocation": "SCHEMA_LOCATION",
          "qps": "QPS",
          "topic": "PUBSUB_TOPIC"
      },
      "containerSpecGcsPath": "gs://dataflow-templates-LOCATION/VERSION/flex/",
   }
}
  

更改下列內容:

  • PROJECT_ID:您要執行 Dataflow 工作的 Google Cloud 專案 ID
  • LOCATION:您要部署 Dataflow 工作的區域,例如 us-central1
  • JOB_NAME:您選擇的不重複工作名稱
  • VERSION:您要使用的範本版本

    您可以使用下列值:

  • SCHEMA_LOCATION:Cloud Storage 中結構定義檔案的路徑。例如:gs://mybucket/filename.json
  • QPS:每秒發布的訊息數量
  • PUBSUB_TOPIC:輸出 Pub/Sub 主題。例如:projects/my-project-id/topics/my-topic-id

後續步驟