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

如何创建自定义Kafka授权器并在本地使用?求Authorizer接口实现示例

Kafka自定义授权器实现与本地部署

一、实现org.apache.kafka.server.authorizer.Authorizer接口示例

首先确保项目中引入Kafka Server依赖(以Maven为例):

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-server</artifactId>
    <version>2.8.0</version>
    <scope>provided</scope>
</dependency>

以下是一个基础的自定义授权器实现,核心逻辑为:允许用户admin拥有所有资源的全部权限,其他用户仅允许访问以user-为前缀的Topic:

import org.apache.kafka.common.acl.AclOperation;
import org.apache.kafka.common.resource.ResourceType;
import org.apache.kafka.server.authorizer.Authorizer;
import org.apache.kafka.server.authorizer.AuthorizationRequestContext;
import org.apache.kafka.server.authorizer.AuthorizationResult;

import java.util.Map;
import java.util.concurrent.CompletableFuture;

public class CustomKafkaAuthorizer implements Authorizer {

    @Override
    public void configure(Map<String, ?> configs) {
        // 加载自定义配置,比如从configs中读取允许的用户列表、Topic前缀等
        String allowedTopicPrefix = (String) configs.get("custom.authorizer.topic.prefix");
        // 初始化逻辑
    }

    @Override
    public CompletableFuture<AuthorizationResult> authorize(AuthorizationRequestContext requestContext, Map<Integer, ? extends AclOperation> operations) {
        // 获取请求的用户主体
        String userPrincipal = requestContext.principal().getName();
        
        // 权限判断逻辑
        if ("admin".equals(userPrincipal)) {
            // admin用户允许所有操作
            return CompletableFuture.completedFuture(AuthorizationResult.ALLOWED);
        }

        // 遍历所有请求的操作和资源
        for (Map.Entry<Integer, ? extends AclOperation> entry : operations.entrySet()) {
            AclOperation operation = entry.getValue();
            ResourceType resourceType = requestContext.resourceType(entry.getKey());
            String resourceName = requestContext.resourceName(entry.getKey());

            // 仅允许非admin用户访问以"user-"为前缀的Topic,且操作限于读写
            if (resourceType == ResourceType.TOPIC && resourceName.startsWith("user-") 
                && (operation == AclOperation.READ || operation == AclOperation.WRITE)) {
                return CompletableFuture.completedFuture(AuthorizationResult.ALLOWED);
            }
        }

        // 不符合条件则拒绝
        return CompletableFuture.completedFuture(AuthorizationResult.DENIED);
    }

    @Override
    public void close() {
        // 资源释放逻辑,比如关闭数据库连接等
    }

    // 其他可选方法(如addAcls、removeAcls等)可根据需求实现
}

二、本地Kafka环境配置使用

1. 打包自定义授权器

将上述代码编译打包为jar包(例如custom-kafka-authorizer.jar),复制到本地Kafka安装目录的libs文件夹下。

2. 修改Kafka配置文件

打开Kafka安装目录下的config/server.properties,添加以下配置:

# 指定自定义授权器的全类名
authorizer.class.name=com.your.package.CustomKafkaAuthorizer
# 自定义授权器的配置参数(对应configure方法中的configs)
custom.authorizer.topic.prefix=user-
# 关闭默认的全访问权限,严格执行ACL检查
allow.everyone.if.no.acl.found=false

3. 配置用户认证(可选但建议)

为了让授权器能获取到用户主体,建议开启SASL认证。修改config/server.properties添加:

listeners=PLAINTEXT://localhost:9092,SASL_PLAINTEXT://localhost:9093
security.inter.broker.protocol=SASL_PLAINTEXT
sasl.enabled.mechanisms=PLAIN
sasl.mechanism.inter.broker.protocol=PLAIN

同时创建config/kafka_server_jaas.conf文件,内容如下:

KafkaServer {
    org.apache.kafka.common.security.plain.PlainLoginModule required
    username="admin"
    password="admin-secret"
    user_admin="admin-secret"
    user_user1="user1-secret";
};

设置环境变量指定JAAS配置文件(Windows用set,Linux/macOS用export):

export KAFKA_OPTS="-Djava.security.auth.login.config=/path/to/kafka/config/kafka_server_jaas.conf"

4. 重启Kafka服务

停止现有Kafka服务后,重新启动:

# 先停止ZooKeeper(如果使用内置ZK)
bin/zookeeper-server-stop.sh config/zookeeper.properties
# 停止Kafka
bin/kafka-server-stop.sh
# 启动ZooKeeper
bin/zookeeper-server-start.sh -daemon config/zookeeper.properties
# 启动Kafka
bin/kafka-server-start.sh -daemon config/server.properties

5. 测试授权效果

测试admin用户(拥有全权限)

创建producer-admin.properties:

bootstrap.servers=localhost:9093
security.protocol=SASL_PLAINTEXT
sasl.mechanism=PLAIN
sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required username="admin" password="admin-secret";

启动生产者发送消息到任意Topic:

bin/kafka-console-producer.sh --broker-list localhost:9093 --topic test-topic --producer.config producer-admin.properties

测试普通用户(仅允许访问user-前缀Topic)

创建producer-user1.properties:

bootstrap.servers=localhost:9093
security.protocol=SASL_PLAINTEXT
sasl.mechanism=PLAIN
sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required username="user1" password="user1-secret";

尝试发送消息到非user-前缀Topic,会收到权限拒绝错误:

bin/kafka-console-producer.sh --broker-list localhost:9093 --topic test-topic --producer.config producer-user1.properties

发送消息到user-testTopic则可以成功:

bin/kafka-console-producer.sh --broker-list localhost:9093 --topic user-test --producer.config producer-user1.properties

内容的提问来源于stack exchange,提问作者VIKAS TADGE

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 17:42:11