初始化
更新时间:2026-09-09
确认 Endpoint
请根据 Kafka 资源所在地域,从服务域名中获取对应的 Endpoint。
Plain Text
1{{地域名称}}:{{Kafka Endpoint}}
Java SDK 的地址解析规则与其他语言不同,需要特别注意:
- 显式调用
setEndpoint(...)时,SDK 使用该地址。地址不含://时按配置的协议自动补全前缀。 - 未设置 Endpoint 时,SDK 按
<protocol>://kafka.<region>.baidubce.com拼接。默认协议为HTTP,默认区域为bj。 com.baidubce.Region当前只定义了CN_N1(对应bj),setRegion(...)没有其他可选值。
因此访问 bj 以外的地域必须显式设置 Endpoint,不能通过 setRegion(...) 切换。生产环境应同时显式设置 Protocol.HTTPS。
获取访问密钥
调用 Kafka OpenAPI 需要使用 Access Key ID(AK)和 Secret Access Key(SK)进行身份认证。
获取方式请参见获取 AK/SK。
请勿将 AK、SK 明文写入代码,建议通过环境变量或密钥管理服务读取。
新建 Kafka Client
Java
1package com.example.kafka;
2
3import com.baidubce.Protocol;
4import com.baidubce.auth.DefaultBceCredentials;
5import com.baidubce.services.kafka.KafkaClient;
6import com.baidubce.services.kafka.KafkaClientConfiguration;
7
8public class KafkaClientFactory {
9
10 public static KafkaClient create() {
11 KafkaClientConfiguration config = new KafkaClientConfiguration();
12 config.setCredentials(new DefaultBceCredentials(
13 System.getenv("BCE_ACCESS_KEY_ID"),
14 System.getenv("BCE_SECRET_ACCESS_KEY")));
15 config.setEndpoint(System.getenv("KAFKA_ENDPOINT"));
16 config.setProtocol(Protocol.HTTPS);
17 return new KafkaClient(config);
18 }
19}
KafkaClient 可在多线程间复用,不需要为每个请求重复创建。进程退出前调用 client.shutdown() 释放连接池。
无参构造 new KafkaClient() 使用默认配置且不含凭证,仅适用于匿名或后续自行注入凭证的场景。
使用临时安全凭证(可选)
Java
1import com.baidubce.auth.DefaultBceSessionCredentials;
2
3KafkaClientConfiguration config = new KafkaClientConfiguration();
4config.setCredentials(new DefaultBceSessionCredentials(
5 System.getenv("BCE_STS_ACCESS_KEY_ID"),
6 System.getenv("BCE_STS_SECRET_ACCESS_KEY"),
7 System.getenv("BCE_STS_SESSION_TOKEN")));
8config.setEndpoint(System.getenv("KAFKA_ENDPOINT"));
9config.setProtocol(Protocol.HTTPS);
10
11KafkaClient client = new KafkaClient(config);
STS 凭证过期后,需要使用新凭证重新创建客户端。
所有请求对象也支持 withRequestCredentials(...),为单次请求指定凭证并覆盖客户端级配置:
Java
1ListClustersRequest request = new ListClustersRequest();
2request.withRequestCredentials(new DefaultBceCredentials("<ak>", "<sk>"));
createUser 与 resetUserPassword 的密码加密所用 SK 遵循同一优先级:存在请求级凭证时使用请求级 SK,否则使用客户端配置中的 SK。
自定义客户端配置
KafkaClientConfiguration 继承 BceClientConfiguration,未显式配置时使用以下默认值:
| 配置项 | 默认值 |
|---|---|
region |
CN_N1(bj) |
endpoint |
<protocol>://kafka.<region>.baidubce.com |
protocol |
HTTP |
connectionTimeoutInMillis |
50000 毫秒 |
socketTimeoutInMillis |
50000 毫秒 |
maxConnections |
50 |
retryPolicy |
SDK 默认退避重试策略 |
常用配置项:
| 方法 | 说明 |
|---|---|
setEndpoint(String) |
Kafka 服务地址,优先级高于 region |
setRegion(Region) |
Endpoint 为空时用于拼接服务地址 |
setCredentials(BceCredentials) |
AK/SK 或 STS 凭证 |
setProtocol(Protocol) |
HTTP 或 HTTPS |
setRetryPolicy(RetryPolicy) |
请求重试策略 |
setProxyHost(String) / setProxyPort(int) |
HTTP 代理 |
setConnectionTimeoutInMillis(int) |
建立连接超时,单位毫秒 |
setSocketTimeoutInMillis(int) |
读取数据超时,单位毫秒 |
setMaxConnections(int) |
连接池最大连接数 |
setUserAgent(String) |
自定义 User-Agent |
请求签名与路径
- 所有接口路径统一带
/v2前缀。 - 参与签名的 header 固定为
host与x-bce-date。 POST/PUT请求体为 JSON,字符集为 UTF-8。
评价此篇文章
