消费增量数据

更新时间:
复制 MD 格式

Go SDK 可通过 Stream API 消费数据表的写入、更新和删除等增量变更数据。

前提条件

  • 安装Tablestore Go SDK并初始化客户端。

  • 数据表已开启 Stream。可在创建数据表时通过 StreamSpec 开启,也可通过 UpdateTable 更新 Stream 配置。

功能说明

Stream 将数据表的增量变更按 Shard 组织。消费流程包含以下步骤:

  1. 调用 ListStream 获取目标数据表的 Stream ID。

  2. 调用 DescribeStream 获取 Stream 状态和 Shard 列表。

  3. 为每个 Shard 调用 GetShardIterator 获取读取迭代值。

  4. 调用 GetStreamRecord 拉取增量记录,并使用返回的 NextShardIterator 继续读取。

func (client *TableStoreClient) ListStream(req *ListStreamRequest) (*ListStreamResponse, error)
func (client *TableStoreClient) DescribeStream(req *DescribeStreamRequest) (*DescribeStreamResponse, error)
func (client *TableStoreClient) GetShardIterator(req *GetShardIteratorRequest) (*GetShardIteratorResponse, error)
func (client TableStoreClient) GetStreamRecord(req *GetStreamRecordRequest) (*GetStreamRecordResponse, error)

以下示例获取 example_table 的第一个 Shard,并读取一批增量记录。实际应用中应遍历并持续消费所有 Shard。

tableName := "example_table"
listResponse, err := client.ListStream(
    &tablestore.ListStreamRequest{TableName: &tableName},
)
if err != nil {
    log.Fatal(err)
}
streamID := listResponse.Streams[0].Id

describeResponse, err := client.DescribeStream(
    &tablestore.DescribeStreamRequest{StreamId: streamID},
)
if err != nil {
    log.Fatal(err)
}
shardID := describeResponse.Shards[0].SelfShard

iteratorResponse, err := client.GetShardIterator(
    &tablestore.GetShardIteratorRequest{
        StreamId: streamID,
        ShardId:  shardID,
    },
)
if err != nil {
    log.Fatal(err)
}

recordResponse, err := client.GetStreamRecord(
    &tablestore.GetStreamRecordRequest{
        ShardIterator: iteratorResponse.ShardIterator,
    },
)
if err != nil {
    log.Fatal(err)
}

for _, record := range recordResponse.Records {
    fmt.Println(record)
}
nextIterator := recordResponse.NextShardIterator
fmt.Println(nextIterator)

参数说明

列出 Stream

ListStreamRequest 包含以下参数。

名称

类型

说明

TableName(可选)

*string

数据表名称。未设置时返回当前实例中所有已开启 Stream 的数据表;设置后仅返回指定表。

查询 Stream

DescribeStreamRequest 包含以下参数。

名称

类型

说明

StreamId(必选)

*StreamId

Stream ID,由 ListStream 返回。

InclusiveStartShardId(可选)

*ShardId

本次返回 Shard 列表的起始 Shard ID。

ShardLimit(可选)

*int32

本次最多返回的 Shard 数量。

获取 Shard 读取迭代值

GetShardIteratorRequest 包含以下参数。

名称

类型

说明

StreamId(必选)

*StreamId

Stream ID。

ShardId(必选)

*ShardId

Shard ID,由 DescribeStream 返回。

Timestamp(可选)

*int64

读取起点的时间戳,单位为微秒。未设置时从 Shard 起始位置读取。

Token(可选)

*string

用于继续获取读取迭代值的令牌。

读取增量记录

GetStreamRecordRequest 包含以下参数。

名称

类型

说明

ShardIterator(必选)

*ShardIterator

Shard 读取迭代值,由 GetShardIterator 或上一次 GetStreamRecord 返回。

Limit(可选)

*int32

本次最多返回的增量记录数量。

TableName(可选)

*string

Shard 所属的数据表名称。

返回值

Stream 列表

ListStreamResponse.Streams 的类型为 []Stream。每个元素包含 Stream ID、数据表名称和创建时间。

Stream 信息

DescribeStreamResponse 包含以下业务信息。

字段

类型

说明

StreamId

*StreamId

Stream ID。

TableName

*string

数据表名称。

CreationTime

int64

Stream 创建时间,单位为微秒。

ExpirationTime

int32

增量日志保留时长,单位为小时。

Status

StreamStatus

Stream 状态。

Shards

[]*StreamShard

本次返回的 Shard 列表。

NextShardId

*ShardId

下一次查询的起始 Shard ID。值为 nil 时表示已返回全部 Shard。

Shard 读取迭代值

GetShardIteratorResponse.ShardIterator 是首次读取指定 Shard 时使用的迭代值。

增量记录

GetStreamRecordResponse 包含以下业务信息。

字段

类型

说明

Records

[]*StreamRecord

本次返回的增量记录。

NextShardIterator

*ShardIterator

下一次读取使用的迭代值。值为 nil 时表示当前 Shard 已读取完。

MayMoreRecord

*bool

当前 Shard 是否可能还有更多增量记录。

场景示例

持续读取增量数据

以下示例使用 NextShardIterator 持续读取一个 Shard。

iterator := iteratorResponse.ShardIterator
for iterator != nil {
    response, err := client.GetStreamRecord(
        &tablestore.GetStreamRecordRequest{ShardIterator: iterator},
    )
    if err != nil {
        log.Fatal(err)
    }

    for _, record := range response.Records {
        fmt.Println(record)
    }
    iterator = response.NextShardIterator
}