Data JDBC Spark Bersamaan Dibaca
Pernahkah Anda melalui proses penerapan Spark dalam proyek Anda, menentukan jumlah partisi pengacakan yang optimal, alokasi memori untuk instance driver dan pelaksana, jumlah inti pelaksana dan semua hal menyenangkan itu hanya untuk membaca data dari sumber JDBC seperti ini ?
jdbcDF = spark.read \
.format("jdbc") \
.option("driver", "org.postgresql.Driver") \
.option("url", "jdbc:postgresql:dbserver") \
.option("user", os.environ['user']) \
.option("password", os.environ['pass']) \
.option("query", query) \
.load()
Apa masalah dengan kode di atas? Yah…ini sangat lambat :) Jika Anda menjalankan kueri SQL di Spark seperti pada contoh di atas, Anda hanya akan menggunakan satu utas (Anda dapat melihat di Spark UI bahwa hanya satu tugas yang dijalankan saat mengeksekusi kode) dan kerangka data yang dibuat akan memiliki satu partisi.
Sebagian besar dari kita menggunakan Spark di pipa ETL/ELT untuk pemrosesan paralel, jadi membaca data dari sumber JDBC hanya menggunakan satu utas mungkin bukan yang kita inginkan. Belum lagi kita harus mempartisi ulang kerangka data setelah menjalankan kueri karena data sekarang disimpan dalam 1 partisi. Jadi bagaimana sebenarnya kita membaca data secara bersamaan di Spark?
# pass the original query to dbtable as a subquery
final_query = f'({query}) as q'
# depends on your usecase, how many executor cores you have in the cluster
# what do you plan to do with the dataframe downstream, data size etc
partitionsCount = 100
# I'm hardcoding the date values here, but I'll show below how to identify
# these values at run time if you don't know them upfront
min_date = '2020-01-01'
max_date = '2022-12-31'
jdbcDF = spark.read \
.format("jdbc") \
.option("driver", "org.postgresql.Driver") \
.option("url", "jdbc:postgresql:dbserver") \
.option("user", os.environ['user']) \
.option("password", os.environ['pass']) \
.option("dbtable", final_query) \
.option("numPartitions", partitionsCount) \
.option("partitionColumn", "date") \
.option("lowerBound", f"{min_date}") \
.option("upperBound", f"{max_date}") \
.load()
- dbtable:
Berdasarkan dokumentasi resmi dari Spark, jika kita ingin membaca data secara bersamaan dan mempartisi berdasarkan kebutuhan kita, kita tidak dapat lagi menggunakan opsi kueri dan sebagai gantinya harus meneruskan kueri SQL kita melalui opsi dbtable . Satu-satunya perbedaan antara kedua opsi ini (selain nama) adalah kueri sekarang harus diapit dalam tanda kurung dan memiliki alias — pada dasarnya meneruskan kueri asli sebagai subkueri. - numPartitions:
Jumlah maksimum partisi yang dapat digunakan untuk paralelisme dalam pembacaan dan penulisan tabel. Ini juga menentukan jumlah maksimum koneksi JDBC bersamaan. Jika jumlah partisi yang akan ditulis melebihi batas ini, kami menguranginya hingga batas ini dengan memanggil penggabungan(numPartitions) sebelum menulis. ( sumber ) Saya memiliki set saya ke 100 di cluster EMR, tetapi praktik yang baik adalah mengatur nomor partisi antara 1 dan 4 kali jumlah inti. Ini bukan angka yang ditetapkan, jadi selalu uji dengan nilai yang berbeda untuk melihat apa yang cocok untuk Anda. - partitionColumn:
Ini mungkin parameter yang paling penting dan salah satu yang akan membuat perbedaan terbesar dalam waktu proses kueri Anda. Beberapa hal yang perlu diingat:
- kolom yang dimaksud harus numerik, tanggal, atau stempel waktu
- untuk performa terbaik, nilai kolom harus terdistribusi secara merata dan memiliki kardinalitas tinggi. Kami ingin menghindari kolom miring data — saya akan memperluas ini di bawah.
- jika Anda memiliki banyak kolom dengan properti di atas, pilih kolom yang diindeks. - lowerbound, upperbound :
Perhatikan bahwa lowerBound dan upperBound hanya digunakan untuk menentukan langkah partisi, bukan untuk memfilter baris dalam tabel. Jadi semua baris dalam tabel akan dipartisi dan dikembalikan. ( sumber ) Di bawah ini Anda akan melihat bagaimana Spark menghasilkan setiap partisi. Inilah sebabnya mengapa Anda tidak menginginkan kolom kardinalitas rendah sebagai kolom partisi Anda, tetapi saya akan memberikan contoh di bawah jika tidak masuk akal sekarang.
Bayangkan kolom yang ingin Anda gunakan untuk mempartisi hanya memiliki nilai 0 dan 1 — memang contoh ekstrem, tetapi skenario umum dalam Rekayasa Data di mana nilai boolean Benar dan Salah diubah menjadi nilai numerik.
Karena Anda hanya memiliki 0 dan 1, Anda menggunakan 0 sebagai batas bawah dan 1 sebagai batas atas. Berdasarkan bagaimana partisi dihasilkan — lihat gambar di atas — semua baris dengan nilai 0 akan didorong ke partisi pertama (klausa WHERE dari partisi pertama adalah satu-satunya klausa di mana 0 tidak disaring) dan semua baris dengan nilai 1 akan didorong ke partisi terakhir ( klausa WHERE dari partisi terakhir adalah satu-satunya klausa di mana 1 tidak disaring) .
Jadi, terlepas dari berapa banyak partisi yang Anda buat, hanya dua partisi yang akan diisi dan hanya dua inti eksekutor yang akan melakukan semua pekerjaan, yang lainnya akan dalam keadaan diam jika tidak ada pekerjaan lain yang dijadwalkan untuk mereka. Karena itu, kinerja kueri Anda akan terpukul.
Kemiringan Data
Bayangkan kita memiliki kumpulan data dengan 20 juta baris dan 30 partisi, batas bawah dan atas adalah 2020–01–01 dan 2022–12–31. Untuk tujuan artikel ini, saya akan menggunakan Spark Session dengan 30 core eksekutor (artinya kita dapat menjalankan hingga 30 tugas sekaligus). Perlu diingat bahwa Spark menetapkan satu tugas per partisi, sehingga setiap partisi akan diproses oleh satu inti pelaksana.
Data senilai setiap tahun akan dibagi dalam 10 partisi dan mengingat bahwa 1 tugas akan ditetapkan ke 1 partisi, setiap tugas yang ditetapkan ke tahun 2020 dan 2021 akan bertanggung jawab untuk memproses 100 ribu baris, tetapi setiap tugas yang ditetapkan ke tahun 2022 harus diproses 1,8 juta baris. Itu adalah peningkatan 1700% data yang harus diproses oleh satu tugas.
Hasilnya adalah 20 tugas pertama akan menyelesaikan pemrosesan partisi yang ditetapkan dalam waktu singkat dan setelah itu 20 inti pelaksana yang ditugaskan untuk tugas yang diselesaikan akan menganggur (jika tidak ada pekerjaan lain di mana inti pelaksana yang menganggur dapat digunakan) hingga yang terakhir 10 tugas selesai memproses 1700% lebih banyak catatan. Ini adalah pemborosan sumber daya dan peningkatan biaya pengoperasian, terutama jika Anda menggunakan Spark di cluster EMR, Glue Job, atau Databricks.
Bagaimana kita memperbaikinya? Solusi paling sederhana: gunakan kolom yang berbeda untuk mempartisi. Namun jika sangat penting untuk menggunakan kolom ini, kita dapat menggunakan dua kueri dan dua kerangka data, bukan satu. Anda masih akan mengalami masalah dengan data miring ke hilir (masalah berulang dalam sistem komputasi paralel), tetapi kami menangani satu masalah dalam satu waktu.
Kueri pertama akan memuat dua tahun pertama (30 inti eksekutor akan dibagi antara 2020 dan 2021), lalu kueri kedua akan memuat tahun terakhir (sehingga semua 30 inti eksekutor sekarang akan digunakan untuk 2022). Karena itu, tugas yang diberikan terkait dengan 2022 sekarang akan bertanggung jawab untuk memproses 600 ribu baris, bukan 1,8 juta baris. Itu adalah penurunan 66% pada baris yang harus diproses oleh satu inti eksekutor versus implementasi asli dan di atas itu, kami sekarang menggunakan semua inti eksekutor dari Spark Session kami alih-alih menganggur.
Bagaimana jika saya tidak memiliki kolom numerik, tanggal, atau stempel waktu, atau kolom memiliki data miring?
Bukan masalah. Sebagai Insinyur Data, menggunakan CTE dan Fungsi Window untuk menghasilkan kolom numerik row_number()seharusnya sangat mudah. Beberapa overhead akan diperkenalkan saat kolom baru dibuat, tetapi bahkan dengan overhead ini, kueri terakhir akan dieksekusi jauh lebih cepat dengan menggunakan kolom yang baru dibuat daripada mencoba memuat semua data dalam satu partisi menggunakan satu inti pelaksana.
with cte as
(
SELECT f1, f2, f3
FROM table
)
select *, row_number() over (order by f1) as rn from cte
original_query = "SELECT f1, f2, f3 FROM table"
final_query = f'(with cte as ({original_query}) select *, row_number() over (order by f1) as rn from cte) as q'
Kami mengeksekusi dua kueri alih-alih satu:
- kueri pertama akan dieksekusi untuk menemukan nilai min dan maks dari kolom yang akan kami gunakan sebagai kolom partisi
- dalam kueri kedua kami akan meneruskan hasil dari kueri pertama sebagai argumen batas bawah dan batas atas .
query_min_max = f"SELECT min(date), max(date) from ({query}) q"
final_query = f"({query}) as q"
df_min_max = spark.read \
.format("jdbc") \
.option("driver", "org.postgresql.Driver") \
.option("url", "jdbc:postgresql:dbserver") \
.option("user", os.environ['user']) \
.option("password", os.environ['pass']) \
.option("query", query_min_max) \
.load()
min = df_min_max.first()["min"]
max = df_min_max.first()["max"]
jdbcDF = spark.read \
.format("jdbc") \
.option("driver", "org.postgresql.Driver") \
.option("url", "jdbc:postgresql:dbserver") \
.option("user", os.environ['user']) \
.option("password", os.environ['pass']) \
.option("dbtable", final_query) \
.option("numPartitions", partitionsCount) \
.option("partitionColumn", "date") \
.option("lowerBound", f"{min}") \
.option("upperBound", f"{max}") \
.load()
- Jika Anda menjalani proses penerapan Spark di aplikasi Anda, coba gunakan sepenuhnya dan baca data secara bersamaan. Menjalankan kueri SQL ke spark.read() seperti yang ditunjukkan dalam dokumentasi JDBC Spark akan mendorong semua data dalam satu partisi dan hanya menggunakan satu inti pelaksana, terlepas dari berapa banyak inti yang telah Anda siapkan di Sesi Spark Anda.
- Jika Anda menggunakan Spark hari ini untuk membaca data dari sumber data JDBC menggunakan pendekatan utas tunggal, cobalah untuk memfaktorkan ulang kueri Anda dan meningkatkan kinerjanya, Anda bahkan mungkin mendapatkan promosi saat melakukannya :)
- Pada akhirnya, waktu adalah uang. Jika Anda mengkhawatirkan upaya menulis ulang kueri Anda, saya dapat berbicara dari pengalaman — beberapa kueri Spark yang telah saya refactored menggunakan teknik yang dijelaskan di atas sekarang berjalan > 90% lebih cepat dari sebelumnya. Anda dapat membayangkan betapa positifnya peningkatan ini di tim saya.
Jibril

![Apa itu Linked List? [Bagian 1]](https://post.nghiatu.com/assets/images/m/max/724/1*Xokk6XOjWyIGCBujkJsCzQ.jpeg)



































