接口调用示例
接口调用示例
以下"返回示例"给出的是响应结构体的字段结构,字段名与类型来自 SDK 模型定义,具体取值需以实际调用结果为准。示例均属于
kafkademo包,共享新建 Kafka Client中定义的newKafkaClient()。
集群管理
CreateCluster 创建集群
功能说明
创建 Kafka 集群。方法只校验 Name 非空,Mode、Type 与 Provisioned 内部字段的合法组合由服务端校验。创建预置型集群时应完整填写 Provisioned,其中不带 omitempty 的值类型字段需显式赋值。
创建集群会产生计费。使用 NewBilling() 构造时 IsAutoPay 默认为 true,会直接扣费。
请求示例
1package kafkademo
2
3import (
4 "fmt"
5 "log"
6
7 "github.com/baidubce/bce-sdk-go/services/kafka"
8)
9
10func createCluster() {
11 client := newKafkaClient()
12
13 billing := kafka.NewBilling()
14 billing.Payment = "Prepaid"
15 billing.TimeLength = 1
16 billing.AutoRenewEnabled = false
17
18 storageMeta := kafka.NewStorageMeta()
19 storageMeta.StorageType = kafka.StorageTypeEnhancedSSDPL1
20 storageMeta.StorageSize = 100
21
22 configMeta := kafka.NewConfigMeta()
23 configMeta.Context["auto.create.topics.enable"] = "true"
24 configMeta.Context["log.retention.hours"] = "168"
25
26 request := &kafka.CreateClusterRequest{
27 Name: "kafka-demo",
28 Mode: kafka.ModeHP,
29 Type: kafka.TypeProvisioned,
30 Tags: []kafka.Tag{
31 {TagKey: "KAFKA-Cluster", TagValue: "prod"},
32 },
33 Provisioned: &kafka.Provisioned{
34 KafkaVersion: "2.7.2",
35 NodeType: "kafka.g4.c2m8",
36 NumberOfBrokerNodes: 3,
37 StorageMeta: storageMeta,
38 VPCID: "{{VPC ID}}",
39 SubnetIds: []string{"{{子网 ID}}"},
40 SecurityGroupIds: []string{"{{安全组 ID}}"},
41 PublicIPEnabled: false,
42 IntranetIPEnabled: false,
43 ACLEnabled: true,
44 DeploySetEnabled: true,
45 Authentications: []kafka.Authentication{
46 {Mode: string(kafka.AuthModeSASLSCRAM)},
47 },
48 Billing: billing,
49 ConfigMeta: configMeta,
50 },
51 }
52
53 response, err := client.CreateCluster(request)
54 if err != nil {
55 log.Fatalf("failed to create cluster: %v", err)
56 }
57 fmt.Printf("clusterId=%s
58", response.ClusterID)
59}
返回示例
CreateClusterResponse 字段结构:
1ClusterID : string 集群 ID
集群创建是异步过程,返回 ClusterID 只表示创建请求已受理,需继续查询集群状态确认就绪。
相关 API
GetClusterDetail 查询集群详情
功能说明
按 ClusterID 查询集群详情。响应中的 ClusterDetail.Provisioned 包含规格、网络、存储、认证和维护窗口等完整配置。
查询已删除集群请使用 GetClusterDeletion,其请求类型为 GetClusterDeletionRequest,但返回类型复用 *GetClusterDetailResponse。
GetClusterDetailResponse.Cluster 是指针,使用前必须判空。
请求示例
1package kafkademo
2
3import (
4 "fmt"
5 "log"
6
7 "github.com/baidubce/bce-sdk-go/services/kafka"
8)
9
10func getClusterDetail() {
11 client := newKafkaClient()
12
13 request := &kafka.GetClusterDetailRequest{
14 ClusterID: "{{集群 ID}}",
15 }
16
17 response, err := client.GetClusterDetail(request)
18 if err != nil {
19 log.Fatalf("failed to get cluster detail: %v", err)
20 }
21
22 detail := response.Cluster
23 if detail == nil {
24 log.Println("cluster not found")
25 return
26 }
27 fmt.Printf("name=%s state=%s region=%s
28",
29 detail.Name, detail.State, detail.Region)
30
31 if detail.Provisioned != nil {
32 fmt.Printf("kafkaVersion=%s nodeType=%s brokers=%d
33",
34 detail.Provisioned.KafkaVersion,
35 detail.Provisioned.NodeType,
36 detail.Provisioned.NumberOfBrokerNodes)
37 }
38}
返回示例
GetClusterDetailResponse.Cluster(*ClusterDetail)字段结构:
1ClusterID : string 集群 ID
2ClusterSid : string 集群短 ID
3Name : string 集群名称
4Region : string 地域
5Type : string 集群类型
6Mode : string 集群模式
7VPCMode : string VPC 资源管理模式
8State : string 集群状态
9Provisioned : *Provisioned 预置模式配置
10Tags : []Tag 标签列表
11CreateTime : string 创建时间
12DeleteTime : string 删除时间
相关 API
StopCluster 停止集群
功能说明
停止运行中的集群。停止后集群不再提供服务,客户端连接会中断,属于影响业务的操作,执行前需确认目标集群。
接口返回 ActionID,需通过任务接口确认最终状态。
请求示例
1package kafkademo
2
3import (
4 "fmt"
5 "log"
6
7 "github.com/baidubce/bce-sdk-go/services/kafka"
8)
9
10func stopCluster() {
11 client := newKafkaClient()
12
13 request := &kafka.StopClusterRequest{
14 ClusterID: "{{集群 ID}}",
15 }
16
17 response, err := client.StopCluster(request)
18 if err != nil {
19 log.Fatalf("failed to stop cluster: %v", err)
20 }
21 fmt.Printf("clusterId=%s actionId=%s
22",
23 response.ClusterID, response.ActionID)
24}
返回示例
StopClusterResponse 字段结构:
1ClusterID : string 集群 ID
2ActionID : string 任务 ID,用于后续查询任务状态
相关 API
主题管理
CreateTopic 创建主题
功能说明
在指定集群下创建主题。方法只校验 ClusterID 与 TopicName 非空。
PartitionNum 与 ReplicationFactor 是不带 omitempty 的 int,零值也会发送,必须显式赋值。
可用的主题级参数及取值范围通过 ListTopicConfigOptions() 查询,该方法不接受任何参数。
请求示例
1package kafkademo
2
3import (
4 "fmt"
5 "log"
6
7 "github.com/baidubce/bce-sdk-go/services/kafka"
8)
9
10func createTopic() {
11 client := newKafkaClient()
12
13 request := &kafka.CreateTopicRequest{
14 ClusterID: "{{集群 ID}}",
15 TopicName: "{{主题名称}}",
16 PartitionNum: 3,
17 ReplicationFactor: 3,
18 OtherConfigs: map[string]string{
19 "cleanup.policy": "delete",
20 "retention.ms": "604800000",
21 },
22 }
23
24 response, err := client.CreateTopic(request)
25 if err != nil {
26 log.Fatalf("failed to create topic: %v", err)
27 }
28 fmt.Printf("topicName=%s
29", response.TopicName)
30}
返回示例
CreateTopicResponse 字段结构:
1TopicName : string 主题名称
相关 API
UpdateTopic 变更主题
功能说明
变更主题分区数或主题级配置。PartitionNum 与 OtherConfigs 至少设置一项,否则返回 request partitionNum and otherConfigs should not be both empty。分区数只能增加,不能减少。
注意 UpdateTopicRequest.PartitionNum 的类型是 string,而 CreateTopicRequest.PartitionNum 是 int,容易混淆。传入整数时需先用 strconv.Itoa 转换,或直接写字符串字面量。
请求示例
1package kafkademo
2
3import (
4 "fmt"
5 "log"
6 "strconv"
7
8 "github.com/baidubce/bce-sdk-go/services/kafka"
9)
10
11func updateTopic() {
12 client := newKafkaClient()
13
14 request := &kafka.UpdateTopicRequest{
15 ClusterID: "{{集群 ID}}",
16 TopicName: "{{主题名称}}",
17 PartitionNum: strconv.Itoa(6),
18 OtherConfigs: map[string]string{
19 "retention.ms": "1209600000",
20 },
21 }
22
23 response, err := client.UpdateTopic(request)
24 if err != nil {
25 log.Fatalf("failed to update topic: %v", err)
26 }
27 fmt.Printf("topicName=%s
28", response.TopicName)
29}
返回示例
UpdateTopicResponse 字段结构:
1TopicName : string 主题名称
相关 API
ListTopic 查询主题列表
功能说明
查询集群下的主题列表。TopicName 可选,按前缀匹配过滤,其 json tag 为 -,由方法内部转成查询参数。该接口不分页,返回集群下符合条件的全部主题。
Topic 中的 PartitionNum、ReplicaNum、ReadOnly 都是指针类型,读取前需判空。
请求示例
1package kafkademo
2
3import (
4 "fmt"
5 "log"
6
7 "github.com/baidubce/bce-sdk-go/services/kafka"
8)
9
10func listTopic() {
11 client := newKafkaClient()
12
13 request := &kafka.ListTopicRequest{
14 ClusterID: "{{集群 ID}}",
15 TopicName: "order-",
16 }
17
18 response, err := client.ListTopic(request)
19 if err != nil {
20 log.Fatalf("failed to list topics: %v", err)
21 }
22
23 for _, topic := range response.Topics {
24 partitionNum := 0
25 if topic.PartitionNum != nil {
26 partitionNum = *topic.PartitionNum
27 }
28 replicaNum := 0
29 if topic.ReplicaNum != nil {
30 replicaNum = *topic.ReplicaNum
31 }
32 fmt.Printf("topicName=%s partitionNum=%d replicaNum=%d
33",
34 topic.TopicName, partitionNum, replicaNum)
35 }
36}
返回示例
ListTopicResponse.Topics 元素(Topic)字段结构:
1TopicName : string 主题名称
2CreateTime : string 创建时间
3ReadOnly : *bool 是否只读
4PartitionNum : *int 分区数
5ReplicaNum : *int 副本数
相关 API
消费组管理
ListConsumerGroup 查询消费组列表
功能说明
查询集群下的消费组列表。GroupName 可选,用于过滤。
Group.GroupCoordinatorID 是 *int,读取前需判空。
请求示例
1package kafkademo
2
3import (
4 "fmt"
5 "log"
6
7 "github.com/baidubce/bce-sdk-go/services/kafka"
8)
9
10func listConsumerGroup() {
11 client := newKafkaClient()
12
13 request := &kafka.ListConsumerGroupRequest{
14 ClusterID: "{{集群 ID}}",
15 }
16
17 response, err := client.ListConsumerGroup(request)
18 if err != nil {
19 log.Fatalf("failed to list consumer groups: %v", err)
20 }
21
22 for _, group := range response.Groups {
23 coordinator := -1
24 if group.GroupCoordinatorID != nil {
25 coordinator = *group.GroupCoordinatorID
26 }
27 fmt.Printf("groupName=%s state=%s coordinator=%d
28",
29 group.GroupName, group.State, coordinator)
30 }
31}
返回示例
ListConsumerGroupResponse.Groups 元素(Group)字段结构:
1GroupName : string 消费组名称
2UpdateTime : string 更新时间
3State : string 消费组状态
4GroupCoordinatorID : *int Coordinator 所在 Broker ID
相关 API
ResetConsumerGroup 重置消费组位点
功能说明
重置消费组在指定主题分区上的消费位点。TopicName、Partitions、ResetStrategy 都是必填,Partitions 为空切片时返回 request partitions should not be null or empty。
重置位点会直接改变消费进度,可能导致消息重复消费或跳过,执行前必须确认目标消费组、主题和分区。建议在消费组无活跃成员时执行。
请求示例
1package kafkademo
2
3import (
4 "fmt"
5 "log"
6
7 "github.com/baidubce/bce-sdk-go/services/kafka"
8)
9
10func resetConsumerGroup() {
11 client := newKafkaClient()
12
13 request := &kafka.ResetConsumerGroupRequest{
14 ClusterID: "{{集群 ID}}",
15 GroupName: "{{消费组名称}}",
16 TopicName: "{{主题名称}}",
17 Partitions: []int{0, 1, 2},
18 ResetStrategy: "{{重置策略}}",
19 ResetValue: "{{重置目标值}}",
20 }
21
22 response, err := client.ResetConsumerGroup(request)
23 if err != nil {
24 log.Fatalf("failed to reset consumer group: %v", err)
25 }
26 fmt.Printf("groupName=%s
27", response.GroupName)
28}
返回示例
ResetConsumerGroupResponse 字段结构:
1GroupName : string 消费组名称
相关 API
用户管理
CreateUser 创建用户
功能说明
在集群下创建 SASL 用户。SDK 内部取客户端凭证 SK 的前 16 字节作为密钥,以 AES-128 ECB + PKCS5Padding 加密密码后按大写十六进制字符串发送,调用方传入原始密码即可。
SASLMechanisms 可选,取值示例为 SCRAM-SHA-256、SCRAM-SHA-512、PLAIN。
加密在发送前对请求结构体的副本上进行(payload := *request),调用方传入的请求对象不会被修改,因此同一个请求对象可以安全地重复提交,不存在密码被二次加密的问题。ResetUserPassword 采用相同实现。
不要记录原始密码、加密后的密码或完整请求体。
请求示例
1package kafkademo
2
3import (
4 "fmt"
5 "log"
6 "os"
7
8 "github.com/baidubce/bce-sdk-go/services/kafka"
9)
10
11func createUser() {
12 client := newKafkaClient()
13
14 request := &kafka.CreateUserRequest{
15 ClusterID: "{{集群 ID}}",
16 Username: "{{用户名}}",
17 Password: os.Getenv("KAFKA_USER_PASSWORD"),
18 SASLMechanisms: []string{"SCRAM-SHA-256", "SCRAM-SHA-512"},
19 }
20
21 response, err := client.CreateUser(request)
22 if err != nil {
23 log.Fatalf("failed to create user: %v", err)
24 }
25 fmt.Printf("username=%s
26", response.Username)
27}
返回示例
CreateUserResponse 字段结构:
1Username : string 用户名
相关 API
ListUsers 查询用户列表
功能说明
查询集群下的用户列表。ListUsers 的参数类型是 interface{},接受三种形式:
string:集群 IDkafka.ListUsersRequest:值类型*kafka.ListUsersRequest:指针类型
传入其他类型返回 request should be string or ListUsersRequest,传入 nil 指针返回 request should not be nil。ListAcls 采用同样的设计。
三种形式行为完全一致,建议统一使用请求结构体指针,便于后续追加过滤条件。
请求示例
1package kafkademo
2
3import (
4 "fmt"
5 "log"
6
7 "github.com/baidubce/bce-sdk-go/services/kafka"
8)
9
10func listUsers() {
11 client := newKafkaClient()
12
13 // 形式一:直接传集群 ID
14 response, err := client.ListUsers("{{集群 ID}}")
15 if err != nil {
16 log.Fatalf("failed to list users: %v", err)
17 }
18
19 // 形式二:传请求结构体指针
20 response, err = client.ListUsers(&kafka.ListUsersRequest{
21 ClusterID: "{{集群 ID}}",
22 })
23 if err != nil {
24 log.Fatalf("failed to list users: %v", err)
25 }
26
27 for _, user := range response.Users {
28 fmt.Printf("username=%s createTime=%s mechanisms=%v
29",
30 user.Username, user.CreateTime, user.SASLMechanisms)
31 }
32}
返回示例
ListUserResponse.Users 元素(User)字段结构:
1Username : string 用户名
2CreateTime : string 创建时间
3SASLMechanisms : []string 已启用的 SASL 机制
相关 API
权限管理
CreateAcl 创建权限
功能说明
为指定用户创建 ACL 授权。Username、PatternType、ResourceType、ResourceName 必填,Operations 为空切片时返回 request operation should not be null or empty。
注意 CreateAcl 使用 Operations([]string,可一次授予多个操作),而 DeleteAcl 使用 Operation(单个 string)。
CreateAclRequest 实现了 MarshalJSON,非 nil 的空 Operations 也会被发送。
请求示例
1package kafkademo
2
3import (
4 "fmt"
5 "log"
6
7 "github.com/baidubce/bce-sdk-go/services/kafka"
8)
9
10func createAcl() {
11 client := newKafkaClient()
12
13 request := &kafka.CreateAclRequest{
14 ClusterID: "{{集群 ID}}",
15 Username: "{{用户名}}",
16 PatternType: "{{匹配模式}}",
17 ResourceType: "{{资源类型}}",
18 ResourceName: "{{资源名称}}",
19 Operations: []string{"{{操作类型}}"},
20 }
21
22 response, err := client.CreateAcl(request)
23 if err != nil {
24 log.Fatalf("failed to create acl: %v", err)
25 }
26 fmt.Printf("username=%s
27", response.Username)
28}
返回示例
CreateAclResponse 字段结构:
1Username : string 用户名
相关 API
ListAcls 查询权限列表
功能说明
查询集群下的 ACL 列表,支持按 Username、PatternType、ResourceType、ResourceName 过滤。这些字段的 json tag 均为 -,由方法内部转成查询参数。
与 ListUsers 一致,参数类型为 interface{},接受集群 ID 字符串、ListAclRequest 值或指针。
请求示例
1package kafkademo
2
3import (
4 "fmt"
5 "log"
6
7 "github.com/baidubce/bce-sdk-go/services/kafka"
8)
9
10func listAcls() {
11 client := newKafkaClient()
12
13 request := &kafka.ListAclRequest{
14 ClusterID: "{{集群 ID}}",
15 Username: "{{用户名}}",
16 }
17
18 response, err := client.ListAcls(request)
19 if err != nil {
20 log.Fatalf("failed to list acls: %v", err)
21 }
22
23 for _, acl := range response.Acls {
24 fmt.Printf("username=%s patternType=%s resourceType=%s resourceName=%s operation=%s
25",
26 acl.Username, acl.PatternType, acl.ResourceType,
27 acl.ResourceName, acl.Operation)
28 }
29}
返回示例
ListAclResponse.Acls 元素(Acl)字段结构:
1Username : string 用户名
2PatternType : string 匹配模式
3ResourceType : string 资源类型
4ResourceName : string 资源名称
5Operation : string 操作类型
相关 API
任务管理
资源变更、主题重新分区等操作可能异步执行。SDK 调用成功仅表示任务创建成功,不代表操作已经完成。
请继续查询任务状态,直至任务进入成功或失败状态。多数变更接口的响应包含 ActionID 字段,即任务 ID。
1package kafkademo
2
3import (
4 "fmt"
5 "log"
6 "time"
7
8 "github.com/baidubce/bce-sdk-go/services/kafka"
9)
10
11func increaseBrokerAndWait() {
12 client := newKafkaClient()
13 clusterID := "{{集群 ID}}"
14
15 // 1. 提交异步变更,取回 actionId
16 brokerNodes := 6
17 increaseResponse, err := client.IncreaseBrokerCount(&kafka.IncreaseBrokerCountRequest{
18 ClusterID: clusterID,
19 NumberOfBrokerNodes: &brokerNodes,
20 })
21 if err != nil {
22 log.Fatalf("failed to increase broker count: %v", err)
23 }
24
25 actionID := increaseResponse.ActionID
26 fmt.Printf("actionId=%s
27", actionID)
28
29 // 2. 轮询任务状态,直至进入终态
30 for {
31 jobResponse, err := client.GetJob(&kafka.GetJobDetailRequest{
32 ClusterID: clusterID,
33 ActionID: actionID,
34 })
35 if err != nil {
36 log.Fatalf("failed to get job: %v", err)
37 }
38
39 job := jobResponse.Job
40 if job == nil {
41 log.Fatal("job not found")
42 }
43
44 fmt.Printf("status=%s
45", job.Status)
46 for _, operation := range job.Operations {
47 process := 0
48 if operation.Process != nil {
49 process = *operation.Process
50 }
51 fmt.Printf(" operationId=%s type=%s state=%s process=%d
52",
53 operation.OperationID, operation.Type, operation.State, process)
54 }
55
56 if isFinished(job.Status) {
57 break
58 }
59 time.Sleep(10 * time.Second)
60 }
61}
62
63func isFinished(status string) bool {
64 // 终态取值请以服务端返回为准
65 return status == "{{成功状态}}" || status == "{{失败状态}}"
66}
GetOperation 返回 *OperationDetail,包含 Process、Schedule、Groups、SourceContext、TargetContext 等更细粒度的进度信息。注意 OperationDetail.Process 是值类型 int,而 Operation.Process 是 *int。
StartJob、CancelJob、SuspendJob、ResumeJob 用于控制任务执行,四者都只返回 ActionID。
轮询间隔建议不小于 10 秒,避免高频调用触发限流。
相关 API:
评价此篇文章
