Tabel Hash Clustering menggunakan properti shuffle dan sort untuk mengorganisasi data. MaxCompute memanfaatkan properti ini guna mengoptimalkan rencana eksekusi, meningkatkan efisiensi, dan menghemat sumber daya. Topik ini menjelaskan cara menggunakan tabel Hash Clustering di MaxCompute.
Informasi latar belakang
Menggabungkan tabel merupakan skenario umum dalam kueri MaxCompute. Misalnya, kueri berikut melakukan inner join sederhana dengan menggabungkan tabel t1 dan tabel t2 berdasarkan kolom id.
SELECT t1.a, t2.b FROM t1 JOIN t2 ON t1.id = t2.id;MaxCompute menggunakan tiga metode utama untuk mengimplementasikan join:
Broadcast Hash Join
Ketika salah satu tabel dalam join berukuran kecil, MaxCompute menggunakan metode ini untuk menyiarkan (broadcast) tabel kecil ke semua instans task join, lalu melakukan hash join dengan tabel besar.
Shuffle Hash Join
Jika tabel yang di-join berukuran besar sehingga tidak dapat disiarkan, MaxCompute melakukan hash shuffle pada kedua tabel berdasarkan kunci gabungan (join key). Catatan dengan nilai kunci yang sama menghasilkan hash yang identik, sehingga catatan tersebut dikirim ke instans task join yang sama. Setiap instans kemudian membangun tabel hash untuk dataset yang lebih kecil dan melakukan join pembacaan sekuensial dengan dataset yang lebih besar.
Sort Merge Join
Metode Shuffle Hash Join tidak dapat digunakan jika tabel yang di-join terlalu besar karena tidak tersedia cukup memori untuk membangun tabel hash. Metode ini pertama-tama melakukan hash shuffle pada kunci gabungan, mengurutkan data berdasarkan kunci gabungan, lalu menggabungkan kedua sisi join. Gambar berikut menunjukkan proses ini.
Untuk volume dan skala data yang umum di MaxCompute, Sort Merge Join digunakan dalam sebagian besar kasus. Namun, operasi ini sangat mahal. Seperti yang ditunjukkan pada gambar, operasi shuffle memerlukan perhitungan, dan hasil antara harus ditulis ke disk. Reducer berikutnya kemudian harus membaca dan mengurutkan data tersebut. Untuk skenario dengan Mmapper danRreducer, hal ini menghasilkanM × Roperasi baca I/O. Rencana eksekusi fisik Fuxi yang sesuai ditunjukkan di bawah ini. Rencana ini memerlukan dua tahap Mapper dan satu tahap Join. Bagian yang berwarna merah menunjukkan operasi shuffle dan sort.
Selain itu, beberapa join mungkin terjadi berulang kali. Misalnya, jika kueri diubah menjadi:SELECT t1.c, t2.d FROM t1 JOIN t2 ON t1.id = t2.id;Meskipun kolom yang dipilih berbeda, operasi join-nya identik. Seluruh proses shuffle dan sort juga tetap sama.
Atau, jika kueri diubah menjadi:
SELECT t1.c, t3.d FROM t1 JOIN t3 ON t1.id = t3.id;Kueri ini menggabungkan tabel t1 dan tabel t3. Untuk tabel t1, seluruh proses shuffle dan sort tetap sama.
Oleh karena itu, jika data tabel awal disimpan menggunakan metode hash shuffle dan sort, kueri berikutnya dapat menghindari pengacakan (shuffle) dan pengurutan ulang data. Keuntungannya adalah biaya satu kali saat pembuatan tabel dapat menghemat biaya shuffle dan join berulang pada kueri selanjutnya. Rencana eksekusi fisik Fuxi untuk join kemudian berubah seperti yang ditunjukkan pada gambar berikut. Perubahan ini tidak hanya menghemat operasi shuffle dan sort, tetapi juga mengurangi jumlah tahap kueri dari tiga menjadi satu.

Catatan penggunaan
Membuat tabel Hash Clustering
Anda dapat menggunakan pernyataan berikut untuk membuat tabel Hash Clustering. Anda harus menentukan kunci kluster (cluster key), yaitu kunci hash, dan jumlah bucket hash. Pengurutan bersifat opsional. Namun, untuk kinerja optimal, Anda sebaiknya mengatur kunci pengurutan (sort key) agar sama dengan kunci kluster dalam sebagian besar kasus.
Sintaks
CREATE TABLE [IF NOT EXISTS] <table_name> [(<col_name> <data_type> [comment <col_comment>], ...)] [comment <table_comment>] [PARTITIONED BY (<col_name> <data_type> [comment <col_comment>], ...)] [CLUSTERED BY (<col_name> [, <col_name>, ...]) [SORTED BY (<col_name> [ASC | DESC] [, <col_name> [ASC | DESC] ...])] INTO <number_of_buckets> BUCKETS] [AS <select_statement>]Contoh
Tabel non-partisi
CREATE TABLE T1 (a string, b string, c bigint) CLUSTERED BY (c) SORTED by (c) INTO 1024 BUCKETS;Tabel partisi
CREATE TABLE T1 (a string, b string, c bigint) PARTITIONED BY (dt string) CLUSTERED BY (c) SORTED by (c) INTO 1024 BUCKETS;
Properti
CLUSTERED BY
Menentukan kunci hash. MaxCompute melakukan operasi hash pada kolom yang ditentukan dan mendistribusikan data ke bucket berdasarkan nilai hash. Untuk menghindari kesenjangan data (data skew), mencegah hot spot, dan mencapai eksekusi paralel yang baik, pilih kolom dengan rentang nilai besar dan sedikit nilai kunci duplikat untuk klausa `CLUSTERED BY`. Untuk mengoptimalkan join, Anda juga dapat memilih kunci join atau agregasi yang sering digunakan, yang mirip dengan primary key pada database tradisional.
SORTED BY
Menentukan urutan pengurutan bidang dalam satu bucket. Untuk kinerja yang lebih baik, atur kunci `SORTED BY` agar sama dengan kunci `CLUSTERED BY`. Saat klausa `SORTED BY` ditentukan, MaxCompute secara otomatis membuat indeks dan menggunakannya untuk mempercepat kueri.
INTO number_of_buckets BUCKETS
Menentukan jumlah bucket hash. Jumlah ini wajib ditentukan dan bergantung pada volume data. Jumlah bucket yang lebih besar meningkatkan konkurensi dan dapat mempersingkat waktu proses pekerjaan (job runtime). Namun, terlalu banyak bucket dapat menghasilkan jumlah file kecil yang berlebihan, dan konkurensi tinggi dapat meningkatkan waktu CPU. Anda sebaiknya mengatur jumlah bucket sehingga setiap bucket berukuran 500 MB hingga 1 GB. Untuk tabel yang sangat besar, jumlah ini dapat lebih besar. Untuk mengoptimalkan join dengan menghilangkan langkah shuffle dan sort, jumlah bucket untuk kedua tabel harus merupakan kelipatan satu sama lain, misalnya
256dan512. Anda sebaiknya menggunakan pangkat 2 untuk jumlah bucket, seperti 512, 1.024, 2.048, atau 4.096. Hal ini memungkinkan sistem secara otomatis membagi dan menggabungkan bucket hash serta menghilangkan langkah shuffle dan sort.
Mengubah properti Hash Clustering tabel
Anda dapat menggunakan pernyataan ALTER TABLE untuk menambahkan atau menghapus properti Hash Clustering pada tabel partisi.
Pernyataan
-- Mengubah tabel menjadi tabel Hash Clustering ALTER TABLE <table_name> [CLUSTERED BY (<col_name> [, <col_name>, ...]) [SORTED BY (<col_name> [ASC | DESC] [, <col_name> [ASC | DESC] ...])] INTO <number_of_buckets> BUCKETS]; -- Mengubah tabel Hash Clustering menjadi tabel non-Hash Clustering ALTER TABLE <table_name> NOT CLUSTERED;Catatan
Pernyataan `ALTER TABLE` hanya mengubah properti kluster untuk tabel partisi. Untuk tabel non-partisi, properti kluster tidak dapat diubah setelah ditetapkan.
Pernyataan `ALTER TABLE` hanya berlaku untuk partisi baru pada tabel partisi, termasuk partisi yang dihasilkan oleh `INSERT OVERWRITE`. Partisi baru disimpan dengan properti kluster yang baru, sedangkan partisi data yang sudah ada tetap tidak berubah.
Jangan tentukan klausa `PARTITION` dalam pernyataan karena pernyataan ini hanya berlaku untuk partisi baru.
Pernyataan `ALTER TABLE` cocok untuk tabel yang sudah ada. Setelah Anda menambahkan properti kluster baru, partisi baru akan disimpan menggunakan Hash Clustering.
Memverifikasi properti tabel
Setelah membuat tabel Hash Clustering, Anda dapat menjalankan perintah berikut untuk melihat propertinya. Properti Hash Clustering ditampilkan pada bagian Extended Info.
DESC EXTENDED <table_name>;Contoh hasil yang dikembalikan adalah sebagai berikut.
| Owner: ALIYUN$ | Project:
| TableComment:
|
| CreateTime: 2017-06-19 14:10:55
| LastDDLTime: 2017-06-19 14:10:55
| LastModifiedTime: 2017-06-19 14:13:13
|
| InternalTable: YES | Size: 21680295746
|
| Native Columns:
|
| Field | Type | Label | Comment
|
| l_orderkey | bigint | |
| l_partkey | bigint | |
| l_suppkey | bigint | |
| l_linenumber | bigint | |
| l_quantity | double | |
| l_extendedprice | double | |
| l_discount | double | |
| l_tax | double | |
| l_returnflag | string | |
| l_linestatus | string | |
| l_shipdate | string | |
| l_commitdate | string | |
| l_receiptdate | string | |
| l_shipinstruct | string | |
| l_shipmode | string | |
| l_comment | string | |
|
| Extended Info:
|
| TableID:
| IsArchived: false
| PhysicalSize: 65040887238
| FileNum: 1001
| ClusterType: hash
| BucketNum: 1000
| ClusterColumns: [l_orderkey]
| SortColumns: [l_orderkey ASC]Untuk tabel partisi, setelah melihat properti tabel, Anda dapat menjalankan perintah berikut untuk melihat properti partisi.
DESC EXTENDED <table_name> partition(<pt_spec>);Contoh hasil yang dikembalikan adalah sebagai berikut.
| PartitionSize: 754
| CreateTime: 2017-07-07 14:01:03
| LastDDLTime: 2017-07-07 14:01:03
| LastModifiedTime: 2017-07-07 14:01:03
| IsExstore: false
| IsArchived: false
| PhysicalSize: 2262
| FileNum: 2
| ClusterType: hash
| BucketNum: 500
| ClusterColumns: [c1]
| SortColumns: [c1 ASC]Manfaat Hash Clustering
Pemangkasan bucket dan optimasi indeks
CREATE TABLE t1 (id bigint,
a string,
b string)
CLUSTERED BY (id)
SORTED BY (id) into 1000 BUCKETS;
...
SELECT t1.a, t1.b FROM t1 WHERE t1.id=12345;idid
Kueri menemukan bucket hash yang sesuai untuk nilai
12345. Hal ini hanya memerlukan pemindaian satu bucket, bukan semua 1.000 bucket. Proses ini dikenal sebagai bucket pruning.Karena data dalam bucket diurutkan berdasarkan
id, MaxCompute secara otomatis membuat indeks. MaxCompute kemudian menggunakan pencarian berbasis indeks untuk langsung menemukan catatan yang relevan.
Optimasi ini tidak hanya mengurangi jumlah mapper secara signifikan, tetapi juga memungkinkan mapper langsung menemukan halaman data menggunakan indeks, sehingga secara drastis mengurangi jumlah data yang dimuat dan dibaca.
Sebagai contoh, sebuah tugas data besar menjalankan 1.111 mapper dan membaca 42,7 miliar catatan untuk menemukan 26 catatan yang cocok. Waktu proses total adalah 1 menit 48 detik. Dengan tabel Hash Clustering, kueri yang sama pada data yang sama dapat langsung menemukan satu bucket dan menggunakan indeks untuk hanya membaca halaman yang berisi data kueri. Proses ini hanya menggunakan 4 mapper, membaca 10.000 catatan, dan membutuhkan waktu hanya 6 detik.
Optimasi agregasi
Untuk kueri berikut:
SELECT department, SUM(salary) FROM employee GROUP BY (department);Umumnya, kueri ini melakukan shuffle dan sort data pada kolom department, lalu melakukan stream aggregation untuk menghitung setiap kelompok department. Namun, jika data tabel sudah dikluster dan diurutkan berdasarkan `department`, operasi shuffle dan sort tidak lagi diperlukan.
Optimasi penyimpanan
Bahkan tanpa mempertimbangkan optimasi komputasi, hanya dengan melakukan shuffle dan sort data tabel untuk penyimpanan saja sudah dapat menghemat ruang secara signifikan. MaxCompute menggunakan column store di lapisan dasar. Pengurutan menempatkan catatan dengan nilai kunci yang sama atau mirip secara berdekatan, sehingga meningkatkan efektivitas kompresi dan encoding serta menghasilkan rasio kompresi yang lebih tinggi. Dalam pengujian, tabel yang diurutkan dapat menggunakan ruang penyimpanan hingga 50% lebih sedikit dibandingkan tabel yang tidak diurutkan dalam beberapa kasus ekstrem. Untuk tabel dengan siklus hidup panjang, menggunakan Hash Clustering untuk penyimpanan merupakan optimasi yang layak dipertimbangkan.
Eksperimen berikut menggunakan tabel lineitem berukuran 100 GB dari dataset TPC-H. Tabel ini berisi berbagai tipe data, seperti int, double, dan string. Dengan data dan metode kompresi yang sama, kami membandingkan ukuran penyimpanan tabel dengan dan tanpa Hash Clustering. Tabel dengan Hash Clustering menggunakan ruang penyimpanan sekitar 10% lebih sedikit, seperti yang ditunjukkan pada gambar berikut.
Tanpa Hash Clustering
odps@xxx>desc tpch_lineitem; +------------------------------------------------------------------------------------+ | Owner: xxx | Project: xxx | | TableComment: | +------------------------------------------------------------------------------------+ | CreateTime: 2016-04-17 21:48:08 | | LastDDLTime: 2016-04-17 21:48:08 | | LastModifiedTime: 2016-04-17 21:50:10 | +------------------------------------------------------------------------------------+ | InternalTable: YES | Size: 23573055432 | +------------------------------------------------------------------------------------+ | Native Columns: | +------------------------------------------------------------------------------------+ | Field | Type | Label | Comment | +------------------------------------------------------------------------------------+ | l_orderkey | bigint | | | | l_partkey | bigint | | | | l_suppkey | bigint | | | | l_linenumber | bigint | | | | l_quantity | double | | | | l_extendedprice | double | | | | l_discount | double | | | | l_tax | double | | | | l_returnflag | string | | | | l_linestatus | string | | | | l_shipdate | string | | | | l_commitdate | string | | | | l_receiptdate | string | | | | l_shipinstruct | string | | | | l_shipmode | string | | | | l_comment | string | | | +------------------------------------------------------------------------------------+Dengan Hash Clustering
odps@ xxx >desc tpch_lineitem_hash_500; | Owner: xxx | Project: xxx | TableComment: | CreateTime: 2017-07-13 14:40:11 | LastDDLTime: 2017-07-13 14:40:11 | LastModifiedTime: 2017-07-13 15:05:04 | InternalTable: YES | Size: 21658913950 | Native Columns: | Field | Type | Label | Comment | l_orderkey | bigint | | | l_partkey | bigint | | | l_suppkey | bigint | | | l_linenumber | bigint | | | l_quantity | double | | | l_extendedprice | double | | | l_discount | double | | | l_tax | double | | | l_returnflag | string | | | l_linestatus | string | | | l_shipdate | string | | | l_commitdate | string | | | l_receiptdate | string | | | l_shipinstruct | string | | | l_shipmode | string | | | l_comment | string | |
Data uji dan analisis
Manfaat kinerja keseluruhan dari Hash Clustering diukur menggunakan set uji standar TPC-H. Pengujian menggunakan data 1 TB dan 500 bucket untuk semua tabel. Kecuali dua tabel kecil, nation dan region, semua tabel lain menggunakan kolom pertama sebagai kunci kluster dan kunci pengurutan. Hasil pengujian keseluruhan menunjukkan bahwa setelah menggunakan Hash Clustering, total waktu CPU berkurang sekitar 17,3%, dan total waktu proses pekerjaan berkurang sekitar 12,8%.
Perlu diperhatikan bahwa tidak semua kueri dalam TPC-H dapat memanfaatkan properti kluster. Secara khusus, dua kueri dengan waktu proses terlama tidak dapat menggunakan properti ini. Oleh karena itu, peningkatan efisiensi keseluruhan tidak terlalu dramatis. Namun, untuk kueri yang dapat memanfaatkan properti kluster, manfaatnya sangat signifikan. Sebagai contoh, Q4 menjadi sekitar 68% lebih cepat, Q12 menjadi sekitar 62% lebih cepat, dan Q10 menjadi sekitar 47% lebih cepat.
Gambar berikut menunjukkan rencana eksekusi Fuxi untuk TPC-H Q4 pada tabel standar:
Gambar berikut menunjukkan rencana eksekusi setelah Hash Clustering digunakan. Seperti yang terlihat, Directed Acyclic Graph (DAG) menjadi jauh lebih sederhana. Inilah alasan utama peningkatan kinerja yang signifikan.