Template Streaming Data Generator membuat pesan atau rekaman sintetis dan mengirimkannya ke sink tujuan. Anda dapat mengonfigurasi skema rekaman dan kecepatan pembuatan rekaman.
Template ini mendukung tujuan berikut:
- Topik Apache Kafka
- Tabel BigQuery
- Bucket Cloud Storage
- Endpoint Java Database Connectivity (JDBC)
- Topik Pub/Sub
- Tabel Spanner
Berikut adalah beberapa kemungkinan kasus penggunaan:
- Mensimulasikan publikasi peristiwa real-time skala besar ke topik Pub/Sub untuk mengukur dan menentukan jumlah serta ukuran konsumen yang diperlukan untuk memproses peristiwa yang dipublikasikan.
- Buat data sintetis untuk mengevaluasi tolok ukur performa atau berfungsi sebagai bukti konsep.
- Memvalidasi pipeline end-to-end. Misalnya, mengirimkan rekaman ke topik Kafka, yang kemudian dibaca oleh konsumen hilir.
Persyaratan pipeline
- Akun layanan worker
memerlukan peran yang ditetapkan Dataflow Worker (
roles/dataflow.worker). Untuk informasi selengkapnya, lihat Pengantar IAM. - Sink tujuan harus sudah ada.
- Konfigurasi template untuk sink tujuan. Untuk mengetahui informasi selengkapnya, lihat Menentukan sink tujuan.
- Tentukan skema data.
- Template menghasilkan data JSON secara default. Untuk membuat Avro atau Parquet, lihat Menentukan format output.
Menentukan skema rekaman
Template ini menyediakan skema yang telah ditentukan sebelumnya untuk data yang dihasilkan. Untuk menggunakan
skema ini, tetapkan parameter template schemaTemplate ke
GAME_EVENT.
Atau, Anda dapat memberikan skema data Anda sendiri seperti berikut:
-
Buat file skema yang berisi template JSON untuk data yang dihasilkan. Template ini menggunakan library JSON Data Generator, yang mendukung berbagai fungsi untuk mengacak data. Contoh:
{ "id": {{integer(0,1000)}}, "name": "{{uuid()}}", "isInStock": {{bool()}} }
Untuk mengetahui informasi selengkapnya, lihat dokumentasi json-data-generator.
- Upload file skema ke bucket Cloud Storage.
- Tetapkan parameter template
schemaLocationke Cloud Storage URI file template.
Tentukan format output
Secara default, template menghasilkan data JSON. Untuk beberapa tujuan, template juga mendukung format Avro atau Parquet:
- Avro: Didukung untuk Cloud Storage, Apache Kafka, dan Pub/Sub
- Parquet: Didukung untuk Cloud Storage.
Untuk menghasilkan format Avro atau Parquet, lakukan hal berikut:
- Tetapkan parameter template
outputTypekeAVROuntuk format Avro, atauPARQUETuntuk format Parquet. - Buat file skema Avro.
- Upload file skema ke Cloud Storage.
- Tetapkan parameter template
avroSchemaLocationke Cloud Storage URI file skema.
Tentukan sink tujuan
Bagian berikut menjelaskan cara mengonfigurasi template untuk setiap jenis tujuan.
- Topik Apache Kafka
- BigQuery table
- Bucket Cloud Storage
- Endpoint Java Database Connectivity (JDBC)
- Topik Pub/Sub
- Tabel Spanner
Topik Apache Kafka
Untuk menulis ke topik Kafka, tetapkan parameter template berikut:
sinkType:KAFKA.bootstrapServer: Alamat bootstrap cluster Kafka.kafkaTopic: Topik Kafka yang akan ditulisi.
Jika Anda menulis ke cluster Google Cloud Managed Service untuk Apache Kafka, berikan peran Managed Kafka Client
(roles/managedkafka.client) kepada
akun layanan pekerja.
Tabel BigQuery
Untuk menulis ke tabel BigQuery, tetapkan parameter template berikut:
sinkType:BIGQUERY.outputTableSpec: Tabel BigQuery yang akan ditulisi. Format parameter ini sebagai berikut:PROJECT_ID:DATASET.TABLE.
Parameter berikut bersifat opsional:
outputDeadletterTable: Nama tabel tempat pipeline menulis catatan yang gagal. Jika tidak ditentukan, pipeline akan membuat tabel bernamaOUTPUT_TABLE_error_records, denganOUTPUT_TABLEadalah nama tabel output.writeDisposition: Menentukan cara menulis ke tabel yang ada. Nilai-nilai berikut didukung:WRITE_APPEND. Menambahkan baris ke tabel yang ada.WRITE_TRUNCATE. Memangkas baris yang ada.WRITE_EMPTY. Tulis hanya jika tabel kosong. Jika tabel sudah memiliki data, tugas akan gagal.
Nilai defaultnya adalah
WRITE_APPEND.
Berikan peran
akun layanan pekerja peran BigQuery Data Editor
(roles/bigquery.dataEditor).
Cloud Storage
Untuk menulis ke bucket Cloud Storage, tetapkan parameter template berikut:
sinkType:GCS.outputDirectory: Jalur folder Cloud Storage yang akan ditulis.
Parameter berikut bersifat opsional:
numShards: Jumlah maksimum shard. Nilai yang lebih tinggi dapat memungkinkan throughput yang lebih tinggi, tetapi berpotensi menimbulkan biaya agregasi data yang lebih tinggi. Jika nilainya 0, Dataflow akan memilih jumlah shard. Nilai defaultnya adalah 0.outputFilenamePrefix: Awalan nama file. Nilai defaultnya adalahoutput-.windowDuration: Interval saat pipeline menulis file ke Cloud Storage. Format yang diizinkan adalahNs(detik),Nm(menit), danNh(jam). Nilai defaultnya adalah1m(1 menit).
Berikan peran Storage Object Admin (roles/storage.objectAdmin) ke
akun layanan pekerja.
Endpoint JDBC
Untuk menulis ke endpoint JDBC, tetapkan parameter template berikut:
sinkType:JDBC.driverClassName: Class driver JDBC yang akan digunakan. Contoh:com.mysql.jdbc.Driver.connectionUrl: String koneksi untuk terhubung ke sumber JDBC.-
statement: Pernyataan SQLINSERT INTOyang digunakan untuk menulis ke database. Pernyataan harus menentukan kolom tabel yang akan ditulis, dengan karakter placeholder'?'untuk klausaVALUES. Pipeline mengganti placeholder dengan nilai kolom yang sesuai dari data JSON.Contoh:
INSERT INTO tableName (column1, column2) VALUES (?,?).
Parameter berikut bersifat opsional:
username: Nama pengguna untuk koneksi JDBC.password: Sandi untuk koneksi JDBC.connectionProperties: String properti untuk koneksi JDBC. Contoh:unicode=true;characterEncoding=UTF-8.
Topik Pub/Sub
Untuk menulis ke topik Pub/Sub, tetapkan parameter template berikut:
sinkType:PUBSUB.topic: Topik Pub/Sub yang akan ditulis.
Berikan peran Pub/Sub Publisher
(roles/pubsub.publisher) kepada
akun layanan pekerja.
Tabel Spanner
Untuk menulis ke tabel Spanner, tetapkan parameter template berikut:
sinkType:SPANNER.projectId: ID project yang berisi tabel Spanner.spannerInstanceName: Nama instance Spanner.spannerDatabaseName: Nama database Spanner.spannerTableName: Nama tabel Spanner.
Parameter berikut bersifat opsional:
maxNumMutations: Jumlah maksimum sel yang diubah per batch.maxNumRows: Jumlah maksimum baris yang diubah per batch.batchSizeBytes: Jumlah maksimum byte yang diubah per batch.commitDeadlineSeconds: Batas waktu untuk panggilan API commit, dalam detik.
Berikan peran Cloud Spanner Database User
(roles/spanner.databaseUser) kepada
akun layanan pekerja.
Parameter template
Parameter yang diperlukan
- qps: Menunjukkan kecepatan pesan per detik yang akan dipublikasikan ke Pub/Sub.
Parameter opsional
- schemaTemplate: Template skema yang sudah ada untuk digunakan. Nilainya harus berupa salah satu dari: [GAME_EVENT].
- schemaLocation: Jalur Cloud Storage lokasi skema. Contoh,
gs://<bucket-name>/prefix. - topic: Nama topik tempat pipeline harus memublikasikan data. Contoh,
projects/<project-id>/topics/<topic-name>. - messagesLimit: Menunjukkan jumlah maksimum pesan output yang akan dibuat. 0 berarti tidak terbatas. Nilai defaultnya adalah: 0.
- outputType: Jenis Output pesan. Default-nya adalah JSON.
- avroSchemaLocation: Jalur Cloud Storage lokasi skema Avro. Wajib diisi jika jenis output adalah AVRO atau PARQUET. Contoh,
gs://your-bucket/your-path/schema.avsc. - sinkType: Jenis Sink pesan. Defaultnya adalah PUBSUB.
- outputTableSpec: Tabel BigQuery output. Wajib diisi jika sinkType adalah BIGQUERY. Misalnya,
<project>:<dataset>.<table_name>. - writeDisposition: BigQuery WriteDisposition. Misalnya, WRITE_APPEND, WRITE_EMPTY, atau WRITE_TRUNCATE. Nilai defaultnya: WRITE_APPEND.
- outputDeadletterTable: Pesan yang gagal mencapai tabel output karena berbagai alasan (misalnya, skema tidak cocok, JSON salah bentuk) akan ditulis ke tabel ini. Jika tidak ada, tabel akan dibuat selama eksekusi pipeline. Contoh,
your-project-id:your-dataset.your-table-name. - windowDuration: Durasi/ukuran jendela tempat data akan ditulis ke Cloud Storage. Format yang diizinkan adalah: Ns (untuk detik, contoh: 5s), Nm (untuk menit, contoh: 12m), Nh (untuk jam, contoh: 2h). Contoh,
1m. Nilai defaultnya adalah: 1 menit. - outputDirectory: Jalur dan awalan nama file untuk menulis file output. Harus diakhiri dengan garis miring. Pemformatan DateTime digunakan untuk mengurai jalur direktori untuk pemformat tanggal & waktu. Contoh,
gs://your-bucket/your-path/. - outputFilenamePrefix: Awalan yang akan ditempatkan pada setiap file berwindow. Contoh,
output-. Nilai defaultnya: output-. - numShards: Jumlah maksimum shard output yang dihasilkan saat menulis. Jumlah shard yang lebih tinggi berarti throughput yang lebih tinggi untuk menulis ke Cloud Storage, tetapi berpotensi menimbulkan biaya agregasi data yang lebih tinggi di seluruh shard saat memproses file Cloud Storage output. Nilai default ditentukan oleh Dataflow.
- driverClassName: Nama class driver JDBC yang akan digunakan. Contoh,
com.mysql.jdbc.Driver. - connectionUrl: String koneksi URL untuk terhubung ke sumber JDBC. Contoh,
jdbc:mysql://some-host:3306/sampledb. - username: Nama pengguna yang akan digunakan untuk koneksi JDBC.
- password: Sandi yang akan digunakan untuk koneksi JDBC.
- connectionProperties: String properti yang akan digunakan untuk koneksi JDBC. Format string harus [propertyName=property;]*. Contoh,
unicode=true;characterEncoding=UTF-8. - statement: Pernyataan SQL yang akan dieksekusi untuk menulis ke database. Pernyataan harus menentukan nama kolom tabel dalam urutan apa pun. Hanya nilai nama kolom yang ditentukan yang akan dibaca dari JSON dan ditambahkan ke pernyataan. Contoh,
INSERT INTO tableName (column1, column2) VALUES (?,?). - projectId: Project ID GCP tempat tabel Spanner berada.
- spannerInstanceName: Nama instance Cloud Spanner.
- spannerDatabaseName: Nama database Cloud Spanner.
- spannerTableName: Nama tabel Cloud Spanner.
- maxNumMutations: Menentukan batas mutasi sel (jumlah maksimum sel yang dimutasi per batch). Nilai defaultnya adalah 5000.
- maxNumRows: Menentukan batas mutasi baris (jumlah maksimum baris yang dimutasi per batch). Nilai defaultnya adalah 1.000.
- batchSizeBytes: Menentukan batas ukuran batch (jumlah maksimum byte yang diubah per batch). Nilai defaultnya adalah 1 MB.
- commitDeadlineSeconds: Menentukan batas waktu dalam hitungan detik untuk panggilan API Commit.
- bootstrapServer: Server Bootstrap Kafka. Misalnya,
localhost:9092. - kafkaTopic: Topik Kafka yang akan ditulisi. Contoh,
topic.
Menjalankan template
Konsol
- Buka halaman Dataflow Membuat tugas dari template. Buka Membuat tugas dari template
- Di kolom Nama tugas, masukkan nama tugas yang unik.
- Opsional: Untuk Endpoint regional, pilih nilai dari menu drop-down. Region default-nya adalah
us-central1.Untuk mengetahui daftar region tempat Anda dapat menjalankan tugas Dataflow, lihat Lokasi Dataflow.
- Dari menu drop-down Template Dataflow, pilih template Streaming Data Generator.
- Di kolom parameter yang disediakan, masukkan nilai parameter Anda.
- Klik Jalankan tugas.
gcloud
Di shell atau terminal Anda, jalankan template:
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
Ganti kode berikut:
PROJECT_ID: Google Cloud Project ID tempat Anda ingin menjalankan tugas DataflowREGION_NAME: region tempat Anda ingin men-deploy tugas Dataflow—misalnya,us-central1JOB_NAME: nama tugas unik pilihan AndaVERSION: versi template yang ingin Anda gunakanAnda dapat menggunakan nilai berikut:
latestuntuk menggunakan template versi terbaru, yang tersedia di folder induk tanpa tanggal di bucket—gs://dataflow-templates-REGION_NAME/latest/- nama versi, seperti
2023-09-12-00_RC00, untuk menggunakan template versi tertentu, yang dapat ditemukan bertingkat di masing-masing folder induk yang diberi tanggal dalam bucket—gs://dataflow-templates-REGION_NAME/
SCHEMA_LOCATION: jalur ke file skema di Cloud Storage. Contoh:gs://mybucket/filename.json.QPS: jumlah pesan yang akan dipublikasikan per detikPUBSUB_TOPIC: topik Pub/Sub output. Contoh:projects/my-project-id/topics/my-topic-id.
API
Untuk menjalankan template menggunakan REST API, kirim permintaan POST HTTP. Untuk mengetahui informasi selengkapnya tentang API dan cakupan otorisasinya, lihat 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/", } }
Ganti kode berikut:
PROJECT_ID: Google Cloud Project ID tempat Anda ingin menjalankan tugas DataflowLOCATION: region tempat Anda ingin men-deploy tugas Dataflow—misalnya,us-central1JOB_NAME: nama tugas unik pilihan AndaVERSION: versi template yang ingin Anda gunakanAnda dapat menggunakan nilai berikut:
latestuntuk menggunakan template versi terbaru, yang tersedia di folder induk tanpa tanggal di bucket—gs://dataflow-templates-REGION_NAME/latest/- nama versi, seperti
2023-09-12-00_RC00, untuk menggunakan template versi tertentu, yang dapat ditemukan bertingkat di masing-masing folder induk yang diberi tanggal dalam bucket—gs://dataflow-templates-REGION_NAME/
SCHEMA_LOCATION: jalur ke file skema di Cloud Storage. Contoh:gs://mybucket/filename.json.QPS: jumlah pesan yang akan dipublikasikan per detikPUBSUB_TOPIC: topik Pub/Sub output. Contoh:projects/my-project-id/topics/my-topic-id.
Langkah berikutnya
- Pelajari template Dataflow.
- Lihat daftar template yang disediakan Google.