基于Strimzi部署的Kafka:Java应用通过Admin Client运行时创建用户并关联ACL
可行性与实现方案
完全可行。Strimzi Operator部署的Kafka完全兼容原生Kafka的Admin API,同时也支持通过Kubernetes自定义资源(KafkaUser CR)管理用户与ACL,Java应用可以根据运行环境选择以下两种方案:
方案一:原生Kafka AdminClient直接操作(无需K8s权限)
适合Java应用能直接访问Kafka集群,且持有超级用户凭证(如Strimzi默认创建的admin用户)的场景。
实现步骤
- 配置AdminClient连接参数,包含Kafka bootstrap地址、超级用户的SASL认证信息(Strimzi默认启用SCRAM-SHA-512认证)。
- 调用AdminClient API创建SCRAM认证用户。
- 为用户添加指定资源的ACL规则。
代码示例
import org.apache.kafka.clients.admin.AdminClient; import org.apache.kafka.clients.admin.AdminClientConfig; import org.apache.kafka.clients.admin.CreateAclsResult; import org.apache.kafka.clients.admin.CreateUsersResult; import org.apache.kafka.common.acl.AclBinding; import org.apache.kafka.common.acl.AclOperation; import org.apache.kafka.common.acl.AclPermissionType; import org.apache.kafka.common.resource.ResourcePattern; import org.apache.kafka.common.resource.ResourceType; import org.apache.kafka.common.security.scram.ScramMechanism; import org.apache.kafka.common.security.scram.UserScramCredential; import java.util.Collections; import java.util.Properties; import java.util.concurrent.ExecutionException; public class KafkaUserAdmin { public static void main(String[] args) throws ExecutionException, InterruptedException { // 1. 初始化AdminClient配置 Properties adminProps = new Properties(); adminProps.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "my-kafka-cluster-kafka-bootstrap:9093"); adminProps.put(AdminClientConfig.SECURITY_PROTOCOL_CONFIG, "SASL_SSL"); adminProps.put("sasl.mechanism", "SCRAM-SHA-512"); // 超级用户凭证从Strimzi生成的Secret中获取(如my-kafka-cluster-admin的secret) adminProps.put("sasl.jaas.config", "org.apache.kafka.common.security.scram.ScramLoginModule required username=\"admin\" password=\"your-admin-password\";"); try (AdminClient adminClient = AdminClient.create(adminProps)) { String targetUser = "business-user-001"; String userPassword = "secure-pass-2024"; // 2. 创建SCRAM认证用户 CreateUsersResult createUserResult = adminClient.createUsers(Collections.singletonList( new UserScramCredential(targetUser, ScramMechanism.SCRAM_SHA_512, userPassword.toCharArray()) )); createUserResult.all().get(); // 等待操作完成 // 3. 配置ACL:允许用户读写指定业务Topic AclBinding readAcl = new AclBinding( new ResourcePattern(ResourceType.TOPIC, "business-order-topic", ResourcePattern.LITERAL), new org.apache.kafka.common.acl.AccessControlEntry("User:" + targetUser, "*", AclOperation.READ, AclPermissionType.ALLOW) ); AclBinding writeAcl = new AclBinding( new ResourcePattern(ResourceType.TOPIC, "business-order-topic", ResourcePattern.LITERAL), new org.apache.kafka.common.acl.AccessControlEntry("User:" + targetUser, "*", AclOperation.WRITE, AclPermissionType.ALLOW) ); CreateAclsResult createAclResult = adminClient.createAcls(Collections.singletonList(readAcl)); createAclResult.all().get(); createAclResult = adminClient.createAcls(Collections.singletonList(writeAcl)); createAclResult.all().get(); System.out.printf("用户 %s 已创建并配置ACL%n", targetUser); } } }
注意事项
- 超级用户的凭证需要从Strimzi生成的对应Secret中提取(K8s中执行
kubectl get secret my-kafka-cluster-admin -o jsonpath='{.data.password}' | base64 -d获取)。 - 确保Java应用的JDK版本与Kafka客户端版本兼容(建议使用Kafka 2.8+版本的客户端)。
方案二:通过Kubernetes API创建KafkaUser CR(K8s环境内应用)
适合Java应用运行在Kubernetes集群内,且拥有操作KafkaUser自定义资源的RBAC权限的场景,由Strimzi Operator自动完成用户创建与ACL配置。
实现步骤
- 引入Fabric8 Kubernetes Client与Strimzi API依赖。
- 构建KafkaUser CR对象,指定认证方式(SCRAM)与ACL规则。
- 调用K8s API提交KafkaUser资源,由Operator后续处理。
依赖配置(Maven)
<dependencies> <!-- Fabric8 Kubernetes Client --> <dependency> <groupId>io.fabric8</groupId> <artifactId>kubernetes-client</artifactId> <version>6.8.1</version> </dependency> <!-- Strimzi Kafka Custom Resource API --> <dependency> <groupId>io.strimzi</groupId> <artifactId>strimzi-kafka-api</artifactId> <version>0.34.0</version> <scope>provided</scope> </dependency> </dependencies>
代码示例
import io.fabric8.kubernetes.client.DefaultKubernetesClient; import io.fabric8.kubernetes.client.KubernetesClient; import io.strimzi.api.kafka.model.KafkaUser; import io.strimzi.api.kafka.model.KafkaUserBuilder; import io.strimzi.api.kafka.model.acl.AclRule; import io.strimzi.api.kafka.model.acl.AclRuleBuilder; import io.strimzi.api.kafka.model.authentication.ScramSha512AuthenticationBuilder; import java.util.Collections; public class StrimziUserManager { public static void main(String[] args) { String kafkaNamespace = "kafka"; String targetUser = "business-user-002"; // 初始化K8s客户端(自动读取集群内Pod的服务账号权限) try (KubernetesClient k8sClient = new DefaultKubernetesClient()) { // 构建ACL规则:允许用户读写指定Topic AclRule topicReadRule = new AclRuleBuilder() .withResourceType("topic") .withResourceName("business-payment-topic") .withOperation("read") .withPermissionType("allow") .build(); AclRule topicWriteRule = new AclRuleBuilder() .withResourceType("topic") .withResourceName("business-payment-topic") .withOperation("write") .withPermissionType("allow") .build(); // 构建KafkaUser CR KafkaUser kafkaUser = new KafkaUserBuilder() .withNewMetadata() .withName(targetUser) .withNamespace(kafkaNamespace) .endMetadata() .withNewSpec() .withAuthentication(new ScramSha512AuthenticationBuilder().build()) .withAcls(Collections.singletonList(topicReadRule)) .addToAcls(topicWriteRule) .endSpec() .build(); // 创建或更新KafkaUser资源 k8sClient.resources(KafkaUser.class) .inNamespace(kafkaNamespace) .createOrReplace(kafkaUser); System.out.printf("KafkaUser %s 已提交,Strimzi Operator将自动创建用户与ACL%n", targetUser); } } }
注意事项
- 需要为Java应用的服务账号配置RBAC权限,示例如下:
apiVersion: rbac.authorization.k8s.io/v1 kind: Role metadata: namespace: kafka name: kafka-user-manager rules: - apiGroups: ["kafka.strimzi.io"] resources: ["kafkausers"] verbs: ["create", "get", "update", "delete"] --- apiVersion: rbac.authorization.k8s.io/v1 kind: RoleBinding metadata: namespace: kafka name: kafka-user-manager-binding subjects: - kind: ServiceAccount name: your-app-service-account namespace: kafka roleRef: kind: Role name: kafka-user-manager apiGroup: rbac.authorization.k8s.io - Strimzi Operator会自动生成用户密码,并存储在与KafkaUser同名的Secret中。
通用注意事项
- 幂等性:建议在创建用户前先检查是否已存在,避免重复操作(AdminClient可调用
describeUsers,K8s API可调用get方法)。 - 权限最小化:执行操作的主体(超级用户或K8s服务账号)应仅拥有必要权限,避免过度授权。
- 环境一致性:确保Java应用使用的认证方式与Kafka集群的配置完全匹配(如SASL机制、SSL配置)。
内容的提问来源于stack exchange,提问作者raj kumar
相关产品推荐
相关产品推荐

