全部產品
Search
文件中心

Tablestore:並發匯出資料

更新時間:Jul 30, 2026

使用 Tablestore Java SDK 可並發掃描多元索引中的匹配資料,並在不要求結果順序時匯出完整結果集。

前提條件

安裝Tablestore Java SDK並初始化用戶端。

並發匯出資料要求使用 5.6.0 及以上版本的 Tablestore Java SDK。版本資訊請參見版本歷史

功能說明

並發匯出資料用於全量掃描多元索引中滿足查詢條件的資料。掃描結果不保證全域順序,且不支援排序和統計彙總。如果需要對結果排序、進行統計彙總或面向終端使用者返回檢索結果,請使用 Search 介面。

單並發掃描配置簡單;多並發掃描可同時讀取多個分區,通常能夠獲得比單並發更高的掃描輸送量。

完整的並發掃描流程如下:

  1. 調用 computeSplits 擷取多元索引的最大並發度 splitsSize 和任務會話標識 sessionId

  2. 配置 ParallelScanRequest。單並發掃描可省略 maxParallelcurrentParallelId;多並發掃描時,各任務使用相同的查詢條件、sessionIdmaxParallel,並分別設定不同的 currentParallelId

  3. 調用 createParallelScanIterator 自動讀取全部分頁資料,或調用 parallelScan 並使用 nextToken 手動翻頁。

  4. 等待所有並發任務完成,併合並各任務的無序結果。

同一 sessionId 下的掃描任務在首次調用 parallelScan 時確定資料快照。任務運行期間對資料的新增或更新不會進入該快照。sessionId 可省略,但掃描期間服務端發生負載平衡等變化時,結果可能包含少量重複資料,因此建議先調用 computeSplits 並在後續請求中攜帶返回的 sessionId

重要

動態修改 Schema 觸發切換索引、服務端容錯移轉或負載平衡等操作時,會話可能提前失效,服務端返回 OTSSessionExpired;用戶端網路異常也可能中斷掃描。遇到此類異常時,請丟棄當前任務的不完整結果,重新調用 computeSplits,並從頭重啟整個掃描任務。同一多元索引最多同時運行 10 個並發掃描任務,其他限制請參見多元索引使用限制

調用 computeSplits 計算分區,通過 parallelScan 手動翻頁掃描資料,或通過 createParallelScanIterator 自動讀取全部分頁。

ComputeSplitsResponse computeSplits(ComputeSplitsRequest request)
ParallelScanResponse parallelScan(ParallelScanRequest request)
RowIterator createParallelScanIterator(ParallelScanRequest request)

以下樣本以單並發方式掃描全部資料,並返回 categoryprice 欄位。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 返回的任務會話標識。建議設定,以保證掃描期間使用同一資料快照。

timeoutInMillisecond(可選)

int

請求級逾時時間,單位為毫秒。預設值為 -1,表示不單獨佈建要求逾時時間。

掃描配置

request.scanQuery 的類型為 ScanQuery,包含以下參數。

名稱

類型

說明

query(必選)

Query

掃描範圍對應的查詢條件。支援精確查詢、匹配查詢、範圍查詢、地理位置查詢和巢狀型別查詢等,查詢條件的配置方式與 Search 介面相同。掃描多元索引中的全部資料時,設定為 MatchAllQuery

limit(可選)

Integer

單次請求最多返回的行數,預設值為 2000,建議保持預設值。

maxParallel(可選)

Integer

掃描任務的並發度,不能大於 ComputeSplitsResponse.splitsSize,預設值為 1。

currentParallelId(可選)

Integer

當前並發任務 ID。maxParallel 大於 1 時必須設定,各任務的取值範圍為 [0, maxParallel) 且不能重複。

aliveTime(可選)

Integer

掃描任務在兩次分頁請求之間的最大有效時間,單位為秒,取值範圍為 1~600,預設值為 60。每次成功擷取資料後會重新整理有效時間。

token(可選)

byte[]

分頁憑證。首次請求不設定;手動翻頁時設定為上一次響應的 nextToken。使用 createParallelScanIterator 時由 SDK 自動管理。

說明

服務端允許將 limit 設定為最大 10000,但不建議使用該上限。

返回列配置

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[]

任務會話標識,用於在同一資料快照中掃描資料。

splitsSize

Integer

多元索引支援的最大並發度。

掃描結果

parallelScan 返回 ParallelScanResponse,包含以下欄位。

名稱

類型

說明

rows

List<Row>

本次請求返回的資料行。

nextToken

byte[]

下一頁憑證。值為 null 時,當前並發任務已讀取完畢。

bodyBytes

long

本次響應體的位元組數。

createParallelScanIterator 返回 RowIterator。該迭代器自動使用 nextToken 擷取後續分頁,每次迭代返回一個 Row,不支援擷取匹配總行數。

情境樣本

多並發掃描

以下樣本根據 splitsSize 建立多個掃描任務。每個任務使用唯一的 currentParallelId,全部任務共用同一 sessionIdmaxParallel。線程池大小不超過用戶端的 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);