You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于Strimzi部署的Kafka:Java应用通过Admin Client运行时创建用户并关联ACL

可行性与实现方案

完全可行。Strimzi Operator部署的Kafka完全兼容原生Kafka的Admin API,同时也支持通过Kubernetes自定义资源(KafkaUser CR)管理用户与ACL,Java应用可以根据运行环境选择以下两种方案:

方案一:原生Kafka AdminClient直接操作(无需K8s权限)

适合Java应用能直接访问Kafka集群,且持有超级用户凭证(如Strimzi默认创建的admin用户)的场景。

实现步骤

  1. 配置AdminClient连接参数,包含Kafka bootstrap地址、超级用户的SASL认证信息(Strimzi默认启用SCRAM-SHA-512认证)。
  2. 调用AdminClient API创建SCRAM认证用户。
  3. 为用户添加指定资源的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配置。

实现步骤

  1. 引入Fabric8 Kubernetes Client与Strimzi API依赖。
  2. 构建KafkaUser CR对象,指定认证方式(SCRAM)与ACL规则。
  3. 调用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.17 09:25:03