Memahami paralelisme di Dataflow

Dataflow dirancang untuk menjalankan pipeline pemrosesan data besar dengan mendistribusikan pekerjaan di seluruh kumpulan instance komputasi terkelola. Memahami cara Dataflow memparalelkan pemrosesan akan membantu Anda mendesain pipeline yang efisien, menghindari hambatan performa, dan mengoptimalkan biaya resource.

Halaman ini menjelaskan cara Dataflow memparalelkan pemrosesan data, cara mengelola dan menskalakan eksekusi, faktor umum yang membatasi paralelisme, dan teknik yang dapat Anda gunakan untuk mengoptimalkan throughput pipeline.

Model paralelisme: Horizontal versus vertikal

Dataflow mencapai paralelisme menggunakan dua strategi komplementer:

  • Paralelisme horizontal: Data pipeline dipartisi dan diproses di beberapa instance pekerja (mesin virtual) secara bersamaan. Dataflow dapat menyesuaikan ukuran kumpulan pekerja secara otomatis berdasarkan permintaan workload melalui Penskalaan Otomatis Horizontal. Secara default, Dataflow menetapkan batas resource 4.000 worker per tugas, yang dapat disesuaikan menggunakan permintaan kuota.

  • Paralelisme vertikal: Beberapa thread dan core CPU dalam satu instance pekerja memproses data pipeline secara bersamaan. Setiap worker VM menjalankan proses worker dan thread harness untuk menggunakan resource komputasi yang tersedia. Dengan penskalaan thread dinamis, Dataflow dapat menyesuaikan jumlah thread aktif per pekerja dalam pipeline batch berdasarkan pemakaian CPU dan ruang kosong memori. Di Dataflow Prime, Penskalaan Otomatis Vertikal secara dinamis menskalakan memori dan komputasi yang dialokasikan ke worker.

Unit kerja dan hierarki eksekusi

Untuk mendistribusikan pemrosesan di seluruh pekerja dan thread, Dataflow memecah pipeline Apache Beam menjadi unit kerja diskrit:

  • PCollection dan partisi: PCollection merepresentasikan set data terdistribusi. Untuk data terbatas (pipeline batch), Dataflow membagi set data menjadi beberapa pemisahan atau shard. Untuk data yang tidak terbatas (pipeline streaming), data tiba secara terus-menerus dan diserap sebagai pesan atau partisi streaming.
  • Bundle: Dataflow mengelompokkan elemen ke dalam bundle arbitrer untuk diproses oleh DoFn. Bundle adalah unit kegagalan dan percobaan ulang: jika pemrosesan elemen memunculkan pengecualian yang tidak tertangani, seluruh bundle akan dicoba lagi. Operasi dengan konsumsi memori yang tinggi dapat meningkatkan tekanan memori pekerja dan menyebabkan error kehabisan memori.
  • Penggabungan tahap dan langkah: Selama pengoptimalan grafik, Dataflow menggabungkan transformasi yang berdekatan ke dalam tahap eksekusi gabungan untuk menghilangkan overhead materialisasi data perantara. Dalam tahap gabungan, elemen diproses dalam loop eksekusi yang ketat pada satu thread sebelum diteruskan ke tahap berikutnya atau batas pengacakan.

Untuk mengetahui detail selengkapnya tentang terjemahan pipeline dan pembuatan grafik, lihat Siklus proses pipeline.

Paralelisme dan penskalaan otomatis terkelola

Secara default, Dataflow mengelola paralelisme pipeline secara otomatis tanpa memerlukan penyesuaian partisi manual dengan cara berikut:

  • Penskalaan Otomatis Horizontal:
    • Pipeline batch: Mengevaluasi total perkiraan pekerjaan yang tersisa, backlog sumber, dan penggunaan CPU untuk menskalakan kumpulan worker ke atas atau ke bawah guna menyelesaikan pekerjaan dengan cepat dan hemat biaya.
    • Pipeline streaming: Menganalisis latensi sistem, ukuran backlog, dan penggunaan CPU untuk meningkatkan skala pekerja selama lonjakan throughput dan memperkecil skala selama periode traffic rendah. Untuk mengetahui detailnya, lihat Menyesuaikan Penskalaan Otomatis Horizontal Streaming.
  • Penyeimbangan Ulang Kerja Dinamis (DWR): Dalam pipeline batch, Dataflow memantau progres tugas worker individual. Jika worker selesai lebih awal atau worker lain tertinggal karena kemiringan data (straggler), Dataflow akan membagi secara dinamis sisa tugas yang belum diproses dari worker yang lambat dan menetapkannya kembali ke worker yang tidak aktif. Untuk mengetahui informasi selengkapnya, lihat Penyeimbangan ulang tugas dinamis.
  • Penskalaan thread dinamis: Dalam pipeline batch yang menggunakan Portable Runner, secara otomatis menyesuaikan jumlah thread pemrosesan serentak per pekerja berdasarkan pemanfaatan CPU dan ruang memori. Untuk mengetahui informasi selengkapnya, lihat Penskalaan thread dinamis.
  • Penskalaan Otomatis Vertikal: Di Dataflow Prime, Dataflow menskalakan memori pekerja dan resource komputasi secara dinamis untuk mencegah error kehabisan memori dan mengoptimalkan penggunaan resource. Untuk mengetahui informasi selengkapnya, lihat Penskalaan Otomatis Vertikal.

Faktor yang membatasi paralelisme

Pipeline mungkin tidak mencapai paralelisme yang diharapkan karena karakteristik data atau desain grafik pipeline berikut:

Sumber input yang tidak dapat dibagi

Jika sumber input tidak dapat dibagi menjadi rentang independen, Dataflow akan dipaksa untuk membaca sumber secara berurutan dengan satu thread pekerja:

  • Kompresi file yang tidak dapat dibagi: Format seperti .gz (gzip) atau .bzip2 (tanpa pengindeksan) tidak dapat dibaca secara paralel dari offset byte arbitrer. Membaca satu file besar yang dikompresi membatasi tahap penyerapan ke satu thread hingga data didekompresi dan didistribusikan ulang.
  • Penyelesaian: Simpan data dalam format file yang dapat dibagi (seperti Parquet, Avro, atau format yang dikompresi Snappy) atau bagi data input menjadi beberapa file yang lebih kecil di Cloud Storage.

Penggabungan langkah dan fan-out tinggi

Penggabungan langkah meningkatkan performa dengan mengurangi overhead serialisasi, tetapi dapat secara tidak sengaja membatasi paralelisme dan meningkatkan tekanan memori saat langkah dengan paralelisme rendah menghasilkan sejumlah besar elemen output (operasi "fan-out tinggi"):

  • Contoh: Sumber membaca lima file dan digabungkan dengan transformasi FlatMap yang menghasilkan 1.000.000 elemen output. Jika transformasi FlatMap digabungkan dengan transformasi hilir, semua 1.000.000 elemen akan terus dieksekusi di paling banyak lima thread pekerja, sehingga sangat membatasi throughput hilir. Selain itu, jika transformasi perantara meluas secara signifikan dalam memori sebelum di-commit, paket besar dapat menghabiskan memori pekerja yang tersedia.
  • Resolusi: Sisipkan transformasi Redistribute (atau Reshuffle klasik) antara langkah fan-out tinggi dan transformasi hilir untuk menghentikan penggabungan dan mendistribusikan ulang pekerjaan di seluruh kumpulan pekerja. Untuk men-debug masalah terkait memori, lihat Memecahkan masalah error kehabisan memori.

Penyimpangan kunci dan kunci aktif

Operasi agregasi (GroupByKey, CoGroupByKey, Combine.PerKey) mengelompokkan elemen berdasarkan kunci terkaitnya.

  • Bottleneck hot key: Dataflow merutekan semua elemen dengan kunci yang sama ke satu thread pekerja untuk agregasi. Jika satu kunci berisi persentase besar dari total set data, pekerja tersebut akan menjadi lambat, dan pekerja di upstream mungkin mengalami tekanan balik. Misalnya, kunci null default atau kunci kategori yang sangat populer.
  • Resolusi:
    1. Gunakan Combiner (CombineFn atau Combine.PerKey) dan bukan GroupByKey jika memungkinkan, sehingga Dataflow dapat melakukan kombinasi lokal parsial sebelum pengacakan.
    2. Tambahkan awalan atau akhiran bilangan bulat acak ke tombol cepat (key salting) untuk mendistribusikan ruang kunci di seluruh pekerja, diikuti dengan agregasi tahap kedua untuk menggabungkan hasil yang di-salt.

Throttling sink downstream

Saat menulis output pipeline ke layanan eksternal seperti database atau API pihak ketiga, paralelisme tinggi dapat membebani sistem tujuan:

  • Pembatasan: Ratusan thread pekerja yang mengeluarkan panggilan tulis serentak dapat menyebabkan error batas kecepatan, waktu tunggu koneksi habis, atau penurunan kualitas database.
  • Resolusi:
    • Batasi paralelisme penulisan dengan mengelompokkan elemen dengan GroupByKey atau menggunakan sink batch dengan paralelisme yang terkontrol.
    • Terapkan backoff eksponensial dan logika percobaan ulang sisi klien dalam penerapan DoFn sink.

Strategi pengoptimalan

Untuk mengoptimalkan paralelisme dalam tugas Dataflow, pertimbangkan pendekatan berikut:

Langkah berikutnya