并发导出数据

更新时间:
复制 MD 格式

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

下一页凭证。值非空时继续扫描当前分片。