Gunakan Tablestore SDK untuk Go guna membagi data indeks pencarian dan memindai bagian-bagian tersebut secara paralel demi ekspor skala besar yang efisien.
Prasyarat
Sebelum memulai, lakukan persiapan berikut:
Instal Tablestore Go SDK dan inisialisasi client.
Pemindaian paralel memerlukan Tablestore SDK untuk Go versi 1.6.0 atau lebih baru. Disarankan untuk menggunakan versi terbaru.
Deskripsi
Ekspor paralel dimulai dengan memanggil ComputeSplits untuk membuat sesi pemindaian dan mendapatkan tingkat konkurensi yang direkomendasikan. Setiap worker konkuren kemudian memanggil ParallelScan dengan CurrentParallelID unik dan menggunakan NextToken untuk membaca bagian datanya secara berkelanjutan. Pemindaian paralel tidak menjamin urutan set hasil secara keseluruhan, sehingga cocok untuk ekspor skala besar yang tidak bergantung pada urutan hasil. Anda juga dapat mengatur MaxParallel ke 1 dan CurrentParallelID ke 0 untuk pemindaian satu worker. Pendekatan ini menghasilkan kode yang lebih sederhana, dengan throughput biasanya lebih tinggi daripada Search namun lebih rendah dibandingkan pemindaian multi-worker.
Pemindaian paralel tidak mendukung sorting atau agregasi. Gunakan Search jika Anda perlu mengurutkan hasil, melakukan agregasi, atau menampilkan hasil pencarian kepada pengguna akhir.
Worker yang menggunakan SessionId yang sama akan membuat snapshot data pada panggilan ParallelScan pertama. Data yang dimasukkan atau diperbarui selama tugas tidak termasuk dalam snapshot tersebut. SessionId dapat dihilangkan, tetapi perubahan seperti load balancing di sisi server dapat menyebabkan sejumlah kecil baris duplikat. Disarankan untuk memanggil ComputeSplits terlebih dahulu dan menyertakan SessionId yang dikembalikan dalam permintaan berikutnya.
Maksimal 10 tugas pemindaian paralel dapat berjalan secara konkuren pada indeks pencarian yang sama. Untuk batasan lainnya, lihat Batasan indeks pencarian.
Contoh berikut memulai beberapa goroutine berdasarkan tingkat konkurensi yang direkomendasikan oleh ComputeSplits, memberikan CurrentParallelID unik untuk setiap worker, dan membaca semua bagian data.
splits, err := client.ComputeSplits(
(&tablestore.ComputeSplitsRequest{}).
SetTableName("example_table").
SetSearchIndexSplitsOptions(tablestore.SearchIndexSplitsOptions{
IndexName: "example_index",
}),
)
if err != nil {
log.Fatal(err)
}
var waitGroup sync.WaitGroup
var mutex sync.Mutex
totalRows := 0
errors := make(chan error, splits.SplitsSize)
waitGroup.Add(int(splits.SplitsSize))
for workerID := int32(0); workerID < splits.SplitsSize; workerID++ {
currentWorkerID := workerID
go func() {
defer waitGroup.Done()
scanQuery := search.NewScanQuery().
SetQuery(&search.MatchAllQuery{}).
SetLimit(1000).
SetMaxParallel(splits.SplitsSize).
SetCurrentParallelID(currentWorkerID)
request := (&tablestore.ParallelScanRequest{}).
SetTableName("example_table").
SetIndexName("example_index").
SetScanQuery(scanQuery).
SetSessionId(splits.SessionId).
SetColumnsToGet(&tablestore.ColumnsToGet{
ReturnAllFromIndex: true,
})
for {
response, err := client.ParallelScan(request)
if err != nil {
errors <- err
return
}
// Proses response.Rows di sini.
mutex.Lock()
totalRows += len(response.Rows)
mutex.Unlock()
if len(response.NextToken) == 0 {
return
}
request.SetScanQuery(scanQuery.SetToken(response.NextToken))
}
}()
}
waitGroup.Wait()
close(errors)
for err := range errors {
log.Fatal(err)
}
fmt.Println(totalRows)
Parameter
Buat sesi pemindaian
|
Name |
Type |
Description |
|
TableName (required) |
string |
Nama tabel data. |
|
IndexName (required) |
string |
Nama indeks pencarian. |
Permintaan pemindaian paralel
|
Name |
Type |
Description |
|
TableName (required) |
string |
Nama tabel data. |
|
IndexName (required) |
string |
Nama indeks pencarian. |
|
ScanQuery (required) |
search.ScanQuery |
Kondisi pemindaian dan konfigurasi konkurensi. |
|
SessionId (optional) |
[]byte |
ID sesi yang dikembalikan oleh ComputeSplits. Kami menyarankan Anda menentukan parameter ini agar menggunakan snapshot data yang sama selama pemindaian. |
|
ColumnsToGet (optional) |
*tablestore.ColumnsToGet |
Konfigurasi kolom yang dikembalikan. Jika dihilangkan, hanya kolom kunci primer yang dikembalikan. ParallelScan tidak mendukung ReturnAll. |
|
TimeoutMs (optional) |
*int32 |
Periode timeout permintaan dalam milidetik. |
Konfigurasi pemindaian
|
Name |
Type |
Description |
|
Query (required) |
search.Query |
Kondisi pemindaian. Jenis kueri non-vektor yang didukung oleh Search juga didukung. |
|
Limit (optional) |
int32 |
Jumlah maksimum baris yang dikembalikan per permintaan. Nilai default: 2.000. Kami menyarankan Anda mempertahankan nilai default. Server mengizinkan nilai hingga 10.000, tetapi nilai yang lebih besar meningkatkan latensi permintaan dan penggunaan resource. |
|
MaxParallel (optional) |
int32 |
Jumlah total worker. Nilai default: 1. Nilai ini tidak boleh melebihi SplitsSize yang dikembalikan oleh ComputeSplits. |
|
CurrentParallelID (optional) |
int32 |
ID worker saat ini. Parameter ini wajib ditentukan ketika MaxParallel lebih dari 1. ID worker harus unik dan berada dalam rentang [0, MaxParallel). |
|
Token (optional) |
[]byte |
Nilai NextToken yang dikembalikan oleh respons sebelumnya. |
|
AliveTime (optional) |
int32 |
Interval maksimum antara dua permintaan paginasi, dalam detik. Nilai valid: 1 hingga 600. Nilai default: 60. Kami menyarankan Anda mempertahankan nilai default. Setiap permintaan data yang berhasil akan memperbarui periode validitas. |
Konfigurasi kolom yang dikembalikan
|
Name |
Type |
Description |
|
Columns (optional) |
[]string |
Bidang indeks pencarian yang akan dikembalikan. Bidang yang hanya ada di tabel data dan tidak termasuk dalam indeks pencarian tidak dapat dikembalikan. |
|
ReturnAllFromIndex (optional) |
bool |
Menentukan apakah semua bidang dari indeks pencarian akan dikembalikan. Nilai default: false. Jika parameter ini bernilai true, Anda tidak perlu menentukan Columns. |
|
ReturnAll (optional) |
bool |
Pemindaian paralel tidak mendukung parameter ini. Jangan atur nilainya menjadi true. |
Perubahan skema yang mengganti indeks, failover server, atau load balancing dapat membatalkan sesi lebih awal dan mengembalikan OTSSessionExpired. Error jaringan client juga dapat mengganggu pemindaian. Jika error semacam ini terjadi, buang hasil yang tidak lengkap, panggil ComputeSplits lagi, dan mulai ulang seluruh tugas pemindaian dari awal.
Respons
Informasi bagian
|
Name |
Type |
Description |
|
SessionId |
[]byte |
ID sesi tugas yang digunakan untuk memindai snapshot data yang sama. |
|
SplitsSize |
int32 |
Konkurensi maksimum yang didukung untuk indeks pencarian. |
Hasil pemindaian
|
Name |
Type |
Description |
|
Rows |
[]*tablestore.Row |
Baris yang dikembalikan oleh pemindaian saat ini. |
|
NextToken |
[]byte |
Token untuk halaman berikutnya. Lanjutkan memindai bagian saat ini jika nilainya tidak kosong. |