All Products
Search
Document Center

Tablestore:Pemindaian paralel

Last Updated:Aug 07, 2026

Gunakan Tablestore SDK untuk Python untuk memindai data yang sesuai dalam indeks pencarian secara paralel dan mengekspor set hasil lengkap yang tidak terurut.

Prasyarat

Instal Tablestore SDK untuk Python dan inisialisasi klien.

Deskripsi

Pemindaian paralel mengekspor semua baris yang sesuai dengan kueri dalam indeks pencarian. Hasilnya tidak diurutkan secara global, serta pengurutan dan agregasi tidak didukung. Untuk mengurutkan atau mengagregasi hasil atau menampilkannya kepada pengguna akhir, gunakan API Search. Satu worker lebih sederhana dikonfigurasi, sedangkan beberapa worker membaca beberapa split secara konkuren dan umumnya memberikan throughput yang lebih tinggi.

Alur kerjanya adalah sebagai berikut: Panggil compute_splits untuk mendapatkan jumlah konkurensi maksimum (splits_size) dan pengidentifikasi sesi (session_id). Konfigurasikan setiap worker dengan kueri yang sama, session_id, max_parallel, serta current_parallel_id yang unik. Setiap worker memanggil parallel_scan dan menggunakan next_token untuk membaca halaman berikutnya. Terakhir, tunggu semua worker dan gabungkan hasil yang tidak terurut tersebut.

Penting

Worker dengan session_id yang sama membuat snapshot data saat pemindaian pertama dimulai. Pembaruan skema dinamis, failover, atau load balancing dapat menyebabkan sesi kedaluwarsa lebih awal dan mengembalikan OTSSessionExpired; kegagalan jaringan juga dapat mengganggu pemindaian. Dalam kasus ini, buang hasil yang belum lengkap, panggil kembali compute_splits, lalu mulai ulang semua worker dari awal. Maksimal 10 pekerjaan pemindaian paralel dapat berjalan secara konkuren pada satu indeks pencarian.

Contoh berikut menghitung split lalu memindai seluruh data menggunakan satu worker.

splits = client.compute_splits("example_table", "example_index")
next_token = None
rows = []

while True:
    scan_query = ScanQuery(
        MatchAllQuery(),
        limit=2000,
        next_token=next_token,
        current_parallel_id=0,
        max_parallel=1,
        alive_time=60,
    )
    response = client.parallel_scan(
        "example_table",
        "example_index",
        scan_query,
        splits.session_id,
        ColumnsToGet(return_type=ColumnReturnType.ALL_FROM_INDEX),
    )
    rows.extend(response.rows)
    next_token = response.next_token
    if not next_token:
        break

print(len(rows))

Parameter

Compute splits

compute_splits(table_name, index_name) mencakup parameter berikut.

Nama

Tipe

Deskripsi

table_name (wajib)

str

Nama tabel data.

index_name (wajib)

str

Nama indeks pencarian.

Permintaan pemindaian

parallel_scan mencakup parameter berikut.

Nama

Tipe

Deskripsi

table_name (wajib)

str

Nama tabel data.

index_name (wajib)

str

Nama indeks pencarian.

scan_query (wajib)

ScanQuery

Kondisi pemindaian, konfigurasi paginasi, dan konkurensi.

session_id (wajib)

bytes

Pengidentifikasi sesi yang dikembalikan oleh compute_splits, digunakan untuk mempertahankan snapshot data yang sama.

columns_to_get (opsional)

ColumnsToGet

Konfigurasi kolom yang dikembalikan. Jika dihilangkan, hanya kolom kunci primer yang dikembalikan.

timeout_s (opsional)

int

Timeout permintaan dalam detik.

Konfigurasi pemindaian

scan_query bertipe ScanQuery dan mencakup parameter berikut.

Nama

Tipe

Deskripsi

query (wajib)

Query

Kondisi kueri yang menentukan cakupan pemindaian. Gunakan MatchAllQuery untuk memindai semua baris.

limit (wajib)

int

Jumlah maksimum baris per permintaan. Kami merekomendasikan nilai default 2000. Server mengizinkan hingga 10000, tetapi penggunaan nilai maksimum tidak disarankan.

next_token (wajib)

bytes

Token paginasi. Atur ke None pada permintaan pertama dan gunakan next_token dari respons sebelumnya pada permintaan berikutnya.

current_parallel_id (wajib)

int

ID worker saat ini. Nilai valid: [0, max_parallel). Setiap worker harus menggunakan nilai yang unik.

max_parallel (wajib)

int

Konkurensi pekerjaan, yang tidak boleh melebihi splits_size.

alive_time (opsional)

int

Periode validitas maksimum antara dua permintaan halaman, dalam detik. Nilai valid: 1 hingga 600. Nilai default: 60. Periode ini diperbarui setelah respons berhasil.

Kolom yang dikembalikan

columns_to_get bertipe ColumnsToGet dan mencakup parameter berikut.

Nama

Type

Deskripsi

column_names (opsional)

list[str]

Nama bidang indeks pencarian yang akan dikembalikan. Tentukan parameter ini hanya ketika return_type bernilai SPECIFIED.

return_type (opsional)

ColumnReturnType

Mode kolom yang dikembalikan. Pemindaian paralel mendukung NONE, SPECIFIED, dan ALL_FROM_INDEX, tetapi tidak mendukung ALL.

Respons

Informasi split

compute_splits mengembalikan informasi split.

Bidang

Tipe

Deskripsi

session_id

bytes

Pengidentifikasi sesi pekerjaan.

splits_size

int

Konkurensi maksimum yang didukung oleh indeks pencarian.

Hasil pemindaian

parallel_scan mengembalikan hasil pemindaian.

Bidang

Tipe

Deskripsi

rows

list[Row]

Baris yang dikembalikan oleh permintaan.

next_token

bytes

Token untuk halaman berikutnya. Nilai kosong menunjukkan bahwa worker saat ini telah selesai.

Respons kompatibel tupel

Pemindaian paralel didukung mulai dari Tablestore SDK untuk Python versi 5.2.0, yang mengembalikan objek respons. Pada versi 5.2.1 dan seterusnya, Anda dapat memanggil ComputeSplitsResponse.v1_response() dan ParallelScanResponse.v1_response() untuk mendapatkan tupel. Untuk kode baru, akses atribut objek respons secara langsung.

session_id, splits_size = splits.v1_response()
rows, next_token = response.v1_response()

Contoh

Pemindaian dengan beberapa worker

Contoh berikut membuat worker berdasarkan splits_size. Kolam thread tidak melebihi jumlah core CPU klien, dan semua worker berbagi sesi serta konkurensi maksimum.

from concurrent.futures import ThreadPoolExecutor
import os


def scan_split(parallel_id, max_parallel, session_id):
    rows = []
    next_token = None
    while True:
        scan_query = ScanQuery(
            MatchAllQuery(),
            limit=2000,
            next_token=next_token,
            current_parallel_id=parallel_id,
            max_parallel=max_parallel,
            alive_time=60,
        )
        response = client.parallel_scan(
            "example_table",
            "example_index",
            scan_query,
            session_id,
            ColumnsToGet(return_type=ColumnReturnType.ALL_FROM_INDEX),
        )
        rows.extend(response.rows)
        next_token = response.next_token
        if not next_token:
            return rows


splits = client.compute_splits("example_table", "example_index")
worker_count = min(splits.splits_size, os.cpu_count() or 1)
with ThreadPoolExecutor(max_workers=worker_count) as executor:
    futures = [
        executor.submit(
            scan_split,
            parallel_id,
            splits.splits_size,
            splits.session_id,
        )
        for parallel_id in range(splits.splits_size)
    ]
    all_rows = [row for future in futures for row in future.result()]

print(len(all_rows))