Anda dapat menggunakan Tablestore SDK for Java untuk memindai secara konkuren baris-baris yang sesuai dalam search index dan mengekspor seluruh set hasil ketika urutan hasil tidak diperlukan.
Prasyarat
Instal Tablestore SDK for Java dan inisialisasi client.
Fitur parallel scan memerlukan Tablestore SDK for Java versi 5.6.0 atau lebih baru. Untuk informasi versi, lihat Version history.
Cara kerja
Parallel scan memindai semua baris yang sesuai dengan kueri dalam search index. Pemindaian ini tidak menjamin urutan global hasil dan tidak mendukung sorting atau agregasi. Jika Anda memerlukan hasil terurut, agregasi, atau hasil pencarian untuk pengguna akhir, gunakan operasi Search.
Pemindaian dengan satu worker lebih mudah dikonfigurasi, sedangkan pemindaian multi-worker membaca beberapa split secara konkuren dan umumnya memberikan throughput pemindaian lebih tinggi dibandingkan pemindaian satu worker.
Parallel scan mencakup langkah-langkah berikut:
Panggil
computeSplitsuntuk mendapatkan jumlah maksimum split (splitsSize) dan ID sesi tugas (sessionId) dari search index.Konfigurasikan
ParallelScanRequest. Untuk pemindaian satu worker, Anda dapat mengabaikan parametermaxParalleldancurrentParallelId. Untuk pemindaian multi-worker, gunakan kueri yang sama,sessionId, danmaxParalleluntuk semua worker, serta tetapkan nilaicurrentParallelIdyang berbeda untuk setiap worker.Panggil
createParallelScanIteratoruntuk membaca semua halaman secara otomatis, atau panggilparallelScandan gunakannextTokenuntuk mengambil halaman berikutnya secara manual.Tunggu hingga semua worker selesai dan gabungkan hasil mereka yang tidak terurut.
Snapshot data ditetapkan dalam sessionId saat parallelScan dipanggil pertama kali. Baris yang ditambahkan atau diperbarui selama tugas berjalan tidak termasuk dalam snapshot tersebut. Anda dapat mengabaikan sessionId, namun jika terjadi load balancing di sisi server atau perubahan serupa selama pemindaian, hasilnya mungkin berisi sejumlah kecil baris duplikat. Oleh karena itu, disarankan untuk memanggil computeSplits terlebih dahulu dan menyertakan sessionId yang dikembalikan dalam permintaan berikutnya.
Sesi dapat kedaluwarsa lebih awal jika perubahan skema dinamis mengganti indeks atau terjadi failover atau load balancing di sisi server. Dalam kasus tersebut, server mengembalikan error OTSSessionExpired. Gangguan jaringan di sisi client juga dapat mengganggu pemindaian. Jika salah satu kondisi ini terjadi, buang hasil yang tidak lengkap, panggil kembali computeSplits, dan mulai ulang seluruh tugas pemindaian dari awal. Maksimal 10 tugas parallel scan dapat berjalan secara konkuren pada satu search index. Untuk batasan lainnya, lihat Search index limits.
Panggil computeSplits untuk menghitung split. Gunakan parallelScan untuk mengambil halaman secara manual, atau gunakan createParallelScanIterator untuk mengambil semua halaman secara otomatis.
ComputeSplitsResponse computeSplits(ComputeSplitsRequest request)
ParallelScanResponse parallelScan(ParallelScanRequest request)
RowIterator createParallelScanIterator(ParallelScanRequest request)
Contoh berikut memindai semua baris menggunakan satu worker dan mengembalikan bidang category dan price. RowIterator secara otomatis mengambil halaman berikutnya.
String tableName = "example_table";
String indexName = "example_index";
ComputeSplitsRequest splitsRequest = ComputeSplitsRequest.newBuilder()
.tableName(tableName)
.splitsOptions(new SearchIndexSplitsOptions(indexName))
.build();
ComputeSplitsResponse splitsResponse =
client.computeSplits(splitsRequest);
ScanQuery scanQuery = ScanQuery.newBuilder()
.query(QueryBuilders.matchAll())
.limit(2000)
.build();
ParallelScanRequest request = ParallelScanRequest.newBuilder()
.tableName(tableName)
.indexName(indexName)
.scanQuery(scanQuery)
.addColumnsToGet("category", "price")
.sessionId(splitsResponse.getSessionId())
.build();
RowIterator iterator = client.createParallelScanIterator(request);
while (iterator.hasNext()) {
Row row = iterator.next();
System.out.println(row);
}
Parameter
Permintaan split
splitsRequest bertipe ComputeSplitsRequest dan berisi parameter berikut.
|
Name |
Type |
Description |
|
tableName (required) |
String |
Nama tabel data. |
|
splitsOptions (required) |
SplitsOptions |
Konfigurasi split. Untuk search index, atur parameter ini ke |
Konfigurasi split search index
splitsRequest.splitsOptions bertipe SearchIndexSplitsOptions dan berisi parameter berikut.
|
Name |
Type |
Description |
|
indexName (required) |
String |
Nama search index. |
Permintaan pemindaian
request bertipe ParallelScanRequest dan berisi parameter berikut.
|
Name |
Type |
Description |
|
tableName (required) |
String |
Nama tabel data. |
|
indexName (required) |
String |
Nama search index. |
|
scanQuery (required) |
ScanQuery |
Kondisi pemindaian, jumlah baris yang dikembalikan per permintaan, dan konfigurasi parallelisme. |
|
columnsToGet (optional) |
SearchRequest.ColumnsToGet |
Kolom yang akan dikembalikan. Jika parameter ini tidak ditentukan, hanya kolom kunci primer yang dikembalikan. |
|
sessionId (optional) |
|
ID sesi tugas yang dikembalikan oleh |
|
timeoutInMillisecond (optional) |
int |
Periode timeout tingkat permintaan dalam milidetik. Nilai default adalah |
Konfigurasi pemindaian
request.scanQuery bertipe ScanQuery dan berisi parameter berikut.
|
Name |
Type |
Description |
|
query (required) |
Query |
Kondisi kueri yang menentukan cakupan pemindaian. Parallel scan mendukung kueri term, match, range, geo, nested, dan lainnya. Konfigurasikan kueri ini sama seperti pada operasi Search. Untuk memindai semua baris dalam search index, atur tipe kueri ke |
|
limit (optional) |
Integer |
Jumlah maksimum baris yang dikembalikan per permintaan. Nilai default: 2000. Kami menyarankan Anda menggunakan nilai default. |
|
maxParallel (optional) |
Integer |
Parallelisme tugas pemindaian. Nilainya tidak boleh melebihi |
|
currentParallelId (optional) |
Integer |
ID worker saat ini. Parameter ini wajib ditentukan jika |
|
aliveTime (optional) |
Integer |
Interval maksimum antara dua permintaan halaman untuk tugas pemindaian, dalam detik. Nilai valid: 1 hingga 600. Nilai default: 60. Periode validitas diperbarui setiap kali baris berhasil diambil. |
|
token (optional) |
|
Token paginasi. Abaikan parameter ini pada permintaan pertama. Untuk paginasi manual, atur ke |
Nilai maksimum limit yang didukung server adalah 10.000. Disarankan agar Anda tidak mengatur limit ke nilai maksimum tersebut.
Kolom yang dikembalikan
request.columnsToGet bertipe SearchRequest.ColumnsToGet dan berisi parameter berikut.
|
Name |
Type |
Description |
|
columns (optional) |
|
Nama bidang search index yang akan dikembalikan. Bidang yang hanya ada di tabel data tetapi tidak termasuk dalam search index tidak dapat dikembalikan. Bidang Date, Geo-point, IP, Vector, JSON/Nested, dan array dapat dikembalikan. |
|
returnAllFromIndex (optional) |
boolean |
Menentukan apakah semua bidang dalam search index dikembalikan. Nilai default: |
|
returnAll (optional) |
boolean |
Parallel scan tidak mendukung parameter ini. Jangan atur ke |
Nilai kembalian
Informasi split
computeSplits mengembalikan objek ComputeSplitsResponse yang berisi bidang-bidang berikut.
|
Name |
Type |
Description |
|
sessionId |
|
ID sesi tugas yang digunakan untuk memindai baris dalam snapshot data yang sama. |
|
splitsSize |
Integer |
Parallelisme maksimum yang didukung oleh search index. |
Hasil pemindaian
parallelScan mengembalikan objek ParallelScanResponse yang berisi bidang-bidang berikut.
|
Name |
Type |
Description |
|
rows |
|
Baris yang dikembalikan dalam respons saat ini. |
|
nextToken |
|
Token untuk halaman berikutnya. Jika nilai ini |
|
bodyBytes |
long |
Ukuran badan respons dalam byte. |
createParallelScanIterator mengembalikan RowIterator. Iterator ini secara otomatis menggunakan nextToken untuk mengambil halaman berikutnya dan mengembalikan satu objek Row per iterasi. Iterator ini tidak mendukung pengambilan jumlah total baris yang sesuai.
Contoh
Pemindaian dengan banyak worker
Contoh berikut membuat beberapa tugas pemindaian berdasarkan nilai splitsSize. Setiap tugas menggunakan currentParallelId yang unik, dan semua tugas berbagi sessionId serta maxParallel yang sama. Ukuran kolam thread tidak melebihi jumlah core CPU pada client untuk menghindari beban berlebih pada client.
String tableName = "example_table";
String indexName = "example_index";
ComputeSplitsResponse splitsResponse = client.computeSplits(
ComputeSplitsRequest.newBuilder()
.tableName(tableName)
.splitsOptions(new SearchIndexSplitsOptions(indexName))
.build());
int maxParallel = splitsResponse.getSplitsSize();
int workerCount = Math.min(
maxParallel, Runtime.getRuntime().availableProcessors());
ExecutorService executor = Executors.newFixedThreadPool(workerCount);
List<Future<Integer>> futures = new ArrayList<Future<Integer>>();
try {
for (int parallelId = 0; parallelId < maxParallel; parallelId++) {
final int currentParallelId = parallelId;
futures.add(executor.submit(new Callable<Integer>() {
@Override
public Integer call() {
ScanQuery scanQuery = ScanQuery.newBuilder()
.query(QueryBuilders.matchAll())
.limit(2000)
.maxParallel(maxParallel)
.currentParallelId(currentParallelId)
.build();
ParallelScanRequest request =
ParallelScanRequest.newBuilder()
.tableName(tableName)
.indexName(indexName)
.scanQuery(scanQuery)
.returnAllColumnsFromIndex(true)
.sessionId(splitsResponse.getSessionId())
.build();
int rowCount = 0;
RowIterator iterator =
client.createParallelScanIterator(request);
while (iterator.hasNext()) {
Row row = iterator.next();
System.out.println(row);
rowCount++;
}
return rowCount;
}
}));
}
long totalRows = 0;
for (Future<Integer> future : futures) {
totalRows += future.get();
}
System.out.println("Total rows: " + totalRows);
} finally {
executor.shutdown();
}
Mengambil halaman secara manual
Contoh berikut langsung memanggil parallelScan dan meneruskan nextToken dari setiap respons ke permintaan berikutnya hingga worker saat ini telah mengambil semua baris.
String tableName = "example_table";
String indexName = "example_index";
ComputeSplitsResponse splitsResponse = client.computeSplits(
ComputeSplitsRequest.newBuilder()
.tableName(tableName)
.splitsOptions(new SearchIndexSplitsOptions(indexName))
.build());
ScanQuery scanQuery = ScanQuery.newBuilder()
.query(QueryBuilders.matchAll())
.limit(2000)
.maxParallel(1)
.currentParallelId(0)
.build();
ParallelScanRequest request = ParallelScanRequest.newBuilder()
.tableName(tableName)
.indexName(indexName)
.scanQuery(scanQuery)
.addColumnsToGet("category", "price")
.sessionId(splitsResponse.getSessionId())
.build();
long totalRows = 0;
do {
ParallelScanResponse response = client.parallelScan(request);
for (Row row : response.getRows()) {
System.out.println(row);
totalRows++;
}
scanQuery.setToken(response.getNextToken());
} while (scanQuery.getToken() != null);
System.out.println("Total rows: " + totalRows);