Penyeimbangan ulang tugas dinamis

Fitur Penyeimbangan Ulang Kerja Dinamis pada layanan Dataflow memungkinkan layanan mempartisi ulang pekerjaan secara dinamis berdasarkan kondisi runtime. Kondisi ini dapat mencakup hal berikut:

  • Ketidakseimbangan dalam penugasan kerja
  • Pekerja membutuhkan waktu lebih lama dari yang diperkirakan untuk menyelesaikan tugas
  • Pekerja menyelesaikan tugas lebih cepat dari yang diperkirakan

Layanan Dataflow secara otomatis mendeteksi kondisi ini dan dapat secara dinamis menetapkan tugas ke worker yang tidak digunakan atau kurang digunakan untuk mengurangi waktu pemrosesan keseluruhan tugas Anda. Untuk mengetahui informasi selengkapnya tentang cara Dataflow mendistribusikan dan menskalakan eksekusi, lihat Memahami paralelisme di Dataflow.

Batasan

Penyeimbangan ulang kerja dinamis hanya terjadi saat layanan Dataflow memproses beberapa data input secara paralel: saat membaca data dari sumber input eksternal, saat bekerja dengan PCollection perantara yang diwujudkan, atau saat bekerja dengan hasil agregasi seperti GroupByKey. Jika sejumlah besar langkah dalam tugas Anda digabungkan (fused), tugas Anda memiliki lebih sedikit PCollection perantara, dan penyeimbangan ulang tugas dinamis terbatas pada jumlah elemen dalam PCollection yang diwujudkan sumber. Jika Anda ingin memastikan bahwa penyeimbangan ulang tugas dinamis dapat diterapkan ke PCollection tertentu dalam pipeline, Anda dapat mencegah penggabungan dengan beberapa cara berbeda untuk memastikan paralelisme dinamis.

Penyeimbangan ulang tugas dinamis tidak dapat memparalelkan ulang data yang lebih halus dari satu rekaman. Jika data Anda berisi setiap rekaman yang menyebabkan penundaan besar dalam waktu pemrosesan, rekaman tersebut mungkin masih menunda tugas Anda. Dataflow tidak dapat membagi dan mendistribusikan ulang satu rekaman "panas" ke beberapa pekerja.

Java

Jika Anda menetapkan jumlah shard tetap untuk output akhir pipeline (misalnya, dengan menulis data menggunakan TextIO.Write.withNumShards), Dataflow akan membatasi paralelisme berdasarkan jumlah shard yang Anda pilih.

Python

Jika Anda menetapkan jumlah shard tetap untuk output akhir pipeline (misalnya, dengan menulis data menggunakan beam.io.WriteToText(..., num_shards=...)), Dataflow akan membatasi paralelisme berdasarkan jumlah shard yang Anda pilih.

Go

Jika Anda menetapkan jumlah sharding tetap untuk output akhir pipeline, Dataflow membatasi paralelisme berdasarkan jumlah sharding yang Anda pilih.

Bekerja dengan Sumber Data Kustom

Java

Jika pipeline Anda menggunakan sumber data kustom yang Anda berikan, Anda harus menerapkan metode splitAtFraction agar sumber Anda dapat berfungsi dengan fitur penyeimbangan ulang kerja dinamis.

Jika Anda menerapkan splitAtFraction dengan tidak benar, data dari sumber Anda mungkin tampak diduplikasi atau dihapus. Lihat informasi referensi API tentang RangeTracker untuk mendapatkan bantuan dan tips tentang menerapkan splitAtFraction.

Python

Jika pipeline Anda menggunakan sumber data kustom yang Anda berikan, RangeTracker Anda harus menerapkan try_claim, try_split, position_at_fraction, dan fraction_consumed agar sumber Anda dapat berfungsi dengan fitur penyeimbangan ulang tugas dinamis.

Lihat informasi referensi API tentang RangeTracker untuk mengetahui informasi selengkapnya.

Go

Jika pipeline Anda menggunakan sumber data kustom yang Anda berikan, Anda harus menerapkan RTracker yang valid agar sumber Anda dapat berfungsi dengan fitur penyeimbangan ulang kerja dinamis.

Untuk mengetahui informasi selengkapnya, lihat informasi referensi API RTracker.

Penyeimbangan ulang tugas dinamis menggunakan nilai yang ditampilkan dari metode getProgress() sumber kustom Anda untuk diaktifkan. Implementasi default untuk getProgress() menampilkan null. Untuk memastikan penskalaan otomatis diaktifkan, pastikan penggantian sumber kustom getProgress() Anda menampilkan nilai yang sesuai.