并发导出数据

更新时间:
复制 MD 格式

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