Go SDK 可通过 Stream API 消费数据表的写入、更新和删除等增量变更数据。
前提条件
安装Tablestore Go SDK并初始化客户端。
数据表已开启 Stream。可在创建数据表时通过
StreamSpec开启,也可通过UpdateTable更新 Stream 配置。
功能说明
Stream 将数据表的增量变更按 Shard 组织。消费流程包含以下步骤:
调用
ListStream获取目标数据表的 Stream ID。调用
DescribeStream获取 Stream 状态和 Shard 列表。为每个 Shard 调用
GetShardIterator获取读取迭代值。调用
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(可选) |
|
数据表名称。未设置时返回当前实例中所有已开启 Stream 的数据表;设置后仅返回指定表。 |
查询 Stream
DescribeStreamRequest 包含以下参数。
|
名称 |
类型 |
说明 |
|
StreamId(必选) |
|
Stream ID,由 |
|
InclusiveStartShardId(可选) |
|
本次返回 Shard 列表的起始 Shard ID。 |
|
ShardLimit(可选) |
|
本次最多返回的 Shard 数量。 |
获取 Shard 读取迭代值
GetShardIteratorRequest 包含以下参数。
|
名称 |
类型 |
说明 |
|
StreamId(必选) |
|
Stream ID。 |
|
ShardId(必选) |
|
Shard ID,由 |
|
Timestamp(可选) |
|
读取起点的时间戳,单位为微秒。未设置时从 Shard 起始位置读取。 |
|
Token(可选) |
|
用于继续获取读取迭代值的令牌。 |
读取增量记录
GetStreamRecordRequest 包含以下参数。
|
名称 |
类型 |
说明 |
|
ShardIterator(必选) |
|
Shard 读取迭代值,由 |
|
Limit(可选) |
|
本次最多返回的增量记录数量。 |
|
TableName(可选) |
|
Shard 所属的数据表名称。 |
返回值
Stream 列表
ListStreamResponse.Streams 的类型为 []Stream。每个元素包含 Stream ID、数据表名称和创建时间。
Stream 信息
DescribeStreamResponse 包含以下业务信息。
|
字段 |
类型 |
说明 |
|
|
|
Stream ID。 |
|
|
|
数据表名称。 |
|
|
|
Stream 创建时间,单位为微秒。 |
|
|
|
增量日志保留时长,单位为小时。 |
|
|
|
Stream 状态。 |
|
|
|
本次返回的 Shard 列表。 |
|
|
|
下一次查询的起始 Shard ID。值为 |
Shard 读取迭代值
GetShardIteratorResponse.ShardIterator 是首次读取指定 Shard 时使用的迭代值。
增量记录
GetStreamRecordResponse 包含以下业务信息。
|
字段 |
类型 |
说明 |
|
|
|
本次返回的增量记录。 |
|
|
|
下一次读取使用的迭代值。值为 |
|
|
|
当前 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
}