接口调用示例
集群管理
createCluster 创建集群
功能说明
创建 Kafka 集群。CreateClusterRequest 只强制校验 name,mode、type 与 provisioned 内部字段的合法组合由服务端校验。创建预置型集群时应完整填写 Provisioned,其中基本类型字段需显式赋值。
创建集群会产生计费,预付费场景下 Billing.isAutoPay 默认为 true,会直接扣费。
请求示例
1package com.example.kafka;
2
3import java.util.Arrays;
4import java.util.LinkedHashMap;
5import java.util.Map;
6
7import com.baidubce.services.kafka.KafkaClient;
8import com.baidubce.services.kafka.model.cluster.Authentication;
9import com.baidubce.services.kafka.model.cluster.Billing;
10import com.baidubce.services.kafka.model.cluster.ConfigMeta;
11import com.baidubce.services.kafka.model.cluster.CreateClusterRequest;
12import com.baidubce.services.kafka.model.cluster.CreateClusterResponse;
13import com.baidubce.services.kafka.model.cluster.Mode;
14import com.baidubce.services.kafka.model.cluster.Provisioned;
15import com.baidubce.services.kafka.model.cluster.StorageMeta;
16import com.baidubce.services.kafka.model.cluster.StorageType;
17import com.baidubce.services.kafka.model.cluster.Tag;
18import com.baidubce.services.kafka.model.cluster.Type;
19
20public class CreateClusterDemo {
21
22 public static void main(String[] args) {
23 KafkaClient client = KafkaClientFactory.create();
24 try {
25 CreateClusterRequest request = new CreateClusterRequest();
26 request.setName("kafka-demo");
27 request.setMode(Mode.HP);
28 request.setType(Type.PROVISIONED);
29 request.setTags(Arrays.asList(
30 Tag.builder().tagKey("KAFKA-Cluster").tagValue("prod").build()));
31
32 Provisioned provisioned = new Provisioned();
33 provisioned.setKafkaVersion("2.7.2");
34 provisioned.setNodeType("kafka.g4.c2m8");
35 provisioned.setNumberOfBrokerNodes(3);
36 provisioned.setStorageMeta(StorageMeta.builder()
37 .storageType(StorageType.ENHANCED_SSD_PL1)
38 .storageSize(100)
39 .numberOfDisk(1)
40 .build());
41 provisioned.setVpcId("{{VPC ID}}");
42 provisioned.setSubnetIds(Arrays.asList("{{子网 ID}}"));
43 provisioned.setSecurityGroupIds(Arrays.asList("{{安全组 ID}}"));
44 provisioned.setPublicIpEnabled(false);
45 provisioned.setIntranetIpEnabled(false);
46 provisioned.setAclEnabled(true);
47 provisioned.setDeploySetEnabled(true);
48 provisioned.setAuthentications(Arrays.asList(
49 Authentication.builder().mode("SASL_SCRAM").build()));
50 provisioned.setBilling(Billing.builder()
51 .payment("Prepaid")
52 .timeLength(1)
53 .timeUnit("month")
54 .autoRenewEnabled(false)
55 .isAutoPay(true)
56 .build());
57
58 Map<String, String> configContext = new LinkedHashMap<String, String>();
59 configContext.put("auto.create.topics.enable", "true");
60 configContext.put("log.retention.hours", "168");
61 provisioned.setConfigMeta(ConfigMeta.builder().context(configContext).build());
62
63 request.setProvisioned(provisioned);
64
65 CreateClusterResponse response = client.createCluster(request);
66 System.out.println("clusterId=" + response.getClusterId());
67 } finally {
68 client.shutdown();
69 }
70 }
71}
返回示例
CreateClusterResponse 字段结构:
1clusterId : String 集群 ID
集群创建是异步过程,返回 clusterId 只表示创建请求已受理,需继续查询集群状态确认就绪。
相关 API
getClusterDetail 查询集群详情
功能说明
按 clusterId 查询集群详情。响应中的 ClusterDetail.provisioned 包含规格、网络、存储、认证和维护窗口等完整配置。
查询已删除集群请使用 getClusterDeletion,其请求类型为 GetClusterDeletionRequest,但响应复用 GetClusterDetailResponse。
请求示例
1package com.example.kafka;
2
3import com.baidubce.services.kafka.KafkaClient;
4import com.baidubce.services.kafka.model.cluster.ClusterDetail;
5import com.baidubce.services.kafka.model.cluster.GetClusterDetailRequest;
6import com.baidubce.services.kafka.model.cluster.GetClusterDetailResponse;
7
8public class GetClusterDetailDemo {
9
10 public static void main(String[] args) {
11 KafkaClient client = KafkaClientFactory.create();
12 try {
13 GetClusterDetailRequest request = new GetClusterDetailRequest();
14 request.setClusterId("{{集群 ID}}");
15
16 GetClusterDetailResponse response = client.getClusterDetail(request);
17 ClusterDetail detail = response.getCluster();
18 if (detail != null) {
19 System.out.printf("name=%s state=%s region=%s%n",
20 detail.getName(), detail.getState(), detail.getRegion());
21 }
22 } finally {
23 client.shutdown();
24 }
25 }
26}
返回示例
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 : List<Tag> 标签列表
11createTime : String 创建时间
12deleteTime : String 删除时间
相关 API
stopCluster 停止集群
功能说明
停止运行中的集群。停止后集群不再提供服务,客户端连接会中断,属于影响业务的操作,执行前需确认目标集群。
接口返回 actionId,需通过任务接口确认最终状态。
请求示例
1package com.example.kafka;
2
3import com.baidubce.services.kafka.KafkaClient;
4import com.baidubce.services.kafka.model.cluster.StopClusterRequest;
5import com.baidubce.services.kafka.model.cluster.StopClusterResponse;
6
7public class StopClusterDemo {
8
9 public static void main(String[] args) {
10 KafkaClient client = KafkaClientFactory.create();
11 try {
12 StopClusterRequest request = new StopClusterRequest();
13 request.setClusterId("{{集群 ID}}");
14
15 StopClusterResponse response = client.stopCluster(request);
16 System.out.printf("clusterId=%s actionId=%s%n",
17 response.getClusterId(), response.getActionId());
18 } finally {
19 client.shutdown();
20 }
21 }
22}
返回示例
StopClusterResponse 字段结构:
1clusterId : String 集群 ID
2actionId : String 任务 ID,用于后续查询任务状态
相关 API
主题管理
createTopic 创建主题
功能说明
在指定集群下创建主题。partitionNum 与 replicationFactor 是 int 基本类型,不赋值会发送 0,必须显式设置。
可用的主题级参数及取值范围通过 listTopicConfigOptions() 查询,该方法不接受请求对象。
请求示例
1package com.example.kafka;
2
3import java.util.LinkedHashMap;
4import java.util.Map;
5
6import com.baidubce.services.kafka.KafkaClient;
7import com.baidubce.services.kafka.model.topic.CreateTopicRequest;
8import com.baidubce.services.kafka.model.topic.CreateTopicResponse;
9
10public class CreateTopicDemo {
11
12 public static void main(String[] args) {
13 KafkaClient client = KafkaClientFactory.create();
14 try {
15 CreateTopicRequest request = new CreateTopicRequest();
16 request.setClusterId("{{集群 ID}}");
17 request.setTopicName("{{主题名称}}");
18 request.setPartitionNum(3);
19 request.setReplicationFactor(3);
20
21 Map<String, String> otherConfigs = new LinkedHashMap<String, String>();
22 otherConfigs.put("cleanup.policy", "delete");
23 otherConfigs.put("retention.ms", "604800000");
24 request.setOtherConfigs(otherConfigs);
25
26 CreateTopicResponse response = client.createTopic(request);
27 System.out.println("topicName=" + response.getTopicName());
28 } finally {
29 client.shutdown();
30 }
31 }
32}
返回示例
CreateTopicResponse 字段结构:
1topicName : String 主题名称
相关 API
updateTopic 变更主题
功能说明
变更主题分区数或主题级配置。partitionNum 与 otherConfigs 至少设置一项,否则抛出 IllegalArgumentException。分区数只能增加,不能减少。
UpdateTopicRequest.partitionNum 的字段类型是 String,同时提供 setPartitionNum(int) 重载完成转换,推荐使用 int 重载避免手工拼字符串。
请求示例
1package com.example.kafka;
2
3import com.baidubce.services.kafka.KafkaClient;
4import com.baidubce.services.kafka.model.topic.UpdateTopicRequest;
5import com.baidubce.services.kafka.model.topic.UpdateTopicResponse;
6
7public class UpdateTopicDemo {
8
9 public static void main(String[] args) {
10 KafkaClient client = KafkaClientFactory.create();
11 try {
12 UpdateTopicRequest request = new UpdateTopicRequest();
13 request.setClusterId("{{集群 ID}}");
14 request.setTopicName("{{主题名称}}");
15 request.setPartitionNum(6);
16
17 UpdateTopicResponse response = client.updateTopic(request);
18 System.out.println("topicName=" + response.getTopicName());
19 } finally {
20 client.shutdown();
21 }
22 }
23}
返回示例
UpdateTopicResponse 字段结构:
1topicName : String 主题名称
相关 API
listTopic 查询主题列表
功能说明
查询集群下的主题列表。topicName 可选,按前缀匹配过滤。该接口不分页,返回集群下符合条件的全部主题。
请求示例
1package com.example.kafka;
2
3import com.baidubce.services.kafka.KafkaClient;
4import com.baidubce.services.kafka.model.topic.ListTopicRequest;
5import com.baidubce.services.kafka.model.topic.ListTopicResponse;
6import com.baidubce.services.kafka.model.topic.Topic;
7
8public class ListTopicDemo {
9
10 public static void main(String[] args) {
11 KafkaClient client = KafkaClientFactory.create();
12 try {
13 ListTopicRequest request = new ListTopicRequest();
14 request.setClusterId("{{集群 ID}}");
15 request.setTopicName("order-");
16
17 ListTopicResponse response = client.listTopic(request);
18 for (Topic topic : response.getTopics()) {
19 System.out.printf("topicName=%s partitionNum=%d replicaNum=%d%n",
20 topic.getTopicName(), topic.getPartitionNum(), topic.getReplicaNum());
21 }
22 } finally {
23 client.shutdown();
24 }
25 }
26}
返回示例
ListTopicResponse.topics 元素(Topic)字段结构:
1topicName : String 主题名称
2createTime : String 创建时间
3readOnly : Boolean 是否只读
4partitionNum : Integer 分区数
5replicaNum : Integer 副本数
相关 API
消费组管理
listConsumerGroup 查询消费组列表
功能说明
查询集群下的消费组列表。groupName 可选,用于过滤。
请求示例
1package com.example.kafka;
2
3import com.baidubce.services.kafka.KafkaClient;
4import com.baidubce.services.kafka.model.consumer.Group;
5import com.baidubce.services.kafka.model.consumer.ListConsumerGroupRequest;
6import com.baidubce.services.kafka.model.consumer.ListConsumerGroupResponse;
7
8public class ListConsumerGroupDemo {
9
10 public static void main(String[] args) {
11 KafkaClient client = KafkaClientFactory.create();
12 try {
13 ListConsumerGroupRequest request = new ListConsumerGroupRequest();
14 request.setClusterId("{{集群 ID}}");
15
16 ListConsumerGroupResponse response = client.listConsumerGroup(request);
17 for (Group group : response.getGroups()) {
18 System.out.printf("groupName=%s state=%s coordinator=%s%n",
19 group.getGroupName(), group.getState(), group.getGroupCoordinatorId());
20 }
21 } finally {
22 client.shutdown();
23 }
24 }
25}
返回示例
ListConsumerGroupResponse.groups 元素(Group)字段结构:
1groupName : String 消费组名称
2updateTime : String 更新时间
3state : String 消费组状态
4groupCoordinatorId : Integer Coordinator 所在 Broker ID
相关 API
resetConsumerGroup 重置消费组位点
功能说明
重置消费组在指定主题分区上的消费位点。topicName、partitions、resetStrategy 都是必填,partitions 不能为空列表。
重置位点会直接改变消费进度,可能导致消息重复消费或跳过,执行前必须确认目标消费组、主题和分区。建议在消费组无活跃成员时执行。
请求示例
1package com.example.kafka;
2
3import java.util.Arrays;
4
5import com.baidubce.services.kafka.KafkaClient;
6import com.baidubce.services.kafka.model.consumer.ResetConsumerGroupRequest;
7import com.baidubce.services.kafka.model.consumer.ResetConsumerGroupResponse;
8
9public class ResetConsumerGroupDemo {
10
11 public static void main(String[] args) {
12 KafkaClient client = KafkaClientFactory.create();
13 try {
14 ResetConsumerGroupRequest request = new ResetConsumerGroupRequest();
15 request.setClusterId("{{集群 ID}}");
16 request.setGroupName("{{消费组名称}}");
17 request.setTopicName("{{主题名称}}");
18 request.setPartitions(Arrays.asList(0, 1, 2));
19 request.setResetStrategy("{{重置策略}}");
20 request.setResetValue("{{重置目标值}}");
21
22 ResetConsumerGroupResponse response = client.resetConsumerGroup(request);
23 System.out.println("groupName=" + response.getGroupName());
24 } finally {
25 client.shutdown();
26 }
27 }
28}
返回示例
ResetConsumerGroupResponse 字段结构:
1groupName : String 消费组名称
相关 API
用户管理
createUser 创建用户
功能说明
在集群下创建 SASL 用户。SDK 内部使用凭证 SK 的前 16 个字符作为密钥,以 AES/ECB/PKCS5Padding 加密密码后以大写十六进制字符串发送,调用方传入原始密码即可。
saslMechanisms 可选,取值示例为 SCRAM-SHA-256、SCRAM-SHA-512、PLAIN。
加密会直接写回请求对象的 password 字段。 同一个 CreateUserRequest 实例不能重复提交,否则第二次调用会对已加密的值再次加密。每次调用请新建请求对象。同样的约束适用于 ResetUserPasswordRequest。
不要记录原始密码、加密后的密码或完整请求体。
请求示例
1package com.example.kafka;
2
3import java.util.Arrays;
4import java.util.HashSet;
5
6import com.baidubce.services.kafka.KafkaClient;
7import com.baidubce.services.kafka.model.user.CreateUserRequest;
8import com.baidubce.services.kafka.model.user.CreateUserResponse;
9
10public class CreateUserDemo {
11
12 public static void main(String[] args) {
13 KafkaClient client = KafkaClientFactory.create();
14 try {
15 CreateUserRequest request = new CreateUserRequest();
16 request.setClusterId("{{集群 ID}}");
17 request.setUsername("{{用户名}}");
18 request.setPassword(System.getenv("KAFKA_USER_PASSWORD"));
19 request.setSaslMechanisms(new HashSet<String>(
20 Arrays.asList("SCRAM-SHA-256", "SCRAM-SHA-512")));
21
22 CreateUserResponse response = client.createUser(request);
23 System.out.println("username=" + response.getUsername());
24 } finally {
25 client.shutdown();
26 }
27 }
28}
返回示例
CreateUserResponse 字段结构:
1username : String 用户名
相关 API
listUsers 查询用户列表
功能说明
查询集群下的用户列表。listUsers(String clusterId) 是已标注 @Deprecated 的便捷重载,内部委托到 listUsers(ListUsersRequest),新代码应直接使用请求对象。
请求示例
1package com.example.kafka;
2
3import com.baidubce.services.kafka.KafkaClient;
4import com.baidubce.services.kafka.model.user.ListUserResponse;
5import com.baidubce.services.kafka.model.user.ListUsersRequest;
6import com.baidubce.services.kafka.model.user.User;
7
8public class ListUsersDemo {
9
10 public static void main(String[] args) {
11 KafkaClient client = KafkaClientFactory.create();
12 try {
13 ListUsersRequest request = new ListUsersRequest();
14 request.setClusterId("{{集群 ID}}");
15
16 ListUserResponse response = client.listUsers(request);
17 for (User user : response.getUsers()) {
18 System.out.printf("username=%s createTime=%s mechanisms=%s%n",
19 user.getUsername(), user.getCreateTime(), user.getSaslMechanisms());
20 }
21 } finally {
22 client.shutdown();
23 }
24 }
25}
返回示例
ListUserResponse.users 元素(User)字段结构:
1username : String 用户名
2createTime : String 创建时间
3saslMechanisms : Set<String> 已启用的 SASL 机制
相关 API
权限管理
createAcl 创建权限
功能说明
为指定用户创建 ACL 授权。username、patternType、resourceType、resourceName 必填,operations 不能为空列表。
注意 createAcl 使用 operations(List<String>,可一次授予多个操作),而 deleteAcl 使用 operation(单个 String)。
请求示例
1package com.example.kafka;
2
3import java.util.Arrays;
4
5import com.baidubce.services.kafka.KafkaClient;
6import com.baidubce.services.kafka.model.acl.CreateAclRequest;
7import com.baidubce.services.kafka.model.acl.CreateAclResponse;
8
9public class CreateAclDemo {
10
11 public static void main(String[] args) {
12 KafkaClient client = KafkaClientFactory.create();
13 try {
14 CreateAclRequest request = new CreateAclRequest();
15 request.setClusterId("{{集群 ID}}");
16 request.setUsername("{{用户名}}");
17 request.setPatternType("{{匹配模式}}");
18 request.setResourceType("{{资源类型}}");
19 request.setResourceName("{{资源名称}}");
20 request.setOperations(Arrays.asList("{{操作类型}}"));
21
22 CreateAclResponse response = client.createAcl(request);
23 System.out.println("username=" + response.getUsername());
24 } finally {
25 client.shutdown();
26 }
27 }
28}
返回示例
CreateAclResponse 字段结构:
1username : String 用户名
相关 API
listAcls 查询权限列表
功能说明
查询集群下的 ACL 列表,支持按 username、patternType、resourceType、resourceName 过滤。
listAcls(String clusterId) 与 listAcls(ListAclRequest) 两个重载都可用,后者支持过滤条件。
请求示例
1package com.example.kafka;
2
3import com.baidubce.services.kafka.KafkaClient;
4import com.baidubce.services.kafka.model.acl.Acl;
5import com.baidubce.services.kafka.model.acl.ListAclRequest;
6import com.baidubce.services.kafka.model.acl.ListAclResponse;
7
8public class ListAclsDemo {
9
10 public static void main(String[] args) {
11 KafkaClient client = KafkaClientFactory.create();
12 try {
13 ListAclRequest request = ListAclRequest.builder()
14 .clusterId("{{集群 ID}}")
15 .username("{{用户名}}")
16 .build();
17
18 ListAclResponse response = client.listAcls(request);
19 for (Acl acl : response.getAcls()) {
20 System.out.printf("username=%s resourceType=%s resourceName=%s operation=%s%n",
21 acl.getUsername(), acl.getResourceType(),
22 acl.getResourceName(), acl.getOperation());
23 }
24 } finally {
25 client.shutdown();
26 }
27 }
28}
返回示例
ListAclResponse.acls 元素(Acl)字段结构:
1username : String 用户名
2patternType : String 匹配模式
3resourceType : String 资源类型
4resourceName : String 资源名称
5operation : String 操作类型
相关 API
任务管理
资源变更、主题重新分区等操作可能异步执行。SDK 调用成功仅表示任务创建成功,不代表操作已经完成。
请继续查询任务状态,直至任务进入成功或失败状态。多数变更接口的响应包含 actionId 字段,即任务 ID。
1package com.example.kafka;
2
3import com.baidubce.services.kafka.KafkaClient;
4import com.baidubce.services.kafka.model.cluster.IncreaseBrokerCountRequest;
5import com.baidubce.services.kafka.model.cluster.IncreaseBrokerCountResponse;
6import com.baidubce.services.kafka.model.job.GetJobDetailRequest;
7import com.baidubce.services.kafka.model.job.GetJobDetailResponse;
8import com.baidubce.services.kafka.model.job.Job;
9import com.baidubce.services.kafka.model.job.Operation;
10
11public class AsyncJobDemo {
12
13 public static void main(String[] args) throws InterruptedException {
14 KafkaClient client = KafkaClientFactory.create();
15 try {
16 String clusterId = "{{集群 ID}}";
17
18 // 1. 提交异步变更,取回 actionId
19 IncreaseBrokerCountRequest increaseRequest = new IncreaseBrokerCountRequest();
20 increaseRequest.setClusterId(clusterId);
21 increaseRequest.setNumberOfBrokerNodes(6);
22
23 IncreaseBrokerCountResponse increaseResponse =
24 client.increaseBrokerCount(increaseRequest);
25 String actionId = increaseResponse.getActionId();
26 System.out.println("actionId=" + actionId);
27
28 // 2. 轮询任务状态,直至进入终态
29 while (true) {
30 GetJobDetailRequest jobRequest = new GetJobDetailRequest();
31 jobRequest.setClusterId(clusterId);
32 jobRequest.setActionId(actionId);
33
34 GetJobDetailResponse jobResponse = client.getJob(jobRequest);
35 Job job = jobResponse.getJob();
36 System.out.println("status=" + job.getStatus());
37
38 if (job.getOperations() != null) {
39 for (Operation operation : job.getOperations()) {
40 System.out.printf(" operationId=%s type=%s state=%s process=%s%n",
41 operation.getOperationId(), operation.getType(),
42 operation.getState(), operation.getProcess());
43 }
44 }
45
46 if (isFinished(job.getStatus())) {
47 break;
48 }
49 Thread.sleep(10 * 1000L);
50 }
51 } finally {
52 client.shutdown();
53 }
54 }
55
56 private static boolean isFinished(String status) {
57 // 终态取值请以服务端返回为准
58 return "{{成功状态}}".equals(status) || "{{失败状态}}".equals(status);
59 }
60}
getOperation 返回 OperationDetail,包含 process、schedule、groups、sourceContext、targetContext 等更细粒度的进度信息。startJob、cancelJob、suspendJob、resumeJob 用于控制任务执行。
轮询间隔建议不小于 10 秒,避免高频调用触发限流。
相关 API:
评价此篇文章
