Halaman ini memberikan ringkasan siklus proses pipeline dari kode pipeline hingga tugas Dataflow.
Halaman ini menjelaskan konsep berikut:
- Apa itu grafik eksekusi, dan bagaimana pipeline Apache Beam menjadi tugas Dataflow
- Cara Dataflow menangani error
- Cara Dataflow secara otomatis memparalelkan dan mendistribusikan logika pemrosesan di pipeline Anda ke worker yang menjalankan tugas Anda
- Pengoptimalan tugas yang mungkin dilakukan Dataflow
Grafik eksekusi
Saat Anda menjalankan pipeline Dataflow, Dataflow akan membuat grafik eksekusi dari kode yang membangun objek Pipeline, termasuk semua transformasi dan fungsi pemrosesan terkaitnya, seperti objek DoFn. Ini adalah grafik eksekusi pipeline, dan fase ini disebut
waktu pembuatan grafik.
Selama konstruksi grafik, Apache Beam akan menjalankan kode secara lokal dari
titik entri utama kode pipeline, berhenti pada panggilan ke langkah sumber, sink,
atau transformasi, dan mengubah panggilan ini menjadi node grafik.
Oleh karena itu, potongan kode di titik entri pipeline (metode Java dan Go main
atau tingkat teratas skrip Python) dieksekusi secara lokal di mesin yang
menjalankan pipeline. Kode yang sama yang dideklarasikan dalam metode objek DoFn
dieksekusi di pekerja Dataflow.
Misalnya, sampel WordCount yang disertakan dengan Apache Beam SDK berisi serangkaian transformasi untuk membaca, mengekstrak, menghitung, memformat, dan menulis setiap kata dalam kumpulan teks, beserta jumlah kemunculan untuk setiap kata. Diagram berikut menunjukkan cara transformasi dalam pipeline WordCount dikembangkan menjadi grafik eksekusi:

Gambar 1: Grafik eksekusi contoh WordCount
Grafik eksekusi sering kali berbeda dari urutan saat Anda menentukan transformasi saat membuat pipeline. Perbedaan ini ada karena layanan Dataflow melakukan berbagai pengoptimalan dan penggabungan pada grafik eksekusi sebelum dijalankan di resource cloud terkelola. Layanan Dataflow mematuhi dependensi data saat menjalankan pipeline Anda. Namun, langkah-langkah tanpa dependensi data di antaranya dapat dijalankan dalam urutan apa pun.
Untuk melihat grafik eksekusi yang tidak dioptimalkan yang telah dibuat Dataflow untuk pipeline Anda, pilih tugas Anda di antarmuka pemantauan Dataflow. Untuk mengetahui informasi selengkapnya tentang cara melihat tugas, lihat Menggunakan antarmuka pemantauan Dataflow.
Selama pembuatan grafik, Apache Beam memvalidasi bahwa semua resource yang dirujuk oleh pipeline, seperti bucket Cloud Storage, tabel BigQuery, dan topik atau langganan Pub/Sub, benar-benar ada dan dapat diakses. Validasi dilakukan melalui panggilan API standar ke layanan masing-masing, jadi sangat penting agar akun pengguna yang digunakan untuk menjalankan pipeline memiliki konektivitas yang tepat ke layanan yang diperlukan dan diizinkan untuk memanggil API layanan. Sebelum mengirimkan pipeline ke layanan Dataflow, Apache Beam juga memeriksa kesalahan lainnya, dan memastikan bahwa grafik pipeline tidak berisi operasi ilegal.
Grafik eksekusi kemudian diterjemahkan ke dalam format JSON, dan grafik eksekusi JSON dikirimkan ke endpoint layanan Dataflow.
Layanan Dataflow kemudian memvalidasi grafik eksekusi JSON. Setelah grafik divalidasi, grafik akan menjadi tugas di layanan Dataflow. Anda dapat melihat tugas, grafik eksekusi, status, dan informasi log menggunakan antarmuka pemantauan Dataflow.
Java
Layanan Dataflow mengirimkan respons ke mesin tempat Anda menjalankan
program Dataflow. Respons ini dienkapsulasi dalam objek DataflowPipelineJob, yang berisi jobId tugas Dataflow Anda.
Gunakan jobId untuk memantau, melacak, dan memecahkan masalah tugas Anda menggunakan
antarmuka pemantauan Dataflow
dan antarmuka command-line Dataflow.
Untuk mengetahui informasi selengkapnya, lihat
referensi API untuk DataflowPipelineJob.
Python
Layanan Dataflow mengirimkan respons ke mesin tempat Anda menjalankan
program Dataflow. Respons ini dienkapsulasi dalam objek
DataflowPipelineResult, yang berisi job_id tugas Dataflow Anda.
Gunakan job_id untuk memantau, melacak, dan memecahkan masalah tugas Anda
dengan menggunakan
Antarmuka pemantauan Dataflow
dan
Antarmuka command line Dataflow.
Go
Layanan Dataflow mengirimkan respons ke mesin tempat Anda menjalankan
program Dataflow. Respons ini dienkapsulasi dalam objek dataflowPipelineResult, yang berisi jobID tugas Dataflow Anda.
Gunakan jobID untuk memantau, melacak, dan memecahkan masalah tugas Anda
dengan menggunakan
Antarmuka pemantauan Dataflow
dan
Antarmuka command line Dataflow.
Konstruksi grafik juga terjadi saat Anda menjalankan pipeline secara lokal, tetapi grafik tidak diterjemahkan ke JSON atau dikirim ke layanan. Sebagai gantinya, grafik dijalankan secara lokal di mesin yang sama tempat Anda meluncurkan program Dataflow. Untuk mengetahui informasi selengkapnya, lihat Mengonfigurasi PipelineOptions untuk eksekusi lokal.
Penanganan error dan pengecualian
Pipeline Anda mungkin memunculkan pengecualian saat memproses data. Beberapa error ini bersifat sementara, seperti kesulitan sementara dalam mengakses layanan eksternal. Error lainnya bersifat permanen, seperti error yang disebabkan oleh data input yang rusak atau tidak dapat diuraikan, atau pointer null selama komputasi.
Dataflow memproses elemen dalam paket arbitrer, dan mencoba kembali seluruh paket saat error terjadi untuk elemen apa pun dalam paket tersebut. Saat berjalan dalam mode batch, paket yang menyertakan item yang gagal akan dicoba lagi sebanyak empat kali. Pipeline akan gagal sepenuhnya jika satu paket gagal empat kali. Saat berjalan dalam mode streaming, paket yang menyertakan item yang gagal akan dicoba lagi tanpa batas, yang dapat menyebabkan pipeline Anda terhenti secara permanen.
Saat memproses dalam mode batch, Anda mungkin melihat sejumlah besar kegagalan individu sebelum tugas pipeline gagal sepenuhnya, yang terjadi saat paket tertentu gagal setelah empat kali percobaan ulang. Misalnya, jika pipeline Anda mencoba memproses 100 paket, Dataflow dapat menghasilkan beberapa ratus kegagalan individual hingga satu paket mencapai kondisi empat kegagalan untuk keluar.
Error worker saat startup, seperti kegagalan menginstal paket di worker, bersifat sementara. Skenario ini menyebabkan percobaan ulang tanpa batas, dan dapat menyebabkan pipeline Anda terhenti secara permanen.
Paralelisasi dan distribusi
Layanan Dataflow secara otomatis memparalelkan dan mendistribusikan logika pemrosesan di pipeline Anda ke seluruh worker dan thread.
Dataflow menggunakan abstraksi dalam
model pemrograman untuk merepresentasikan
fungsi pemrosesan paralel. Misalnya, transformasi ParDo menyebabkan Dataflow mendistribusikan kode pemrosesan, yang diwakili oleh objek DoFn, ke beberapa pekerja untuk dieksekusi secara bersamaan.
Dataflow mendukung dua dimensi paralelisme yang saling melengkapi:
- Paralelisme horizontal: Data pipeline dibagi dan diproses di beberapa instance worker secara bersamaan, yang dikelola secara dinamis menggunakan Penskalaan Otomatis Horizontal.
- Paralelisme vertikal: Data pipeline diproses di beberapa core dan thread CPU pada setiap VM pekerja, yang dikelola menggunakan penskalaan thread dinamis dan Penskalaan Otomatis Vertikal.
Dataflow secara otomatis mengelola paralelisme tugas, menangani penyeimbangan ulang tugas dinamis, dan mengoptimalkan grafik eksekusi melalui pengoptimalan penggabungan dan kombinasi. Faktor-faktor seperti sumber data yang tidak dapat dibagi, langkah-langkah fan-out yang tinggi, kemiringan kunci, dan batas sink hilir dapat membatasi paralelisme pipeline.
Untuk panduan mendalam tentang cara Dataflow memartisi data, menskalakan pekerja, dan mengatasi hambatan paralelisme, lihat Memahami paralelisme di Dataflow.
Pengoptimalan fusi
Setelah bentuk JSON dari grafik eksekusi pipeline Anda divalidasi, layanan Dataflow dapat mengubah grafik untuk melakukan pengoptimalan.
Pengoptimalan dapat mencakup penggabungan beberapa langkah atau transformasi dalam grafik eksekusi pipeline menjadi satu langkah. Penggabungan langkah-langkah mencegah
layanan Dataflow perlu mewujudkan setiap PCollection perantara
dalam pipeline Anda, yang dapat menimbulkan biaya besar dalam hal memori dan
overhead pemrosesan.
Meskipun semua transformasi yang Anda tentukan dalam konstruksi pipeline dijalankan di layanan, untuk memastikan eksekusi pipeline yang paling efisien, transformasi dapat dijalankan dalam urutan yang berbeda atau sebagai bagian dari transformasi gabungan yang lebih besar. Layanan Dataflow mematuhi dependensi data antara langkah-langkah dalam grafik eksekusi, tetapi langkah-langkah lainnya dapat dieksekusi dalam urutan apa pun.
Contoh penggabungan
Diagram berikut menunjukkan cara grafik eksekusi dari contoh WordCount yang disertakan dengan Apache Beam SDK untuk Java dapat dioptimalkan dan digabungkan oleh layanan Dataflow untuk eksekusi yang efisien:

Gambar 2: Contoh Grafik Eksekusi yang Dioptimalkan WordCount
Mencegah penggabungan
Dalam beberapa kasus, Dataflow mungkin salah menebak cara optimal untuk menggabungkan operasi dalam pipeline, yang dapat membatasi kemampuan Dataflow untuk menggunakan semua worker yang tersedia. Dalam kasus seperti itu,
Anda dapat memberikan petunjuk kepada Dataflow untuk mendistribusikan ulang data, dengan menggunakan
transformasi Redistribute.
Untuk menambahkan transformasi Redistribute, panggil salah satu metode berikut:
Redistribute.arbitrarily: Menunjukkan bahwa data kemungkinan tidak seimbang. Dataflow memilih algoritma terbaik untuk mendistribusikan ulang data.Redistribute.byKey: Menunjukkan bahwaPCollectionpasangan nilai kunci kemungkinan tidak seimbang dan harus didistribusikan ulang berdasarkan kunci. Biasanya, Dataflow menempatkan semua elemen dari satu kunci yang sama pada thread pekerja yang sama. Namun, penempatan bersama kunci tidak dijamin, dan elemen diproses secara independen.
Jika pipeline Anda berisi transformasi Redistribute, Dataflow biasanya mencegah penggabungan langkah-langkah sebelum dan setelah transformasi Redistribute, serta mengacak data sehingga langkah-langkah di hilir transformasi Redistribute memiliki paralelisme yang lebih optimal.
Memantau penggabungan
Anda dapat mengakses grafik yang dioptimalkan dan tahap gabungan di konsol Google Cloud , menggunakan gcloud CLI, atau menggunakan API.
Konsol
Untuk melihat tahap dan langkah yang digabungkan dalam grafik di konsol, di tab Detail eksekusi untuk tugas Dataflow, buka tampilan grafik Alur kerja tahap.
Untuk melihat langkah-langkah komponen yang digabungkan untuk suatu tahap, klik tahap yang digabungkan dalam grafik. Di panel Info tahap, baris Langkah komponen menampilkan tahap gabungan. Terkadang, sebagian dari satu transformasi komposit digabungkan ke dalam beberapa tahap.
gcloud
Untuk mengakses grafik yang dioptimalkan dan tahap gabungan menggunakan
gcloud CLI, jalankan perintah gcloud berikut:
gcloud dataflow jobs describe --full JOB_ID --format json
Ganti JOB_ID dengan ID tugas Dataflow Anda.
Untuk mengekstrak bagian yang relevan, teruskan output perintah gcloud ke jq:
gcloud dataflow jobs describe --full JOB_ID --format json | jq '.pipelineDescription.executionPipelineStage\[\] | {"stage_id": .id, "stage_name": .name, "fused_steps": .componentTransform }'
Untuk melihat deskripsi tahap gabungan dalam file respons output, dalam array
ComponentTransform, lihat objek
ExecutionStageSummary.
API
Untuk mengakses grafik yang dioptimalkan dan tahap gabungan menggunakan API, panggil
project.locations.jobs.get.
Untuk melihat deskripsi tahap gabungan dalam file respons output, dalam array
ComponentTransform, lihat objek
ExecutionStageSummary.
Pengoptimalan gabungan
Operasi agregasi adalah konsep penting dalam pemrosesan data berskala besar.
Agregasi menggabungkan data yang secara konseptual sangat berbeda, sehingga sangat berguna untuk melakukan korelasi. Model pemrograman Dataflow merepresentasikan operasi agregasi sebagai transformasi GroupByKey, CoGroupByKey, dan
Combine.
Operasi agregasi Dataflow menggabungkan data di seluruh set data, termasuk data yang mungkin tersebar di beberapa pekerja. Selama operasi
penggabungan tersebut, sering kali lebih efisien untuk menggabungkan data sebanyak mungkin secara lokal sebelum menggabungkan data di seluruh instance. Saat Anda menerapkan
GroupByKey atau transformasi penggabungan lainnya, layanan Dataflow
secara otomatis melakukan penggabungan parsial secara lokal sebelum operasi
pengelompokan utama.
Saat melakukan penggabungan sebagian atau multi-level, layanan Dataflow membuat keputusan yang berbeda berdasarkan apakah pipeline Anda bekerja dengan data batch atau streaming. Untuk data yang dibatasi, layanan ini mengutamakan efisiensi dan akan melakukan penggabungan lokal sebanyak mungkin. Untuk data tidak terbatas, layanan lebih memilih latensi yang lebih rendah, dan mungkin tidak melakukan penggabungan sebagian, karena dapat meningkatkan latensi.