安装 SDK
安装 SDK
运行环境
- Go 版本:1.11 及以上
- 已创建 Go Module 项目
- 已开通消息服务 for Kafka
- 调用账号具有相关资源的操作权限
安装 SDK
1go get github.com/baidubce/bce-sdk-go@{{SDK版本}}
执行后 go.mod 中出现依赖声明:
1module example.com/kafkademo
2
3go 1.16
4
5require github.com/baidubce/bce-sdk-go {{SDK版本}}
导入 Kafka 子包:
1import "github.com/baidubce/bce-sdk-go/services/kafka"
kafka 包只依赖标准库和同仓库的 auth、bce、http、util 子包,接入方不需要引入额外的第三方依赖。
初始化
确认 Endpoint
请根据 Kafka 资源所在地域,从服务域名中获取对应的 Endpoint。
1{{地域名称}}:{{Kafka Endpoint}}
Go SDK 的地址解析规则:
- 显式传入
Endpoint时,SDK 使用该地址。 Endpoint为空但设置了Region时,SDK 按kafka.<region>.baidubce.com拼接。Endpoint与Region均为空时,Region取默认值bj,即kafka.bj.baidubce.com。
Region 字段是普通 string,可直接填写任意地域标识,不需要枚举常量,也不必为非默认地域改用 Endpoint。
包级常量 kafka.DEFAULT_ENDPOINT 的值为 kafka.bj.baidubce.com。
Endpoint 不含协议前缀时,由底层 HTTP 客户端按默认协议访问;需要 HTTPS 时请传入带 https:// 前缀的完整地址。
获取访问密钥
调用 Kafka OpenAPI 需要使用 Access Key ID(AK)和 Secret Access Key(SK)进行身份认证。
获取方式请参见获取 AK/SK。
请勿将 AK、SK 明文写入代码,建议通过环境变量或密钥管理服务读取。
CreateUser 与 ResetUserPassword 需要用 SK 加密密码,SK 长度必须不小于 16 字节,否则返回 secret access key should be at least 16 bytes。
新建 Kafka Client
1package kafkademo
2
3import (
4 "log"
5 "os"
6
7 "github.com/baidubce/bce-sdk-go/services/kafka"
8)
9
10// newKafkaClient 创建可复用的 Kafka 客户端。
11func newKafkaClient() *kafka.Client {
12 client, err := kafka.NewClient(
13 os.Getenv("BCE_ACCESS_KEY_ID"),
14 os.Getenv("BCE_SECRET_ACCESS_KEY"),
15 os.Getenv("KAFKA_ENDPOINT"),
16 )
17 if err != nil {
18 log.Fatalf("failed to create kafka client: %v", err)
19 }
20 return client
21}
kafka.Client 可在多个 goroutine 间复用,不需要为每个请求重复创建。客户端没有显式关闭方法,连接由底层 http.Client 的连接池管理,进程退出即释放。
NewClient 内部调用 auth.NewBceCredentials 构造凭证,AK 或 SK 为空时返回错误,不存在不含凭证的客户端。需要匿名访问或延迟注入凭证时,使用 NewClientWithConfig 并自行设置 Credentials。
使用临时安全凭证(可选)
1package kafkademo
2
3import (
4 "log"
5 "os"
6
7 "github.com/baidubce/bce-sdk-go/services/kafka"
8)
9
10func newSTSClient() *kafka.Client {
11 client, err := kafka.NewClientWithSTS(
12 os.Getenv("BCE_STS_ACCESS_KEY_ID"),
13 os.Getenv("BCE_STS_SECRET_ACCESS_KEY"),
14 os.Getenv("BCE_STS_SESSION_TOKEN"),
15 os.Getenv("KAFKA_ENDPOINT"),
16 )
17 if err != nil {
18 log.Fatalf("failed to create kafka client with sts: %v", err)
19 }
20 return client
21}
NewClientWithSTS 内部使用 auth.NewSessionBceCredentials 构造带 Session Token 的凭证。STS 凭证过期后,需要使用新凭证重新创建客户端。
凭证只能在客户端层面指定。 请求结构体不支持为单次请求覆盖凭证,因此 CreateUser 与 ResetUserPassword 的密码加密固定使用客户端配置中的 SK(c.Config.Credentials.SecretAccessKey)。需要用不同账号的 SK 加密时,必须另建一个客户端。
自定义客户端配置
KafkaClientConfiguration 是 bce.BceClientConfiguration 的类型别名,未显式配置时使用以下默认值:
| 配置项 | 默认值 |
|---|---|
Region |
bj |
Endpoint |
kafka.<region>.baidubce.com |
UserAgent |
bce-sdk-go/<SDK版本>/<Go版本>/<GOOS>/<GOARCH> |
ConnectionTimeoutInMillis |
50000 毫秒 |
DialTimeout |
与 ConnectionTimeoutInMillis 一致 |
ReadTimeout |
50 秒 |
Retry |
退避重试,最多 3 次,最大延迟 20 秒,基数 300 毫秒 |
常用配置项:
| 字段 | 说明 |
|---|---|
Endpoint string |
Kafka 服务地址,优先级高于 Region |
Region string |
Endpoint 为空时用于拼接服务地址 |
Credentials *auth.BceCredentials |
AK/SK 或 STS 凭证 |
Retry bce.RetryPolicy |
请求重试策略 |
ProxyUrl string |
HTTP 代理地址 |
ConnectionTimeoutInMillis int |
建立连接超时,单位毫秒 |
DialTimeout *time.Duration |
建立连接的超时 |
ReadTimeout *time.Duration |
读取连接的超时 |
HTTPClientTimeout *time.Duration |
整个 HTTP 请求的超时 |
HTTPClient *http.Client |
自定义 HTTP 客户端 |
UserAgent string |
自定义 User-Agent |
NoVerifySSL bool |
跳过 TLS 证书校验,生产环境不建议开启 |
UploadRatelimit *int64 |
上传限速,单位 KB/s |
DownloadRatelimit *int64 |
下载限速,单位 KB/s |
超时类字段是 *time.Duration,需要先声明变量再取地址:
1package kafkademo
2
3import (
4 "log"
5 "os"
6 "time"
7
8 "github.com/baidubce/bce-sdk-go/auth"
9 "github.com/baidubce/bce-sdk-go/bce"
10 "github.com/baidubce/bce-sdk-go/services/kafka"
11)
12
13func newClientWithConfig() *kafka.Client {
14 credentials, err := auth.NewBceCredentials(
15 os.Getenv("BCE_ACCESS_KEY_ID"),
16 os.Getenv("BCE_SECRET_ACCESS_KEY"),
17 )
18 if err != nil {
19 log.Fatalf("failed to build credentials: %v", err)
20 }
21
22 dialTimeout := 10 * time.Second
23 readTimeout := 30 * time.Second
24 httpClientTimeout := 60 * time.Second
25
26 client, err := kafka.NewClientWithConfig(&kafka.KafkaClientConfiguration{
27 Region: "{{地域标识}}",
28 Credentials: credentials,
29 ConnectionTimeoutInMillis: 10 * 1000,
30 DialTimeout: &dialTimeout,
31 ReadTimeout: &readTimeout,
32 HTTPClientTimeout: &httpClientTimeout,
33 Retry: bce.NewBackOffRetryPolicy(3, 20000, 300),
34 })
35 if err != nil {
36 log.Fatalf("failed to create kafka client: %v", err)
37 }
38 return client
39}
NewClientWithConfig(nil) 返回错误 kafka client configuration should not be nil。
SignOption 由 SDK 内部固定设置,参与签名的 header 为 host 与 x-bce-date,调用方传入的值会被覆盖。
请求签名与路径
- 所有接口路径统一带
/v2前缀。 - 参与签名的 header 固定为
host与x-bce-date。 POST/PUT请求体为 JSON,Content-Type 取bce.DEFAULT_CONTENT_TYPE。
评价此篇文章
