すべてのプロダクト
Search
ドキュメントセンター

Tablestore:並列スキャン

最終更新日:Jul 30, 2026

Tablestore SDK for Java を使用すると、結果の順序が不要な場合に、多次元インデックス内の一致する行を同時にスキャンし、完全な結果セットをエクスポートできます。

前提条件

Tablestore SDK for Java をインストールし、クライアントを初期化します。

並列スキャンには、Tablestore SDK for Java 5.6.0 以降が必要です。バージョン情報については、「バージョン履歴」をご参照ください。

仕組み

並列スキャンは、多次元インデックス内のクエリに一致するすべての行をスキャンします。このスキャンは、結果のグローバルな順序を保証せず、ソートや集約をサポートしません。順序付けられた結果、集約、またはエンドユーザー向けの検索結果が必要な場合は、検索操作を使用してください。

シングルワーカースキャンは、設定がより簡単です。マルチワーカースキャンは、複数のスプリットを同時に読み取り、通常、シングルワーカースキャンよりも高いスキャン スループットを提供します。

並列スキャンは、次のステップで構成されます:

  1. computeSplits を呼び出して、検索インデックスの最大並列度 splitsSize とタスクセッション ID sessionId を取得します。

  2. ParallelScanRequest を設定します。シングルワーカースキャンの場合、maxParallel と currentParallelId は省略できます。マルチワーカースキャンの場合、すべてのワーカーで同じクエリ、sessionId、および maxParallel を使用し、各ワーカーに異なる currentParallelId を指定します。

  3. createParallelScanIterator を呼び出してすべてのページを自動的に読み取るか、parallelScan を呼び出して nextToken を使用して後続のページを手動で取得します。

  4. すべてのワーカーが終了するのを待ち、順序付けられていない結果をマージします。

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

分割構成。検索インデックスの場合、このパラメーターをSearchIndexSplitsOptionsに設定します。

検索インデックス分割構成

splitsRequest.splitsOptions は SearchIndexSplitsOptions 型で、次のパラメーターを含みます。

名前

型

説明

indexName (必須)

String

検索インデックスの名前です。

スキャンリクエスト

request は ParallelScanRequest 型で、次のパラメーターを含みます。

名前

型

説明

tableName (必須)

String

データテーブルの名前。

indexName (必須)

String

検索インデックスの名前。

scanQuery (必須)

ScanQuery

スキャン条件、リクエストごとに返される行数、および並列度の設定。

columnsToGet (任意)

SearchRequest.ColumnsToGet

返す列。このパラメーターを指定しない場合、プライマリキー列のみが返されます。

sessionId (任意)

byte[]

computeSplits によって返されるタスクセッション ID。スキャン全体で同じデータスナップショットを使用するために、このパラメーターを指定することを推奨します。

timeoutInMillisecond (任意)

int

リクエストレベルのタイムアウト期間 (ミリ秒単位)。デフォルト値は -1 で、リクエストレベルのタイムアウトが設定されていないことを示します。

スキャン設定

request.scanQuery は ScanQuery 型で、次のパラメーターを含みます。

名前

型

説明

query (必須)

Query

スキャン範囲を定義するクエリ条件です。並列スキャンは、term、マッチ、レンジ、ジオ、ネストなどのクエリをサポートしています。クエリは、Search 操作と同様に設定します。検索インデックス内のすべての行をスキャンするには、クエリタイプを MatchAllQuery に設定します。

limit (任意)

Integer

リクエストごとに返される行の最大数。デフォルト値:2000。デフォルト値を使用することを推奨します。

maxParallel (任意)

Integer

スキャンタスクの並列度。値は ComputeSplitsResponse.splitsSize を超えることはできません。デフォルト値:1。

currentParallelId (任意)

Integer

現在のワーカーの ID。maxParallel が 1 より大きい場合、このパラメーターは必須です。各ワーカーに [0, maxParallel) の範囲で一意の値を割り当てます。

aliveTime (任意)

Integer

スキャンタスクの 2 つのページリクエスト間の最大間隔 (秒単位)。有効値:1~600。デフォルト値:60。有効期間は、行が正常に取得されるたびにリフレッシュされます。

token (任意)

byte[]

ページネーショントークン。最初のリクエストではこのパラメーターを省略します。手動ページネーションの場合、前の応答の nextToken に設定します。createParallelScanIterator を使用する場合、SDK がこのパラメーターを管理します。

説明

サーバーがサポートする limit の最大値は 10000 です。limit をこの値に設定しないことを推奨します。

返される列

request.columnsToGet は SearchRequest.ColumnsToGet 型で、次のパラメーターを含みます。

名前

型

説明

columns (任意)

List<String>

返す多次元インデックスフィールドの名前。データテーブルにのみ存在し、多次元インデックスに含まれていないフィールドは返せません。Date、Geo-point、IP、Vector、JSON/Nested、および配列フィールドを返すことができます。

returnAllFromIndex (任意)

boolean

多次元インデックス内のすべてのフィールドを返すかどうかを指定します。デフォルト値:false。このパラメーターを true に設定した場合、columns を指定する必要はありません。

returnAll (任意)

boolean

並列スキャンはこのパラメーターをサポートしていません。true に設定しないでください。

戻り値

スプリット情報

computeSplits は ComputeSplitsResponse を返します。これには次のフィールドが含まれます。

名前

型

説明

sessionId

byte[]

同じデータスナップショット内の行をスキャンするために使用されるタスクセッション ID。

splitsSize

Integer

検索インデックスがサポートする最大の並列度。

スキャン結果

parallelScan は ParallelScanResponse を返します。これには次のフィールドが含まれます。

名前

型

説明

rows

List<Row>

現在の応答で返された行。

nextToken

byte[]

次のページのトークン。この値が null の場合、現在のワーカーはすべての行を取得済みです。

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