Pengelompokan rentang adalah metode pengelompokan data baru yang mendistribusikan data dalam urutan terurut secara global. Metode ini mencegah masalah kesenjangan data yang mungkin terjadi pada pengelompokan hash dan memungkinkan pembuatan indeks dua tingkat. Pengelompokan rentang cocok untuk skenario seperti kueri rentang berdasarkan kunci kluster dan kueri multi-kunci. Topik ini menjelaskan cara menggunakan pengelompokan rentang di MaxCompute.
Informasi Latar Belakang
Tabel terkluster hash memiliki keunggulan berikut:
-
Jika Anda ingin mengkueri data berdasarkan nilai kolom tertentu, algoritma hash dapat langsung menemukan bucket hash—proses yang disebut pemangkasan bucket. Jika data dalam bucket disimpan secara terurut, indeks tambahan dapat digunakan untuk menemukan data tersebut. Pendekatan ini mengurangi jumlah data yang dipindai dan meningkatkan efisiensi kueri.
-
Jika Anda melakukan join dua tabel berdasarkan kolom tertentu di masing-masing tabel dan salah satu kolom tersebut di-hash, langkah shuffle dapat dihilangkan, sehingga menghemat sumber daya komputasi.
Untuk informasi selengkapnya tentang fitur pengelompokan hash, lihat Hash clustering.
Pengelompokan hash memiliki batasan berikut:
-
Penggunaan algoritma hash untuk membuat bucket dapat menyebabkan masalah kesenjangan data. Mirip dengan skew saat join, masalah ini bersifat inheren dalam algoritma hash. Jika distribusi data masukan tidak merata ke bucket, jumlah data di setiap bucket akan sangat berbeda. Karena setiap bucket umumnya menjadi unit pemrosesan konkuren dalam pengelompokan hash, perbedaan ukuran bucket ini berpotensi menyebabkan long tail.
-
Pemangkasan bucket hanya mendukung kueri equality. Untuk kueri berbasis kondisi non-equality—misalnya, nilai suatu kolom lebih besar dari 0—bucket tempat data tersimpan tidak dapat ditentukan, sehingga semua bucket harus dikueri.
-
Untuk kueri berbasis beberapa kunci kluster, peningkatan performa hanya terjadi jika semua kunci kluster tersedia dan semua kondisi kueri merupakan kondisi equality.
Sebagai contoh, pertimbangkan tabel yang dibuat menggunakan pernyataan berikut. Performa kueri hanya meningkat jika kondisi kueri berupa
C1=x AND C2=y. Jika kondisi kueri hanyaC1=xatauC2=y, kueri tidak dapat dipercepat melalui pengelompokan hash. Hal ini karena nilai hash kunci digabungkan berpasangan saat digunakan untuk kueri; tanpa penggabungan tersebut, bucket tempat data tersimpan tidak dapat ditentukan dan pemangkasan bucket tidak dapat diterapkan.CREATE TABLE T2 (C1 int, C2 int, C3 string) CLUSTERED BY (C1, C2) SORTED by (C1, C2) INTO 1024 BUCKETS;
Untuk mengatasi keterbatasan tersebut, MaxCompute menyediakan metode pengelompokan data baru, yaitu pengelompokan rentang.
Deskripsi fitur
Pengelompokan rentang membagi data menjadi beberapa rentang yang saling lepas berdasarkan pengurutan penuh kunci kluster. Setiap rentang dianggap sebagai bucket dan harus memenuhi kedua kondisi berikut:
-
Nilai duplikat disimpan dalam bucket yang sama.
-
Jumlah nilai di setiap bucket kira-kira sama.
Pernyataan contoh berikut membuat tabel bernama T.
CREATE TABLE T (C1 int)
RANGE CLUSTERED BY (C1)
SORTED BY (c1)
INTO 3 BUCKETS;
Nilai dalam kolom C1 adalah { 1, 8, -3, 2, 4, 1, 1, 3, 8, 20, -8, 9 }.
Setelah pengelompokan rentang diaktifkan, bucket berikut diperoleh.
-
Bucket 0 : { -8, -3, 1, 1, 1 }
-
Bucket 1 : { 2, 3, 4 }
-
Bucket 2 : { 8, 8, 9, 20 }
-
Rentang yang diwakili oleh bucket mungkin saling lepas. Misalnya, rentang Bucket 1 adalah
[2, 4]dan rentang Bucket 2 adalah[8, 20]. Tidak ada nilai dalam rentang(4, 8). -
Pengelompokan rentang bertujuan menyamakan ukuran bucket, bukan ukuran rentang. Dalam pemrosesan data besar, setiap bucket merupakan unit pemrosesan konkuren. Ukuran bucket yang seragam mencegah masalah long tail, meskipun distribusi data dalam setiap rentang mungkin tidak merata. Oleh karena itu, konsistensi ukuran bucket tidak berarti konsistensi ukuran rentang.
Proses pengelompokan rentang diimplementasikan secara otomatis oleh MaxCompute. Anda tidak perlu menentukan rentang secara manual. Dalam skenario data besar, konfigurasi rentang manual tidak efisien atau bahkan tidak layak. MaxCompute secara otomatis mengurutkan dan mengambil sampel data, membuat histogram berdasarkan distribusi data di setiap rentang, lalu menggabungkan dan menghitung histogram tersebut untuk mencapai performa optimal pengelompokan rentang.
Saat membuat tabel, Anda dapat menentukan RANGE CLUSTERED BY dan SORTED BY untuk memastikan data diurutkan secara global. MaxCompute kemudian secara otomatis membuat dua tingkat indeks—indeks global dan indeks file—untuk menemukan dan mencari nilai kunci dengan cepat, seperti yang ditunjukkan pada gambar berikut.
Pengelompokan rentang memiliki keunggulan berikut dibandingkan pengelompokan hash:
-
Mendukung kueri rentang.
Misalnya, jika kondisi kueri adalah
c < 3, sistem dapat mengecualikan Bucket 2 dan Bucket 3 berdasarkan indeks global, lalu hanya mengkueri data dari Bucket 0 dan Bucket 1. Sebaliknya, pengelompokan hash hanya mendukung pemangkasan bucket untuk kueri equality. -
Mendukung kueri multi-kunci.
Misalnya, jika Anda menentukan
RANGE CLUSTERED BY (c1, c2, c3) SORTED BY (c1, c2, c3)saat membuat tabel, pengelompokan rentang dan penyimpanan data diimplementasikan dalam urutan c1, c2, dan c3. Dengan demikian, Anda dapat mengkueri data berdasarkan kondisi kompleks sepertic1 = 100 AND c2 > 0atauc1 = 100 AND c2 = 50 AND c3 < 5—jenis kueri yang tidak dapat dioptimalkan dengan pengelompokan hash.PentingUntuk kueri berbasis beberapa kunci, kunci dalam kondisi kueri harus diurutkan secara berurutan, dan hanya kunci terakhir yang dapat digunakan untuk menentukan rentang nilai.
-
Pengelompokan Rentang merupakan implementasi efisien dari pengurutan global.
Sebelum pengelompokan rentang, MaxCompute hanya dapat menggunakan satu instans untuk mengurutkan data secara global, sehingga efisiensinya rendah. Dengan pengelompokan rentang, data di setiap rentang dapat diurutkan secara konkuren lalu digabungkan, yang secara signifikan meningkatkan efisiensi.
Catatan penggunaan
Sintaks pengelompokan rentang mirip dengan pengelompokan hash. Perbedaannya terletak pada kata kunci RANGE dan fakta bahwa jumlah bucket bersifat opsional dalam pengelompokan rentang.
Buat tabel terkluster rentang
Anda dapat menggunakan pernyataan CREATE TABLE untuk membuat tabel terkluster rentang. Dalam pernyataan ini, parameter RANGE CLUSTERED BY wajib ditentukan, sedangkan INTO number_of_buckets BUCKETS dan SORTED BY bersifat opsional. Umumnya, kami merekomendasikan agar nilai yang ditentukan dalam SORTED BY sama dengan yang ditentukan dalam RANGE CLUSTERED BY untuk mencapai efek optimasi optimal.
-
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>], ...)] [RANGE 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 int) RANGE CLUSTERED BY (c) SORTED by (c) INTO 1024 BUCKETS; -
Tabel partisi
CREATE TABLE T1 (a string, b string, c int) PARTITIONED BY (dt int) RANGE CLUSTERED BY (c) SORTED by (c) INTO 1024 BUCKETS;
-
-
Parameter
-
RANGE CLUSTERED BY
Menentukan kunci untuk pengelompokan rentang. Setelah parameter ini ditentukan, MaxCompute mengurutkan dan mengambil sampel satu atau beberapa kolom data ke rentang yang sesuai berdasarkan jumlah bucket yang Anda tentukan. Untuk mencegah masalah kesenjangan data dan hot spot serta meningkatkan performa kueri konkuren, kami merekomendasikan agar Anda menentukan kolom dengan rentang nilai besar dan sedikit nilai kunci duplikat di RANGE CLUSTERED BY. Untuk mengoptimalkan performa kueri, gunakan kunci agregat atau kunci filter yang umum digunakan.
-
SORTED BY
Menentukan cara mengurutkan bidang dalam bucket. Untuk mencapai performa kueri optimal, kami merekomendasikan agar nilai yang ditentukan dalam SORTED BY sama dengan yang ditentukan dalam RANGE CLUSTERED BY. Setelah SORTED BY ditentukan, MaxCompute secara otomatis menghasilkan indeks global dan indeks file serta menggunakannya untuk mempercepat kueri.
-
INTO number_of_buckets BUCKETS
Berbeda dengan pengelompokan hash,
INTO number_of_buckets BUCKETSbersifat opsional untuk pengelompokan rentang. Jika tidak ditentukan, MaxCompute secara otomatis menentukan jumlah bucket berdasarkan volume data. Secara umum, kami merekomendasikan agar Anda menentukan jumlah bucket sesuai situasi aktual.Seperti halnya pengelompokan hash, kami merekomendasikan agar jumlah bucket ditetapkan berdasarkan ukuran bucket ideal (512 MB hingga 1 GB). Untuk tabel yang sangat besar, diperlukan jumlah bucket yang besar, tetapi sebaiknya tidak melebihi 4.000 bucket per tabel.
-
Ubah properti pengelompokan hash tabel
Untuk tabel partisi, Anda dapat menjalankan pernyataan ALTER TABLE untuk menambahkan atau menghapus properti pengelompokan rentang.
-
Sintaks
-- Ubah tabel menjadi tabel terkluster rentang. ALTER TABLE <table_name> [RANGE CLUSTERED BY (<col_name> [, <col_name>, ...]) [SORTED BY (<col_name> [ASC | DESC] [, <col_name> [ASC | DESC] ...])] [INTO <number_of_buckets> BUCKETS]; -- Ubah tabel terkluster rentang menjadi tabel non-terkluster rentang. ALTER TABLE <table_name> NOT CLUSTERED; -
Catatan penggunaan
-
Pernyataan ALTER TABLE hanya dapat memodifikasi properti pengelompokan tabel partisi. Untuk tabel non-partisi, properti pengelompokan tidak dapat diubah setelah ditetapkan.
-
Pernyataan ALTER TABLE hanya berlaku untuk partisi baru, termasuk partisi yang dihasilkan melalui pernyataan INSERT OVERWRITE. Partisi baru disimpan sesuai properti pengelompokan yang baru, sedangkan format penyimpanan partisi yang sudah ada tetap tidak berubah.
-
Karena pernyataan ALTER TABLE hanya berlaku untuk partisi baru, Anda tidak dapat menentukan partisi tertentu dalam pernyataan ini.
-
Pernyataan ALTER TABLE cocok untuk tabel yang sudah ada. Setelah properti pengelompokan rentang ditambahkan, partisi baru akan disimpan sesuai properti tersebut.
Verifikasi eksplisit properti tabel
Setelah membuat tabel terkluster rentang, Anda dapat menjalankan pernyataan berikut untuk melihat properti tabel. Properti pengelompokan rentang ditampilkan di bagian Extended Info pada hasil yang dikembalikan.
DESC EXTENDED <table_name>;
Contoh output:
Owner: ALIYUNxxx | Project: xxx
TableComment:
+----------------------------+
| CreateTime: 2018-01-15 16:05:30 |
| LastDDLTime: 2018-01-15 16:05:30 |
| LastModifiedTime: 2018-01-15 16:05:30 |
+----------------------------+
| InternalTable: YES | Size: 0 |
+----------------------------+
| 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: xxx |
| IsArchived: false |
| PhysicalSize: 0 |
| FileNum: 0 |
| ClusterType: range |
| BucketNum: 1024 |
| ClusterColumns: [l_orderkey] |
| SortColumns: [l_orderkey ASC] |
+----------------------------+
Pada bagian Extended Info, ClusterType bernilai range, BucketNum bernilai 1024, ClusterColumns bernilai [l_orderkey], dan SortColumns bernilai [l_orderkey ASC], yang mengonfirmasi bahwa tabel telah berhasil dikonfigurasi dengan atribut Range Clustering.
Untuk tabel partisi, Anda juga dapat menjalankan pernyataan berikut untuk melihat properti pengelompokan partisi.
DESC EXTENDED <table_name> partition(<pt_spec>);
Dalam output, empat atribut ClusterType, BucketNum, ClusterColumns, dan SortColumns menunjukkan konfigurasi Range Clustering partisi.
odps@ tpch_100g>desc extended xndai_test_range partition(pt="20180115");
| PartitionSize: 0
| CreateTime: 2018-01-15 16:31:10
| LastDDLTime: 2018-01-15 16:31:10
| LastModifiedTime: 2018-01-15 16:31:10
| IsExstore: false
| IsArchived: false
| PhysicalSize: 0
| FileNum: 0
| ClusterType: range
| BucketNum: 1024
| ClusterColumns: [c1]
| SortColumns: [c1 ASC]
Skenario
Optimasi kueri dengan penyaringan
Jika pengelompokan rentang diaktifkan untuk tabel, data dalam tabel diurutkan secara global. MaxCompute secara otomatis membuat indeks global dan indeks file berdasarkan data yang diurutkan, sehingga meningkatkan efisiensi penyaringan data berdasarkan karakteristik penyimpanan. Anda dapat menggunakan pengelompokan rentang untuk mengoptimalkan kueri equality maupun kueri rentang.
Sebagai contoh, untuk kondisi kueri sederhana id < 3, sistem mengekstrak kondisi dari pengoptimal dan mengonversinya menjadi rentang nilai (-∞, 3). Dalam kasus ini, sistem menggunakan indeks global untuk pemangkasan bucket guna mengecualikan Bucket 2 dan Bucket 3 yang datanya berada di luar rentang tersebut. Selanjutnya, sistem menggunakan indeks file di Bucket 0 dan Bucket 1 untuk menemukan data dengan cepat. Proses ini disebut penurunan predikat, seperti yang ditunjukkan pada gambar berikut.
Pernyataan contoh berikut adalah TPC-H Query 6 yang digunakan untuk mengkueri data dari set data 100 GB setelah pengelompokan rentang diterapkan. Dalam TPC-H Query 6, operasi agregat dilakukan berdasarkan penyaringan rentang. Pengelompokan rentang memanfaatkan indeks dua tingkat untuk menemukan data dengan cepat, sehingga durasi eksekusi kueri serta konsumsi CPU dan memori berkurang secara signifikan.
select sum(l_extendedprice * l_discount) as revenue
from tpch_lineitem l
where l_shipdate >= '1994-01-01'
and l_shipdate < '1995-01-01'
and l_discount >= 0.05
and l_discount <= 0.07
and l_quantity < 24;
Kueri multi-kunci
Dalam contoh ini, pernyataan berikut digunakan untuk mengubah tabel mf_tab menjadi tabel terkluster rentang guna mempermudah pemahaman kueri multi-kunci.
ALTER TABLE mf_project.mf_tab
RANGE CLUSTERED BY (project_name, name)
SORTED BY (project_name, name)
INTO 1024 BUCKETS;
Setelah tabel diubah menjadi tabel terkluster rentang, Anda dapat melakukan kueri agregat di tingkat proyek. Contoh pernyataan:
SELECT COUNT(*)
from mf_project.mf_tab
WHERE project_name="xxxdw"
AND ds="20180115"
AND type="TABLE";
Anda juga dapat menggunakan beberapa kunci untuk menemukan tabel secara tepat. Contoh pernyataan:
SELECT count(*)
from mf_project.mf_tab
WHERE project_name="xxxdw"
AND name="adm_ctu_cle_kba_midun_trade_dd"
AND type="TABLE";
Anda juga dapat menggunakan beberapa kunci untuk kueri rentang. Pernyataan berikut mengkueri tabel yang namanya dimulai dengan adm.
SELECT count(*)
from mf_project.mf_tab
WHERE project_name="xxxdw"
AND name>="adm"
AND name < "adn"
AND type="TABLE";
Semua kueri di atas dapat sepenuhnya memanfaatkan fitur pengurutan global pengelompokan rentang dan menerapkan penurunan predikat untuk mengurangi jumlah operasi I/O saat pemindaian tabel serta menghemat sumber daya CPU dan memori yang digunakan untuk penyaringan dan komputasi data.
Jika beberapa kunci digunakan untuk pengelompokan rentang, persyaratan tertentu harus dipenuhi. Untuk RANGE CLUSTERED BY k0, k1, ..., kn dalam pernyataan pembuatan tabel, jika km digunakan untuk kueri data, maka k0, k1, ..., km-1 harus semuanya ditentukan dalam kondisi kueri dan semua kondisi tersebut harus berupa kondisi equality agar akselerasi kueri berbasis indeks mencapai performa optimal.
Sebagai contoh, k1, k2 adalah kunci kluster dalam tabel bernama T.
-
Jika kondisi kueri adalah
k1 < 5, akselerasi kueri berbasis indeks dapat dicapai. -
Jika kondisi kueri adalah
k1 = 10 AND k2 = 20, akselerasi kueri berbasis indeks dapat dicapai. -
Jika kondisi kueri adalah
k1 = 10 AND k2 < 0, akselerasi kueri berbasis indeks dapat dicapai. -
Jika kondisi kueri adalah
k2 < 0, akselerasi kueri berbasis indeks tidak dapat dicapai karena k1 tidak ditentukan dalam kondisi kueri. -
Jika kondisi kueri adalah
k1 < 0 AND k2 > 0, akselerasi kueri berbasis indeks dapat digunakan untuk memperoleh data yang memenuhi kondisik1 < 0. Namun, untuk data yang memenuhi kondisik2 > 0, diperlukan pemindaian tabel penuh.
Optimasi GROUP BY
Jika pengelompokan rentang diaktifkan untuk tabel, data dalam tabel diurutkan secara global. Kunci dengan nilai yang sama ditempatkan dalam bucket yang sama selama pengelompokan rentang. Properti fisik data ini dapat digunakan untuk menghilangkan langkah shuffle selama operasi agregat.
Sebagai contoh, tabel bernama T dibuat menggunakan pernyataan CREATE TABLE berikut. Saat mengkueri data dari tabel tersebut, Anda dapat melakukan operasi GROUP BY pada data tabel dalam tahap map.
CREATE TABLE T (department int, team string, employee string)
RANGE CLUSTERED BY (department, team)
SORTED BY (c1, c2)
INTO 1024 BUCKETS;
SELECT COUNT(*) from T GROUP BY department, team;
Untuk mencapai performa optimal GROUP BY, Anda harus menentukan kunci yang sama di GROUP BY dan RANGE CLUSTERED BY.
Optimasi agregat
Pernyataan berikut menunjukkan struktur data tabel foo.
create table foo(a bigint, b bigint, c bigint)
range clustered by (a,b)
sorted by(a,b) into 3 buckets;
Data yang disimpan dalam bucket tabel foo berada dalam rentang berikut:
Bucket 0: [1,1 : 3,3]
Bucket 1: [5,5 : 7,7]
Bucket 2: [8,8 : 9,9]
Rentang bucket di atas ditentukan dalam format Bucket N: [nilai batas bawah : nilai batas atas]. Jika data diagregat berdasarkan Kolom a, rentang bucket berikut digunakan sebagai gantinya:
Bucket 0: [1 : 3]
Bucket 1: [5 : 7]
Bucket 2: [8 : 9]
Anda dapat langsung menghasilkan rencana eksekusi untuk operasi agregat berdasarkan Kolom a dan Kolom b, memulai tiga instans untuk mengagregat data di setiap bucket, lalu mengembalikan hasil output.
Namun, jika nilai Kolom a tersebar di beberapa bucket, hasil yang tidak valid akan dikembalikan. Contoh:
Bucket 0: [1,1 : 3,3]
Bucket 1: [3,5 : 7,7]
Bucket 2: [7,8 : 9,9]
Kolom a memiliki dua nilai, yaitu 3 dan 7, yang disimpan secara terpisah di dua bucket. Untuk memperoleh hasil yang valid, tupel dengan nilai Kolom a yang sama harus ditempatkan dalam instans yang sama dan diagregat. Dengan demikian, bucket dibuat ulang, seperti yang ditunjukkan pada gambar berikut. Jarak antara dua garis putus-putus merah menentukan rentang data yang dapat dibaca oleh setiap instans.
Histogram diperlukan untuk pengelompokan rentang. Untuk tabel terkluster rentang, jika kunci kluster dan kunci pengurutan sama, worker yang sesuai dengan setiap bucket mengambil sampel satu tupel untuk setiap 10.000 baris guna memperoleh nilai kunci kluster saat data dimasukkan ke tabel. Nilai yang diperoleh disimpan dalam histogram yang disimpan dalam file metadata kluster. Jenis histogram ini disebut histogram equi-depth.
Tupel hanya diambil sampelnya jika kunci kluster dan kunci pengurutan sama.
Setelah histogram setiap bucket diperoleh, bucket dapat dibuat ulang untuk setiap worker berdasarkan aturan berikut:
-
Tupel yang memiliki kunci pengelompokan yang sama disimpan dalam bucket yang sama.
-
Data didistribusikan secara merata ke bucket.
Berdasarkan nilai batas bawah setiap bucket baru, setiap worker dapat membaca data dalam rentang yang valid dan mengembalikan hasil yang valid.
Konten berikut menggunakan tabel partsupp dengan data 1 TB dalam set data TPC-H untuk menguji peningkatan performa. Jalankan pernyataan berikut untuk mengubah tabel partsupp menjadi tabel terkluster rentang:
CREATE TABLE partsupp ( PS_PARTKEY BIGINT NOT NULL,
PS_SUPPKEY BIGINT NOT NULL,
PS_AVAILQTY BIGINT NOT NULL,
PS_SUPPLYCOST DECIMAL(15,2) NOT NULL,
PS_COMMENT VARCHAR(199) NOT NULL)
RANGE CLUSTERED BY(PS_PARTKEY, PS_SUPPKEY)
SORTED BY(PS_PARTKEY, PS_SUPPKEY) INTO 128 BUCKETS;
Jalankan pernyataan kueri berikut untuk melakukan pengujian:
SELECT ps_partkey, count(*) c FROM partsupp GROUP BY ps_partkey;
-
Jalankan perintah berikut untuk menonaktifkan optimasi:
set odps.optimizer.enable.range.partial.repartitioning=false;Berikut ini menunjukkan log eksekusi pekerjaan saat optimasi dinonaktifkan.
resource cost: cpu 9.17 Core * Min, memory 15.15 GB * Min inputs: yuan_tpch_range_1t.partsupp: 800000000 (3376541880 bytes) outputs: Job run time: 14.000 Job run mode: fuxi job Job run engine: execution engine M1: instance count: 128 run time: 7.000 instance time: min: 2.000, max: 4.000, avg: 2.000 input records: TableScan1: 800000000 (min: 5496040, max: 6780663, avg: 6251772) output records: StreamLineWrite1: 200001977 (min: 1374023, max: 1695182, avg: 1562958) writer dumps: StreamLineWrite1: (min: 0, max: 0, avg: 0) R2_1: instance count: 43 run time: 14.000 instance time: min: 3.000, max: 4.000, avg: 3.000 input records: StreamLineRead1: 200001977 (min: 4647288, max: 4654888, avg: 4651214) output records: AdhocSink1: 200000000 (min: 4647242, max: 4654834, avg: 4651168) reader dumps: StreamLineRead1: (min: 0, max: 0, avg: 0) -
Jalankan perintah berikut untuk mengaktifkan optimasi:
set odps.optimizer.enable.range.partial.repartitioning=true;Contoh output:
resource cost: cpu 4.38 Core * Min, memory 4.38 GB * Min inputs: yuan_tpch_range_1t.partsupp: 800000000 (18493328320 bytes) outputs: Job run time: 6.000 Job run mode: fuxi job Job run engine: execution engine M1: instance count: 128 run time: 6.000 instance time: min: 1.000, max: 3.000, avg: 2.000 input records: TableScan1: 800000000 (min: 5625876, max: 6259956, avg: 6254874) output records: AdhocSink1: 200000000 (min: 1406469, max: 1564989, avg: 1563718)
Hasil pengujian menunjukkan bahwa kecepatan kueri meningkat 57%, pemanfaatan CPU berkurang 52%, dan penggunaan memori berkurang 71% setelah optimasi diaktifkan. Peningkatan performa bervariasi tergantung pada volume data dan jenis kueri.
Optimasi join tabel terkluster rentang
-
Dalam contoh ini, buat dua tabel menggunakan pernyataan berikut:
create table t1(a bigint, b bigint, c bigint, d bigint) range clustered by(a,b,c) sorted by(a,b,c) into 3 buckets; create table t2(a bigint, b bigint, c bigint, d bigint) range clustered by(a,b,c) sorted by(a,b,c) into 3 buckets;Kemudian, masukkan data berbeda ke kedua tabel tersebut.
Untuk dua tabel terkluster hash yang perlu di-join, jika jumlah bucket sama untuk kedua tabel, data dalam bucket dapat langsung di-join. Namun, aturan ini tidak berlaku untuk tabel terkluster rentang. Meskipun jumlah bucket sama, data dalam bucket kedua tabel tidak dapat langsung di-join berdasarkan ID bucket karena batas setiap bucket dalam tabel terkluster rentang mungkin berbeda. Akibatnya, rencana eksekusi yang mencakup langkah shuffle selalu dihasilkan, seperti yang ditunjukkan pada gambar berikut.
Untuk mengoptimalkan join antara dua tabel terkluster rentang, bucket kedua tabel perlu dibuat ulang dengan menyelaraskan batasnya, sehingga rentang data yang dapat dibaca oleh setiap instans didefinisikan ulang. -
Buat dua tabel.
create table t1(a bigint, b bigint, c bigint, d bigint) range clustered by(a,b,c) sorted by(a,b,c) into 5 buckets; create table t2(a bigint, b bigint, c bigint, d bigint) range clustered by(a,b,c) sorted by(a,b,c) into 3 buckets;Setelah sejumlah data dimasukkan ke tabel, batas bucket ditentukan, seperti yang ditunjukkan pada gambar berikut.
Contoh kueri 1:SELECT * FROM t1 JOIN t2 ON t1.a=t2.a AND t1.b=t2.b AND t1.c=t2.c;Pengoptimal menyelaraskan batas tabel yang memiliki lebih banyak bucket dengan batas tabel lainnya, lalu memperoleh batas baru untuk setiap tabel, seperti yang ditunjukkan pada gambar berikut.
Dengan demikian, rencana eksekusi tanpa langkah shuffle dihasilkan.
Contoh kueri 2:SELECT * FROM t1 JOIN t2 ON t1.a=t2.a AND t1.b=t2.b;Pengoptimal membuat bucket untuk setiap tabel berdasarkan Kolom a dan Kolom b, lalu menyelaraskan dan membuat ulang bucket untuk memperoleh batas berdasarkan data yang dibaca dari setiap bucket. Dengan pendekatan ini, rencana eksekusi tanpa langkah shuffle dihasilkan, seperti yang ditunjukkan pada gambar sebelumnya.
-
Pengujian Performa
-
Transformasi Tabel
Gunakan TPC-H Query 2 untuk melakukan pengujian pada dua tabel bernama PART dan PARTSUPP, masing-masing berisi data 1 TB. Ubah kedua tabel menjadi tabel terkluster rentang, sementara tabel lain dibiarkan tidak berubah.
CREATE TABLE PARTSUPP ( PS_PARTKEY BIGINT NOT NULL, PS_SUPPKEY BIGINT NOT NULL, PS_AVAILQTY BIGINT NOT NULL, PS_SUPPLYCOST DECIMAL(15,2) NOT NULL, PS_COMMENT VARCHAR(199) NOT NULL) RANGE CLUSTERED BY(PS_PARTKEY, PS_SUPPKEY) SORTED BY(PS_PARTKEY, PS_SUPPKEY) INTO 128 BUCKETS; CREATE TABLE PART ( P_PARTKEY BIGINT NOT NULL, P_NAME VARCHAR(55) NOT NULL, P_MFGR CHAR(25) NOT NULL, P_BRAND CHAR(10) NOT NULL, P_TYPE VARCHAR(25) NOT NULL, P_SIZE BIGINT NOT NULL, P_CONTAINER CHAR(10) NOT NULL, P_RETAILPRICE DECIMAL(15,2) NOT NULL, P_COMMENT VARCHAR(23) NOT NULL) RANGE CLUSTERED BY(P_PARTKEY) SORTED BY(P_PARTKEY) INTO 64 BUCKETS;Kueri data menggunakan TPC-H Query 2 berikut:
select s_acctbal, s_name, n_name, p_partkey, p_mfgr, s_address, s_phone, s_comment from part, supplier, partsupp, nation, region where p_partkey = ps_partkey and s_suppkey = ps_suppkey and p_size = 15 and p_type like '%BRASS' and s_nationkey = n_nationkey and n_regionkey = r_regionkey and r_name = 'EUROPE' and ps_supplycost = (select min(ps_supplycost) from partsupp, supplier, nation, region where p_partkey = ps_partkey and s_suppkey = ps_suppkey and s_nationkey = n_nationkey and n_regionkey = r_regionkey and r_name = 'EUROPE') order by s_acctbal desc, n_name, s_name, p_partkey limit 100; -
Hasil Pengujian
-
Jalankan perintah berikut untuk menonaktifkan optimasi:
set odps.optimizer.enable.range.partial.repartitioning=false;Contoh output log eksekusi pekerjaan saat optimasi dinonaktifkan:
resource cost: cpu 61.64 Core * Min, memory 41.62 GB * Min inputs: yuan_tpch_range_1t.nation: 25 (1848 bytes) yuan_tpch_range_1t.partsupp: 800000000 (7392850104 bytes) yuan_tpch_range_1t.region: 5 (1040 bytes) yuan_tpch_range_1t.part: 200000000 (1427093008 bytes) yuan_tpch_range_1t.supplier: 10000000 (483851352 bytes) outputs: Job run time: 56.000 Job run mode: fuxi job Job run engine: execution engine J11_13: instance count: 299 run time: 52.000 instance time: min: 1.000, max: 2.000, avg: 1.000 input records: StreamLineRead13: 637969 (min: 1978, max: 2313, avg: 2133) StreamLineRead7: 159971440 (min: 532818, max: 537639, avg: 535020) output records: StreamLineWrite14: 470727 (min: 1473, max: 1676, avg: 1574) -
Jalankan perintah berikut untuk mengaktifkan optimasi:
set odps.optimizer.enable.range.partial.repartitioning=true;Contoh output:
resource cost: cpu 39.81 Core * Min, memory 18.89 GB * Min inputs: yuan_tpch_range_1t.nation: 25 (1848 bytes) yuan_tpch_range_1t.region: 5 (1040 bytes) yuan_tpch_range_1t.part: 200000000 (7544753176 bytes) yuan_tpch_range_1t.partsupp: 800000000 (22722759616 bytes) yuan_tpch_range_1t.supplier: 10000000 (483851352 bytes) outputs: Job run time: 44.000 Job run mode: fuxi job Job run engine: execution engine J11_13: instance count: 135 run time: 40.000 instance time: min: 1.000, max: 3.000, avg: 1.000 input records: StreamLineRead11: 637969 (min: 4516, max: 4962, avg: 4725) StreamLineRead7: 159971440 (min: 1181336, max: 1189183, avg: 1184969) output records: StreamLineWrite12: 470727 (min: 3368, max: 3647, avg: 3486) writer dumps: StreamLineWrite12: (min: 0, max: 0, avg: 0) reader dumps: StreamLineRead11: (min: 0, max: 0, avg: 0) StreamLineRead7: (min: 0, max: 0, avg: 0)
Setelah optimasi, dua tahap dihilangkan, kecepatan kueri meningkat sekitar 21,4%, pemanfaatan CPU berkurang sekitar 35,4%, dan penggunaan memori berkurang sekitar 54,6%.
-
-
Akselerasi pengurutan global
Pengelompokan rentang juga dapat digunakan untuk akselerasi pengurutan global. Dalam skenario umum yang menggunakan ORDER BY, semua data yang diurutkan didistribusikan ke satu instans untuk memastikan pengurutan global, sehingga pemrosesan konkuren tidak dapat dimanfaatkan secara optimal. Dengan pengelompokan rentang, Anda dapat menerapkan pengurutan global konkuren melalui langkah partisi: ambil sampel data, bagi data ke dalam rentang, urutkan data di setiap rentang secara paralel, lalu gabungkan hasilnya untuk memperoleh urutan global.
Setelah pengurutan global selesai, beberapa bucket tetap termasuk dalam tabel saat Anda mengubah properti kluster tabel atau partisi. Selama konsumsi data, data dalam file harus dibaca berdasarkan ID bucket untuk mempertahankan pengurutan global.
Secara default, akselerasi pengurutan global dinonaktifkan untuk tabel terkluster rentang. Untuk mengaktifkannya, jalankan perintah berikut:
set odps.optimizer.distribute.ordering.enable=true;
Batasan dan catatan penggunaan
Dibandingkan dengan pengelompokan hash, pengelompokan rentang memiliki batasan berikut:
-
Biaya pembuatan data pengelompokan rentang lebih tinggi daripada pengelompokan hash. Pengelompokan hash hanya melibatkan operasi hashing dan pengurutan data yang sederhana. Sebaliknya, pengelompokan rentang memerlukan pengambilan sampel data, pengurutan, dan penggabungan histogram. Konsumsi keseluruhan—termasuk durasi eksekusi, biaya CPU, dan biaya memori—lebih tinggi dibandingkan pengelompokan hash. Oleh karena itu, jika pengelompokan hash sudah cukup untuk menyelesaikan masalah, tidak perlu menggunakan pengelompokan rentang.
-
Pengelompokan rentang tidak didukung dalam DYNAMIC PARTITION atau INSERT INTO.
-
Pengelompokan rentang hanya didukung untuk operasi join berikut: inner join, left outer join, right outer join, dan semi join. Fitur ini tidak didukung untuk anti-join atau full outer join.
-
Untuk tabel terkluster rentang, kunci yang ditentukan di RANGE CLUSTERED BY harus sama dengan yang ditentukan di SORTED BY. Misalnya, jika
range clustered by (a,b) sorted by (a,b)ditentukan untuk tabel bernama foo, optimasi yang dijelaskan dalam topik ini dapat diterapkan. Namun, jikarange clustered by(a,b) sorted by (b,a)ditentukan untuk tabel bernama bar, optimasi tersebut tidak dapat diterapkan. -
Kunci yang ditentukan di JOIN atau GROUP BY harus merupakan awalan atau seluruh kunci yang ditentukan di RANGE CLUSTERED BY. Misalnya, jika range clustered by(a,b,c) sorted by(a,b,c) ditentukan dalam pernyataan pembuatan tabel, optimasi hanya dapat diterapkan jika kunci di JOIN atau
GROUP BYberupaa,a,b, ataua,b,c. Optimasi tidak dapat diterapkan jika kunci yang digunakan adalahbataua,c. -
Untuk tabel partisi terkluster rentang, optimasi yang dijelaskan dalam topik ini tidak dapat diterapkan jika data dibaca dari dua atau lebih partisi. Optimasi hanya berlaku untuk tabel partisi dengan satu partisi dan tabel non-partisi.