使用 Tablestore Go SDK 将多元索引数据划分为多个分片并发扫描,以提高大规模数据导出效率。
前提条件
开始前,完成以下准备工作:
安装Tablestore Go SDK并初始化客户端。
并发导出数据需要 Go SDK 1.6.0 及以上版本,建议使用最新版本。
功能说明
并发导出先调用 ComputeSplits 创建扫描会话并获取建议并发数,再为每个并发任务调用 ParallelScan。每个并发任务使用独立的 CurrentParallelID,并通过 NextToken 连续读取该分片。并发扫描不保证整个结果集的顺序,适合不依赖返回顺序的大规模导出。也可以将 MaxParallel 设置为 1、CurrentParallelID 设置为 0 进行单并发扫描;代码更简单,吞吐通常高于 Search,但低于多并发扫描。
并发扫描不支持排序和统计聚合。如果需要对结果排序、执行统计聚合或面向最终用户返回检索结果,请使用 Search 接口。
同一 SessionId 下的扫描任务在首次调用 ParallelScan 时确定数据快照,任务运行期间新增或更新的数据不进入该快照。SessionId 可以省略,但服务端负载均衡等变化可能使结果包含少量重复数据,因此建议先调用 ComputeSplits 并在后续请求中携带返回的 SessionId。
同一多元索引最多同时运行 10 个并发扫描任务,其他限制请参见多元索引使用限制。
以下示例根据 ComputeSplits 返回的建议并发数启动多个 goroutine,并为每个并发任务设置不同的 CurrentParallelID,直至读取完所有分片。
splits, err := client.ComputeSplits(
(&tablestore.ComputeSplitsRequest{}).
SetTableName("example_table").
SetSearchIndexSplitsOptions(tablestore.SearchIndexSplitsOptions{
IndexName: "example_index",
}),
)
if err != nil {
log.Fatal(err)
}
var waitGroup sync.WaitGroup
var mutex sync.Mutex
totalRows := 0
errors := make(chan error, splits.SplitsSize)
waitGroup.Add(int(splits.SplitsSize))
for workerID := int32(0); workerID < splits.SplitsSize; workerID++ {
currentWorkerID := workerID
go func() {
defer waitGroup.Done()
scanQuery := search.NewScanQuery().
SetQuery(&search.MatchAllQuery{}).
SetLimit(1000).
SetMaxParallel(splits.SplitsSize).
SetCurrentParallelID(currentWorkerID)
request := (&tablestore.ParallelScanRequest{}).
SetTableName("example_table").
SetIndexName("example_index").
SetScanQuery(scanQuery).
SetSessionId(splits.SessionId).
SetColumnsToGet(&tablestore.ColumnsToGet{
ReturnAllFromIndex: true,
})
for {
response, err := client.ParallelScan(request)
if err != nil {
errors <- err
return
}
// Process response.Rows here.
mutex.Lock()
totalRows += len(response.Rows)
mutex.Unlock()
if len(response.NextToken) == 0 {
return
}
request.SetScanQuery(scanQuery.SetToken(response.NextToken))
}
}()
}
waitGroup.Wait()
close(errors)
for err := range errors {
log.Fatal(err)
}
fmt.Println(totalRows)
参数说明
创建扫描会话
|
名称 |
类型 |
说明 |
|
TableName(必选) |
string |
数据表名称。 |
|
IndexName(必选) |
string |
多元索引名称。 |
并发扫描请求
|
名称 |
类型 |
说明 |
|
TableName(必选) |
string |
数据表名称。 |
|
IndexName(必选) |
string |
多元索引名称。 |
|
ScanQuery(必选) |
search.ScanQuery |
扫描条件和并发配置。 |
|
SessionId(可选) |
[]byte |
ComputeSplits 返回的会话 ID。建议设置,以保证扫描期间使用同一数据快照。 |
|
ColumnsToGet(可选) |
*tablestore.ColumnsToGet |
返回列配置。未设置时只返回主键列。并发扫描不能使用 ReturnAll。 |
|
TimeoutMs(可选) |
*int32 |
请求超时时间,单位为毫秒。 |
扫描配置
|
名称 |
类型 |
说明 |
|
Query(必选) |
search.Query |
扫描条件,支持与 Search 一致的非向量查询类型。 |
|
Limit(可选) |
int32 |
单次请求返回的最大行数,默认值为 2000,建议保持默认值。服务端允许设置为最大 10000,但较大值会增加单次请求延迟和资源占用。 |
|
MaxParallel(可选) |
int32 |
并发任务总数,默认值为 1,不能超过 ComputeSplits 返回的 SplitsSize。 |
|
CurrentParallelID(可选) |
int32 |
当前并发任务编号。MaxParallel 大于 1 时必须设置,各任务的取值范围为 [0, MaxParallel) 且不能重复。 |
|
Token(可选) |
[]byte |
上一页响应中的 NextToken。 |
|
AliveTime(可选) |
int32 |
两次分页请求之间允许的最大间隔,单位为秒,取值范围为 1~600,默认值为 60,建议保持默认值。每次成功获取数据后会刷新有效时间。 |
返回列配置
|
名称 |
类型 |
说明 |
|
Columns(可选) |
[]string |
要返回的多元索引字段名称。仅存在于数据表但未加入多元索引的字段不能返回。 |
|
ReturnAllFromIndex(可选) |
bool |
是否返回多元索引中的全部字段,默认值为 false。设置为 true 时无需设置 Columns。 |
|
ReturnAll(可选) |
bool |
并发扫描不支持该参数,请勿设置为 true。 |
动态修改 Schema 触发切换索引、服务端故障转移或负载均衡等操作时,会话可能提前失效并返回 OTSSessionExpired;客户端网络异常也可能中断扫描。遇到此类异常时,丢弃当前任务的不完整结果,重新调用 ComputeSplits,并从头重启整个扫描任务。
返回值
分片信息
|
名称 |
类型 |
说明 |
|
SessionId |
[]byte |
任务会话标识,用于在同一数据快照中扫描数据。 |
|
SplitsSize |
int32 |
多元索引支持的最大并发度。 |
扫描结果
|
名称 |
类型 |
说明 |
|
Rows |
[]*tablestore.Row |
本次扫描返回的行。 |
|
NextToken |
[]byte |
下一页凭证。值非空时继续扫描当前分片。 |