使用 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 |
下一頁憑證。值非空時繼續掃描當前分區。 |