并发扫描时序数据

更新时间:
复制 MD 格式

使用 Java SDK 将时序数据扫描任务切分为多个分片并发执行,以快速导出时序表中的大量数据。

前提条件

安装Tablestore Java SDK并初始化时序模型客户端。

功能说明

并发扫描时,先调用 splitTimeseriesScanTask 方法将扫描范围切分为多个相互独立的分片,再为每个分片调用 scanTimeseriesData 方法,并发读取分片中的时序数据。splitCountHint 仅表示期望的分片数,实际分片数由服务端决定。

public SplitTimeseriesScanTaskResponse splitTimeseriesScanTask(SplitTimeseriesScanTaskRequest request) throws TableStoreException, ClientException
public ScanTimeseriesDataResponse scanTimeseriesData(ScanTimeseriesDataRequest request) throws TableStoreException, ClientException

以下示例将 example_timeseries_table 中度量名称为 cpu 的时序数据切分为期望的 4 个分片,使用并行流扫描各分片,并通过 nextToken 读取每个分片的全部数据。

TimeseriesClient timeseriesClient = client.asTimeseriesClient();
String tableName = "example_timeseries_table";

SplitTimeseriesScanTaskRequest splitRequest =
        new SplitTimeseriesScanTaskRequest(tableName, "cpu", 4);
SplitTimeseriesScanTaskResponse splitResponse =
        timeseriesClient.splitTimeseriesScanTask(splitRequest);

List<TimeseriesRow> rows = splitResponse.getSplitInfos().parallelStream()
        .flatMap(splitInfo -> {
            ScanTimeseriesDataRequest scanRequest =
                    new ScanTimeseriesDataRequest(tableName);
            scanRequest.setSplitInfo(splitInfo);
            scanRequest.setLimit(5000);

            List<TimeseriesRow> splitRows = new ArrayList<>();
            do {
                ScanTimeseriesDataResponse scanResponse =
                        timeseriesClient.scanTimeseriesData(scanRequest);
                splitRows.addAll(scanResponse.getRows());
                scanRequest.setNextToken(scanResponse.getNextToken());
            } while (scanRequest.getNextToken() != null);
            return splitRows.stream();
        })
        .collect(Collectors.toList());

System.out.println(rows.size());
重要

扫描过程中可能返回空数据页,但响应中仍包含非空的 nextToken。请始终以 nextToken 是否为 null 判断当前分片是否扫描完成。

参数说明

切分扫描任务

SplitTimeseriesScanTaskRequest 包含以下参数。

名称

类型

说明

timeseriesTableName(必选)

String

时序表名称。

splitCountHint(必选)

int

期望的分片数,取值大于 0。实际返回的分片数由服务端决定。

measurementName(可选)

String

要扫描的度量名称。

扫描分片

ScanTimeseriesDataRequest 包含以下参数。

名称

类型

说明

timeseriesTableName(必选)

String

时序表名称,必须与切分扫描任务时使用的名称一致。

splitInfo(必选)

TimeseriesScanSplitInfo

通过 SplitTimeseriesScanTaskResponse.getSplitInfos() 获取的分片信息。每个并发任务使用一个不同的分片。

beginTimeInUs(可选)

long

扫描范围的起始时间,格式为从 1970-01-01 00:00:00 UTC 开始计算的微秒时间戳,取值大于等于 0。调用 setTimeRange 时与 endTimeInUs 一起设置。

endTimeInUs(可选)

long

扫描范围的结束时间,格式为微秒时间戳,取值大于 beginTimeInUs。调用 setTimeRange 时与 beginTimeInUs 一起设置。

fieldsToGet(可选)

List<Pair<String, ColumnType>>

要返回的数据字段名称和类型。不设置时返回全部数据字段。

limit(可选)

int

单次请求的最大返回行数。默认值和最大值均为 5000

nextToken(可选)

byte[]

当前分片的分页令牌。首次扫描分片时不设置;响应中的 nextToken 不为空时,将其传入该分片的下一次请求。

返回值

切分结果

SplitTimeseriesScanTaskResponse 包含以下业务返回字段。

字段

类型

说明

splitInfos

List<TimeseriesScanSplitInfo>

通过 getSplitInfos() 获取服务端生成的分片信息。列表中的每个元素用于一次独立的分片扫描。

扫描结果

ScanTimeseriesDataResponse 包含以下业务返回字段。

字段

类型

说明

rows

List<TimeseriesRow>

通过 getRows() 获取当前分片本次请求返回的时序数据行。

nextToken

byte[]

通过 getNextToken() 获取当前分片的下一页令牌。值为 null 时,表示该分片已扫描完成。

时序数据行

rows[] 中每个元素的类型为 TimeseriesRow,包含以下字段。

字段

类型

说明

timeseriesKey

TimeseriesKey

通过 getTimeseriesKey() 获取时间线标识。

timeInUs

long

通过 getTimeInUs() 获取数据点时间,单位为微秒。

fields

SortedMap<String, ColumnValue>

通过 getFields() 获取数据字段。