Tablestore SDK for Python を使用して、検索インデックス内の一致するデータを並列でスキャンし、順序付けされていない完全な結果セットをエクスポートします。
前提条件
Tablestore SDK for Python をインストールし、クライアントを初期化します。
説明
並列スキャンは、検索インデックス内のクエリに一致するすべての行をエクスポートします。結果は全体では順序付けされず、ソートと集計はサポートされません。結果をソートまたは集計する場合、または検索結果をエンドユーザーに返す場合は、Search API を使用してください。単一のワーカーは設定が簡単ですが、複数のワーカーは複数のスプリットを同時に読み取り、通常はより高いスループットを提供します。
ワークフローは次のとおりです: compute_splits を呼び出して、最大同時実行数 splits_size とセッション ID session_id を取得します。各ワーカーに、同一のクエリ、session_id、max_parallel、および一意の current_parallel_id を設定します。各ワーカーは parallel_scan を呼び出し、next_token を使用して後続ページを読み取ります。最後に、すべてのワーカーの完了を待機し、それぞれの順序付けされていない結果をマージします。
同一の session_id 配下のワーカーは、最初のスキャン開始時にデータスナップショットを作成します。動的なスキーマ更新、フェールオーバー、または負荷分散によりセッションが早期に期限切れとなり、OTSSessionExpired が返される場合があります。ネットワーク障害によってスキャンが中断されることもあります。これらの場合は、不完全な結果を破棄し、compute_splits を再度呼び出して、すべてのワーカーを最初から再開してください。1 つの検索インデックスで同時に実行できる並列スキャンジョブは最大 10 件です。
次の例では、スプリットを計算した後、1 つのワーカーを使用してすべてのデータをスキャンします。
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))
パラメーター
スプリットの計算
compute_splits(table_name, index_name) には、次のパラメーターが含まれます。
|
名前 |
型 |
説明 |
|
table_name (必須) |
|
データテーブルの名前です。 |
|
index_name (必須) |
|
検索インデックスの名前です。 |
スキャンリクエスト
parallel_scan には、次のパラメーターが含まれます。
|
名前 |
型 |
説明 |
|
table_name (必須) |
|
データテーブルの名前です。 |
|
index_name (必須) |
|
検索インデックスの名前です。 |
|
scan_query (必須) |
|
スキャン条件、ページネーション、および同時実行数の設定です。 |
|
session_id (必須) |
|
|
|
columns_to_get (任意) |
|
返却する列の設定です。省略した場合、プライマリキー列のみが返されます。 |
|
timeout_s (任意) |
|
リクエストのタイムアウト (秒) です。 |
スキャン設定
scan_query は ScanQuery 型で、次のパラメーターが含まれます。
|
名前 |
型 |
説明 |
|
query (必須) |
|
スキャン範囲を定義するクエリ条件です。すべての行をスキャンするには、 |
|
limit (任意) |
|
リクエストあたりの最大行数です。デフォルト値の |
|
next_token (必須) |
|
ページネーショントークン。初回リクエストでは |
|
current_parallel_id (必須) |
|
現在のワーカー ID です。有効な値は |
|
max_parallel (必須) |
|
ジョブの同時実行数は、 |
|
alive_time (任意) |
|
2つのページリクエスト間の最大有効期間 (秒) です。有効値は 1~600、デフォルト値は |
返却列
columns_to_get は ColumnsToGet 型で、次のパラメーターが含まれます。
|
名前 |
型 |
説明 |
|
column_names (任意) |
|
返す検索インデックスフィールド名。このパラメーターは、 |
|
return_type (任意) |
|
返却列モード。パラレルスキャンは |
レスポンス
スプリット情報
compute_splits はスプリット情報を返します。
|
フィールド |
型 |
説明 |
|
session_id |
|
ジョブのセッション ID です。 |
|
splits_size |
|
検索インデックスでサポートされる最大同時実行数です。 |
スキャン結果
parallel_scan はスキャン結果を返します。
|
フィールド |
型 |
説明 |
|
rows |
|
リクエストによって返された行です。 |
|
next_token |
|
次のページのトークンです。空の値は、現在のワーカーが完了したことを示します。 |
タプル互換のレスポンス
並列スキャンは、レスポンスオブジェクトを返す Tablestore SDK for Python 5.2.0 以降でサポートされています。バージョン 5.2.1 以降では、ComputeSplitsResponse.v1_response() と ParallelScanResponse.v1_response() を呼び出してタプルを取得できます。新しいコードでは、レスポンスオブジェクトの属性に直接アクセスしてください。
session_id, splits_size = splits.v1_response()
rows, next_token = response.v1_response()
例
複数ワーカーでのスキャン
以下の例では、splits_size に基づいてワーカーを作成します。スレッドプールはクライアント CPU コアの数を超えず、すべてのワーカーがセッションと最大同時実行数を共有します。
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))