Celeborn memproses data antara untuk meningkatkan stabilitas, fleksibilitas, dan kinerja mesin data besar. Topik ini menjelaskan cara menggunakan layanan Celeborn.
Informasi latar belakang
Solusi shuffle yang ada memiliki kekurangan berikut:
-
Pada skenario dengan volume data besar, operasi tulis shuffle dapat menyebabkan spill data, yang mengakibatkan write amplification.
-
Selama proses baca shuffle, sejumlah besar paket jaringan kecil dapat menyebabkan error connection reset.
-
Proses baca shuffle menghasilkan banyak permintaan I/O kecil dan pembacaan acak (random reads), sehingga memberikan beban tinggi pada disk dan CPU.
-
Ketika jumlah mapper (M) dan reducer (N) mencapai ribuan, jumlah koneksi M × N membuat penyelesaian job hampir tidak mungkin.
-
NodeManager dan Spark External Shuffle Service berjalan dalam proses yang sama. Saat volume data shuffle sangat besar, NodeManager sering restart, yang mengganggu stabilitas penjadwalan YARN.
Celeborn mengatasi permasalahan proses shuffle tersebut dan memberikan manfaat berikut:
-
Menggunakan push-style shuffle alih-alih pull-style shuffle untuk mengurangi tekanan memori pada mapper.
-
Mendukung agregasi I/O, yang mengurangi jumlah koneksi baca shuffle dari M × N menjadi N serta mengubah pembacaan acak menjadi pembacaan berurutan (sequential reads).
-
Mendukung mekanisme dua-replika untuk mengurangi kemungkinan kegagalan fetch.
-
Mendukung arsitektur pemisahan komputasi dan penyimpanan (compute-storage separation), memungkinkan layanan shuffle ditempatkan di lingkungan perangkat keras khusus dan diisolasi dari kluster komputasi.
-
Menghilangkan ketergantungan pada disk lokal saat menjalankan Spark di Kubernetes.
Gambar berikut menunjukkan arsitektur Celeborn.
Prasyarat
Anda telah membuat kluster EMR DataLake atau kluster kustom dan memilih layanan Celeborn. Untuk informasi lebih lanjut tentang cara membuat kluster, lihat Create a cluster.
Batasan
Topik ini hanya berlaku untuk kluster dengan versi berikut.
|
Cluster |
Version |
|
DataLake cluster |
EMR-3.45.0 atau yang lebih baru, dan EMR-5.11.0 atau yang lebih baru. |
|
custom cluster |
EMR-3.45.0 atau yang lebih baru, dan EMR-5.11.0 atau yang lebih baru. |
Prosedur
Konfigurasi Spark
|
Parameter |
Description |
|
spark.shuffle.manager |
|
|
spark.serializer |
Nilainya harus org.apache.spark.serializer.KryoSerializer. |
|
spark.celeborn.push.replicate.enabled |
Menentukan apakah mekanisme dua-replika diaktifkan. Nilai yang valid:
|
|
spark.shuffle.service.enabled |
Atur parameter ini ke false untuk menggunakan Celeborn. Untuk menggunakan Celeborn, Anda harus menonaktifkan External Shuffle Service yang ada. Fitur Dynamic Allocation Spark tetap berfungsi sebagaimana mestinya ketika Celeborn diaktifkan. Catatan
|
|
spark.celeborn.shuffle.writer |
Celeborn mendukung mode writer berikut:
|
|
spark.celeborn.master.endpoints |
Tentukan titik akhir dalam format <celeborn-master-ip>:<celeborn-master-port>. Parameter:
Untuk kluster ketersediaan tinggi, konfigurasikan alamat IP semua node master. |
|
spark.sql.adaptive.enabled |
Celeborn mendukung Adaptive Query Execution (AQE). Untuk kinerja shuffle optimal, nonaktifkan local shuffle reader. Atur parameter-parameter ini masing-masing ke true, false, dan true. |
|
spark.sql.adaptive.localShuffleReader.enabled |
|
|
spark.sql.adaptive.skewJoin.enabled |
Layanan Spark mendukung konfigurasi satu klik untuk menggunakan layanan Celeborn.
-
Untuk EMR-5.11.1 dan yang lebih baru, serta EMR-3.45.1 dan yang lebih baru:
Pada halaman Status layanan Spark, di bagian Service Overview, Anda dapat mengaktifkan/menonaktifkan sakelar enableCeleborn.
-
Untuk EMR-5.11.0 dan EMR-3.45.0:
Pada halaman Status layanan Spark, di bagian Components, temukan SparkThriftServer. Di kolom Actions, pilih atau . Tindakan ini secara otomatis memodifikasi parameter konfigurasi Spark yang tercantum di atas, melakukan restart SparkThriftServer, dan memperbarui file spark-defaults.conf dan spark-thriftserver.conf.
-
Jika Anda memilih , semua job Spark akan menggunakan layanan Celeborn.
-
Jika Anda memilih , tidak ada job Spark yang menggunakan layanan Celeborn.
-
Konfigurasi Celeborn
Lihat dan modifikasi semua parameter konfigurasi Celeborn pada halaman konfigurasi layanan Celeborn.
Nilai parameter bervariasi berdasarkan kelompok node (misalnya, CORE atau TASK).
|
Parameter |
Description |
Default |
|
celeborn.worker.flusher.threads |
Jumlah thread untuk flushing data ke disk (HDD atau SSD). |
|
|
CELEBORN_WORKER_OFFHEAP_MEMORY |
Ukuran memori off-heap worker. |
Dihitung secara otomatis berdasarkan konfigurasi kluster. |
|
celeborn.application.heartbeat.timeout |
Timeout heartbeat aplikasi. Jika timeout tercapai, sistem akan melepaskan sumber daya aplikasi. |
120 s |
|
celeborn.worker.flusher.buffer.size |
Ukuran buffer flushing. Flushing dipicu ketika ukuran ini terlampaui. |
256 KB |
|
celeborn.metrics.enabled |
Menentukan apakah pemantauan diaktifkan. Nilai yang valid:
|
true |
|
CELEBORN_WORKER_MEMORY |
Ukuran memori heap worker. |
1 GB |
|
CELEBORN_MASTER_MEMORY |
Ukuran memori heap master. |
2 GB |
Restart komponen Celeborn
-
Pada halaman Status layanan Celeborn, di kolom Actions untuk komponen CelebornMaster, pilih .
CatatanUntuk kluster non-ketersediaan tinggi, Anda juga dapat mengklik Restart di kolom Actions untuk komponen CelebornMaster.
-
Pada kotak dialog, matikan sakelar Rolling Execution, masukkan alasan eksekusi, lalu klik OK.
-
Pada kotak dialog konfirmasi, klik OK.
> enableCeleborn