如何创建自定义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
相关产品推荐
相关产品推荐

