使用 Tablestore Java SDK 可並發掃描多元索引中的匹配資料,並在不要求結果順序時匯出完整結果集。
前提條件
安裝Tablestore Java SDK並初始化用戶端。
並發匯出資料要求使用 5.6.0 及以上版本的 Tablestore Java SDK。版本資訊請參見版本歷史。
功能說明
並發匯出資料用於全量掃描多元索引中滿足查詢條件的資料。掃描結果不保證全域順序,且不支援排序和統計彙總。如果需要對結果排序、進行統計彙總或面向終端使用者返回檢索結果,請使用 Search 介面。
單並發掃描配置簡單;多並發掃描可同時讀取多個分區,通常能夠獲得比單並發更高的掃描輸送量。
完整的並發掃描流程如下:
調用
computeSplits擷取多元索引的最大並發度splitsSize和任務會話標識sessionId。配置
ParallelScanRequest。單並發掃描可省略maxParallel和currentParallelId;多並發掃描時,各任務使用相同的查詢條件、sessionId和maxParallel,並分別設定不同的currentParallelId。調用
createParallelScanIterator自動讀取全部分頁資料,或調用parallelScan並使用nextToken手動翻頁。等待所有並發任務完成,併合並各任務的無序結果。
同一 sessionId 下的掃描任務在首次調用 parallelScan 時確定資料快照。任務運行期間對資料的新增或更新不會進入該快照。sessionId 可省略,但掃描期間服務端發生負載平衡等變化時,結果可能包含少量重複資料,因此建議先調用 computeSplits 並在後續請求中攜帶返回的 sessionId。
動態修改 Schema 觸發切換索引、服務端容錯移轉或負載平衡等操作時,會話可能提前失效,服務端返回 OTSSessionExpired;用戶端網路異常也可能中斷掃描。遇到此類異常時,請丟棄當前任務的不完整結果,重新調用 computeSplits,並從頭重啟整個掃描任務。同一多元索引最多同時運行 10 個並發掃描任務,其他限制請參見多元索引使用限制。
調用 computeSplits 計算分區,通過 parallelScan 手動翻頁掃描資料,或通過 createParallelScanIterator 自動讀取全部分頁。
ComputeSplitsResponse computeSplits(ComputeSplitsRequest request)
ParallelScanResponse parallelScan(ParallelScanRequest request)
RowIterator createParallelScanIterator(ParallelScanRequest request)
以下樣本以單並發方式掃描全部資料,並返回 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 |
掃描範圍對應的查詢條件。支援精確查詢、匹配查詢、範圍查詢、地理位置查詢和巢狀型別查詢等,查詢條件的配置方式與 Search 介面相同。掃描多元索引中的全部資料時,設定為 |
|
limit(可選) |
Integer |
單次請求最多返回的行數,預設值為 2000,建議保持預設值。 |
|
maxParallel(可選) |
Integer |
掃描任務的並發度,不能大於 |
|
currentParallelId(可選) |
Integer |
當前並發任務 ID。 |
|
aliveTime(可選) |
Integer |
掃描任務在兩次分頁請求之間的最大有效時間,單位為秒,取值範圍為 1~600,預設值為 60。每次成功擷取資料後會重新整理有效時間。 |
|
token(可選) |
|
分頁憑證。首次請求不設定;手動翻頁時設定為上一次響應的 |
服務端允許將 limit 設定為最大 10000,但不建議使用該上限。
返回列配置
request.columnsToGet 的類型為 SearchRequest.ColumnsToGet,包含以下參數。
|
名稱 |
類型 |
說明 |
|
columns(可選) |
|
要返回的多元索引欄位名稱列表。僅資料表中存在但未加入多元索引的欄位不能返回。Date、Geo-point、IP、Vector、JSON/Nested 和數組欄位均可返回。 |
|
returnAllFromIndex(可選) |
boolean |
是否返回多元索引中的全部欄位,預設值為 |
|
returnAll(可選) |
boolean |
並發掃描不支援該參數,請勿設定為 |
傳回值
分區資訊
computeSplits 返回 ComputeSplitsResponse,包含以下欄位。
|
名稱 |
類型 |
說明 |
|
sessionId |
|
任務會話標識,用於在同一資料快照中掃描資料。 |
|
splitsSize |
Integer |
多元索引支援的最大並發度。 |
掃描結果
parallelScan 返回 ParallelScanResponse,包含以下欄位。
|
名稱 |
類型 |
說明 |
|
rows |
|
本次請求返回的資料行。 |
|
nextToken |
|
下一頁憑證。值為 |
|
bodyBytes |
long |
本次響應體的位元組數。 |
createParallelScanIterator 返回 RowIterator。該迭代器自動使用 nextToken 擷取後續分頁,每次迭代返回一個 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);