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