Tablestore SDK for Java を使用すると、結果の順序が不要な場合に、多次元インデックス内の一致する行を同時にスキャンし、完全な結果セットをエクスポートできます。
前提条件
Tablestore SDK for Java をインストールし、クライアントを初期化します。
並列スキャンには、Tablestore SDK for Java 5.6.0 以降が必要です。バージョン情報については、「バージョン履歴」をご参照ください。
仕組み
並列スキャンは、多次元インデックス内のクエリに一致するすべての行をスキャンします。このスキャンは、結果のグローバルな順序を保証せず、ソートや集約をサポートしません。順序付けられた結果、集約、またはエンドユーザー向けの検索結果が必要な場合は、検索操作を使用してください。
シングルワーカースキャンは、設定がより簡単です。マルチワーカースキャンは、複数のスプリットを同時に読み取り、通常、シングルワーカースキャンよりも高いスキャン スループットを提供します。
並列スキャンは、次のステップで構成されます:
computeSplitsを呼び出して、検索インデックスの最大並列度splitsSizeとタスクセッション IDsessionIdを取得します。ParallelScanRequestを設定します。シングルワーカースキャンの場合、maxParallelとcurrentParallelIdは省略できます。マルチワーカースキャンの場合、すべてのワーカーで同じクエリ、sessionId、およびmaxParallelを使用し、各ワーカーに異なるcurrentParallelIdを指定します。createParallelScanIteratorを呼び出してすべてのページを自動的に読み取るか、parallelScanを呼び出してnextTokenを使用して後続のページを手動で取得します。すべてのワーカーが終了するのを待ち、順序付けられていない結果をマージします。
sessionId 内では、parallelScan が初めて呼び出されたときにデータスナップショットが固定されます。タスクの実行中に追加または更新された行は、スナップショットに含まれません。sessionId は省略できます。ただし、スキャン中にサーバー側の負荷分散または同様の変更が発生した場合、結果に少数の重複行が含まれる可能性があります。最初に computeSplits を呼び出し、返された sessionId を後続のリクエストに含めることを推奨します。
動的スキーマの変更によってインデックスが切り替わったり、サーバー側のフェールオーバーや負荷分散が発生したりすると、セッションが早期に有効期限切れになることがあります。この場合、サーバーは OTSSessionExpired を返します。クライアント側のネットワークエラーによってもスキャンが中断されることがあります。これらのエラーのいずれかが発生した場合は、不完全な結果を破棄し、再度 computeSplits を呼び出して、スキャンタスク全体を最初から再起動してください。多次元インデックスでは、最大 10 個の並列スキャンタスクを同時に実行できます。その他の制限については、「多次元インデックスの制限」をご参照ください。
computeSplits を呼び出してスプリットを計算します。parallelScan を呼び出して手動でページを取得するか、createParallelScanIterator を呼び出してすべてのページを自動的に取得します。
ComputeSplitsResponse computeSplits(ComputeSplitsRequest request)
ParallelScanResponse parallelScan(ParallelScanRequest request)
RowIterator createParallelScanIterator(ParallelScanRequest request)
次のサンプルは、1 つのワーカーを使用してすべての行をスキャンし、category と price フィールドを返します。RowIterator は後続のページを自動的に取得します。
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);
}
パラメーター
スプリットリクエスト
splitsRequest は ComputeSplitsRequest 型で、次のパラメーターを含みます。
|
名前 |
型 |
説明 |
|
tableName (必須) |
String |
データテーブルの名前。 |
|
splitsOptions (必須) |
SplitsOptions |
分割構成。検索インデックスの場合、このパラメーターを |
検索インデックス分割構成
splitsRequest.splitsOptions は SearchIndexSplitsOptions 型で、次のパラメーターを含みます。
|
名前 |
型 |
説明 |
|
indexName (必須) |
String |
検索インデックスの名前です。 |
スキャンリクエスト
request は ParallelScanRequest 型で、次のパラメーターを含みます。
|
名前 |
型 |
説明 |
|
tableName (必須) |
String |
データテーブルの名前。 |
|
indexName (必須) |
String |
検索インデックスの名前。 |
|
scanQuery (必須) |
ScanQuery |
スキャン条件、リクエストごとに返される行数、および並列度の設定。 |
|
columnsToGet (任意) |
SearchRequest.ColumnsToGet |
返す列。このパラメーターを指定しない場合、プライマリキー列のみが返されます。 |
|
sessionId (任意) |
|
|
|
timeoutInMillisecond (任意) |
int |
リクエストレベルのタイムアウト期間 (ミリ秒単位)。デフォルト値は |
スキャン設定
request.scanQuery は ScanQuery 型で、次のパラメーターを含みます。
|
名前 |
型 |
説明 |
|
query (必須) |
Query |
スキャン範囲を定義するクエリ条件です。並列スキャンは、term、マッチ、レンジ、ジオ、ネストなどのクエリをサポートしています。クエリは、Search 操作と同様に設定します。検索インデックス内のすべての行をスキャンするには、クエリタイプを |
|
limit (任意) |
Integer |
リクエストごとに返される行の最大数。デフォルト値:2000。デフォルト値を使用することを推奨します。 |
|
maxParallel (任意) |
Integer |
スキャンタスクの並列度。値は |
|
currentParallelId (任意) |
Integer |
現在のワーカーの ID。 |
|
aliveTime (任意) |
Integer |
スキャンタスクの 2 つのページリクエスト間の最大間隔 (秒単位)。有効値:1~600。デフォルト値:60。有効期間は、行が正常に取得されるたびにリフレッシュされます。 |
|
token (任意) |
|
ページネーショントークン。最初のリクエストではこのパラメーターを省略します。手動ページネーションの場合、前の応答の |
サーバーがサポートする limit の最大値は 10000 です。limit をこの値に設定しないことを推奨します。
返される列
request.columnsToGet は SearchRequest.ColumnsToGet 型で、次のパラメーターを含みます。
|
名前 |
型 |
説明 |
|
columns (任意) |
|
返す多次元インデックスフィールドの名前。データテーブルにのみ存在し、多次元インデックスに含まれていないフィールドは返せません。Date、Geo-point、IP、Vector、JSON/Nested、および配列フィールドを返すことができます。 |
|
returnAllFromIndex (任意) |
boolean |
多次元インデックス内のすべてのフィールドを返すかどうかを指定します。デフォルト値: |
|
returnAll (任意) |
boolean |
並列スキャンはこのパラメーターをサポートしていません。 |
戻り値
スプリット情報
computeSplits は ComputeSplitsResponse を返します。これには次のフィールドが含まれます。
|
名前 |
型 |
説明 |
|
sessionId |
|
同じデータスナップショット内の行をスキャンするために使用されるタスクセッション ID。 |
|
splitsSize |
Integer |
検索インデックスがサポートする最大の並列度。 |
スキャン結果
parallelScan は ParallelScanResponse を返します。これには次のフィールドが含まれます。
|
名前 |
型 |
説明 |
|
rows |
|
現在の応答で返された行。 |
|
nextToken |
|
次のページのトークン。この値が |
|
bodyBytes |
long |
レスポンスボディのサイズ (バイト単位)。 |
createParallelScanIterator は RowIterator を返します。イテレータは自動的に nextToken を使用して後続のページを取得し、反復ごとに 1 つの Row を返します。一致する行の総数を取得することはサポートしていません。
例
複数ワーカーによるスキャン
次のサンプルでは、splitsSize に基づいて複数のスキャンタスクを作成します。各タスクは一意の currentParallelId を使用し、すべてのタスクは同じ sessionId と maxParallel を共有します。スレッドプールは、クライアントの過剰な負荷を避けるために、クライアントの CPU コア数を超えません。
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();
}
手動でのページ取得
次のサンプルでは、parallelScan を直接呼び出し、現在のワーカーがすべての行を取得するまで、各応答の nextToken を次のリクエストに渡します。
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);