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.
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) |
|
Nama tabel data. |
|
index_name (wajib) |
|
Nama indeks pencarian. |
Permintaan pemindaian
parallel_scan mencakup parameter berikut.
|
Nama |
Tipe |
Deskripsi |
|
table_name (wajib) |
|
Nama tabel data. |
|
index_name (wajib) |
|
Nama indeks pencarian. |
|
scan_query (wajib) |
|
Kondisi pemindaian, konfigurasi paginasi, dan konkurensi. |
|
session_id (wajib) |
|
Pengidentifikasi sesi yang dikembalikan oleh |
|
columns_to_get (opsional) |
|
Konfigurasi kolom yang dikembalikan. Jika dihilangkan, hanya kolom kunci primer yang dikembalikan. |
|
timeout_s (opsional) |
|
Timeout permintaan dalam detik. |
Konfigurasi pemindaian
scan_query bertipe ScanQuery dan mencakup parameter berikut.
|
Nama |
Tipe |
Deskripsi |
|
query (wajib) |
|
Kondisi kueri yang menentukan cakupan pemindaian. Gunakan |
|
limit (wajib) |
|
Jumlah maksimum baris per permintaan. Kami merekomendasikan nilai default |
|
next_token (wajib) |
|
Token paginasi. Atur ke |
|
current_parallel_id (wajib) |
|
ID worker saat ini. Nilai valid: |
|
max_parallel (wajib) |
|
Konkurensi pekerjaan, yang tidak boleh melebihi |
|
alive_time (opsional) |
|
Periode validitas maksimum antara dua permintaan halaman, dalam detik. Nilai valid: 1 hingga 600. Nilai default: |
Kolom yang dikembalikan
columns_to_get bertipe ColumnsToGet dan mencakup parameter berikut.
|
Nama |
Type |
Deskripsi |
|
column_names (opsional) |
|
Nama bidang indeks pencarian yang akan dikembalikan. Tentukan parameter ini hanya ketika |
|
return_type (opsional) |
|
Mode kolom yang dikembalikan. Pemindaian paralel mendukung |
Respons
Informasi split
compute_splits mengembalikan informasi split.
|
Bidang |
Tipe |
Deskripsi |
|
session_id |
|
Pengidentifikasi sesi pekerjaan. |
|
splits_size |
|
Konkurensi maksimum yang didukung oleh indeks pencarian. |
Hasil pemindaian
parallel_scan mengembalikan hasil pemindaian.
|
Bidang |
Tipe |
Deskripsi |
|
rows |
|
Baris yang dikembalikan oleh permintaan. |
|
next_token |
|
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))