All Products
Search
Document Center

Tablestore:Parallel scan

Last Updated:Jul 30, 2026

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:

  1. Panggil computeSplits untuk mendapatkan jumlah maksimum split (splitsSize) dan ID sesi tugas (sessionId) dari search index.

  2. Konfigurasikan ParallelScanRequest. Untuk pemindaian satu worker, Anda dapat mengabaikan parameter maxParallel dan currentParallelId. Untuk pemindaian multi-worker, gunakan kueri yang sama, sessionId, dan maxParallel untuk semua worker, serta tetapkan nilai currentParallelId yang berbeda untuk setiap worker.

  3. Panggil createParallelScanIterator untuk membaca semua halaman secara otomatis, atau panggil parallelScan dan gunakan nextToken untuk mengambil halaman berikutnya secara manual.

  4. 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.

Penting

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 SearchIndexSplitsOptions.

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)

byte[]

ID sesi tugas yang dikembalikan oleh computeSplits. Kami menyarankan Anda menentukan parameter ini agar menggunakan snapshot data yang sama selama pemindaian.

timeoutInMillisecond (optional)

int

Periode timeout tingkat permintaan dalam milidetik. Nilai default adalah -1, yang berarti tidak ada timeout tingkat permintaan.

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 MatchAllQuery.

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 ComputeSplitsResponse.splitsSize. Nilai default: 1.

currentParallelId (optional)

Integer

ID worker saat ini. Parameter ini wajib ditentukan jika maxParallel lebih besar dari 1. Tetapkan nilai unik untuk setiap worker dalam rentang [0, maxParallel).

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)

byte[]

Token paginasi. Abaikan parameter ini pada permintaan pertama. Untuk paginasi manual, atur ke nextToken dari respons sebelumnya. SDK mengelola parameter ini saat Anda menggunakan createParallelScanIterator.

Catatan

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)

List<String>

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: false. Jika Anda mengatur parameter ini ke true, Anda tidak perlu menentukan columns.

returnAll (optional)

boolean

Parallel scan tidak mendukung parameter ini. Jangan atur ke true.

Nilai kembalian

Informasi split

computeSplits mengembalikan objek ComputeSplitsResponse yang berisi bidang-bidang berikut.

Name

Type

Description

sessionId

byte[]

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

List<Row>

Baris yang dikembalikan dalam respons saat ini.

nextToken

byte[]

Token untuk halaman berikutnya. Jika nilai ini null, worker saat ini telah mengambil semua baris.

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);