📌 本文是「Kafka 3.9 教程」系列第 六 篇 · 系列目录
本文档基于 Apache Kafka 3.9 官方文档整理
目录
6.1 安全架构概述
Kafka 提供三层安全机制:
6.1.1 安全层次
┌─────────────────────────────────────────────────────┐
│ 应用层 │
│ (ACL 授权控制) │
├─────────────────────────────────────────────────────┤
│ 认证层 │
│ (SASL/SSL 认证) │
├─────────────────────────────────────────────────────┤
│ 传输层 │
│ (SSL/TLS 加密) │
└─────────────────────────────────────────────────────┘
6.1.2 安全功能
| 功能 | 说明 | 实现方式 |
|---|---|---|
| 加密 | 保护数据传输安全 | SSL/TLS |
| 认证 | 验证客户端身份 | SASL, SSL |
| 授权 | 控制资源访问 | ACL |
| 审计 | 记录访问日志 | 监控日志 |
6.2 监听器配置
6.2.1 多协议监听器
Kafka 支持同时运行多种协议的监听器。
# 监听器配置
listeners=PLAINTEXT://0.0.0.0:9092,SSL://0.0.0.0:9093,SASL_SSL://0.0.0.0:9094
# 协议映射
listener.security.protocol.map=PLAINTEXT:PLAINTEXT,SSL:SSL,SASL_SSL:SASL_SSL
6.2.2 监听器配置示例
# 内部通信(内网)
listeners=INTERNAL://0.0.0.0:9092
# 外部通信(外网)
listeners=EXTERNAL://0.0.0.0:9093
# 协议映射
listener.security.protocol.map=INTERNAL:PLAINTEXT,EXTERNAL:SASL_SSL
# 对外公告地址
advertised.listeners=INTERNAL://internal.kafka.example.com:9092,EXTERNAL://kafka.example.com:9093
6.3 SSL/TLS 加密
6.3.1 证书管理
生成密钥库
# 生成密钥对和证书
keytool -keystore kafka.keystore.jks -alias kafka -validity 365 -genkey -keyalg RSA
# 生成 CA 证书
openssl req -new -x509 -keyout ca-key -out ca-cert -days 365
# 导入 CA 证书到密钥库
keytool -keystore kafka.keystore.jks -alias CARoot -import -file ca-cert
# 导出证书签名请求
keytool -keystore kafka.keystore.jks -alias kafka -certreq -file kafka.csr
# 用 CA 签署证书
openssl x509 -req -CA ca-cert -CAkey ca-key -in kafka.csr -out kafka-signed.crt -days 365 -CAcreateserial
# 导入 CA 证书到密钥库
keytool -keystore kafka.keystore.jks -alias CARoot -import -file ca-cert
# 导入签署后的证书
keytool -keystore kafka.keystore.jks -import -alias kafka -file kafka-signed.crt
生成信任库
# 从 CA 证书创建信任库
keytool -keystore kafka.truststore.jks -alias CARoot -import -file ca-cert
6.3.2 Broker SSL 配置
# ========== Broker SSL 配置 ==========
# 密钥库
ssl.keystore.location=/etc/kafka/ssl/kafka.keystore.jks
ssl.keystore.password=keystore-password
ssl.key.password=key-password
# 信任库
ssl.truststore.location=/etc/kafka/ssl/kafka.truststore.jks
ssl.truststore.password=truststore-password
# 协议和算法
ssl.protocol=TLSv1.3
ssl.enabled.protocols=TLSv1.3,TLSv1.2
# 客户端认证
ssl.client.auth=required # none, requested, required
# 主机名验证
ssl.hostname.verification=true
6.3.3 客户端 SSL 配置
import java.util.Properties;
import org.apache.kafka.clients.admin.AdminClientConfig;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.producer.ProducerConfig;
public class SSLConfigExample {
public static Properties getSSLConfig() {
Properties props = new Properties();
// Broker 地址
props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9093");
// 信任库
props.put("security.protocol", "SSL");
props.put("ssl.truststore.location", "/path/to/kafka.truststore.jks");
props.put("ssl.truststore.password", "truststore-password");
// 密钥库(客户端认证需要)
props.put("ssl.keystore.location", "/path/to/kafka.keystore.jks");
props.put("ssl.keystore.password", "keystore-password");
props.put("ssl.key.password", "key-password");
// 协议
props.put("ssl.protocol", "TLSv1.3");
props.put("ssl.enabled.protocols", "TLSv1.3,TLSv1.2");
return props;
}
}
6.4 SASL 认证
6.4.1 SASL 机制
| 机制 | 说明 | 适用场景 |
|---|---|---|
| PLAIN | 简单用户名/密码 | 测试/开发环境 |
| SCRAM-SHA-256 | 挑战响应认证 | 生产环境推荐 |
| SCRAM-SHA-512 | 更安全的挑战响应 | 高安全要求 |
| GSSAPI | Kerberos 认证 | 企业级单点登录 |
6.4.2 PLAIN 配置
JAAS 配置文件
# kafka_jaas.conf
KafkaServer {
org.apache.kafka.common.security.plain.PlainLoginModule required
username="admin"
password="admin-secret"
user_admin="admin-secret"
user_alice="alice-secret";
};
Client {
org.apache.kafka.common.security.plain.PlainLoginModule required
username="admin"
password="admin-secret";
};
Broker 配置
# server.properties
listeners=SASL_PLAINTEXT://0.0.0.0:9092
listener.security.protocol.map=SASL_PLAINTEXT:SASL_PLAINTEXT
# JAAS 配置
sasl.mechanism.inter.broker.protocol=PLAIN
sasl.enabled.mechanisms=PLAIN
# JAAS 文件路径
java.security.auth.login.config=/etc/kafka/kafka_jaas.conf
客户端配置
// 生产者/消费者 SASL 配置
props.put("security.protocol", "SASL_PLAINTEXT");
props.put("sasl.mechanism", "PLAIN");
props.put("sasl.jaas.config",
"org.apache.kafka.common.security.plain.PlainLoginModule required " +
"username=\"alice\" password=\"alice-secret\";");
6.4.3 SCRAM 配置
创建 SCRAM 凭证
# 创建用户
kafka-configs.sh --bootstrap-server localhost:9092 \
--alter --add-config 'SCRAM-SHA-256=[password=admin-secret]' \
--entity-type users --entity-name admin
kafka-configs.sh --bootstrap-server localhost:9092 \
--alter --add-config 'SCRAM-SHA-512=[password=alice-secret]' \
--entity-type users --entity-name alice
# 查看用户
kafka-configs.sh --describe --entity-type users --entity-name alice
# 删除用户配置
kafka-configs.sh --bootstrap-server localhost:9092 \
--alter --delete-config 'SCRAM-SHA-256,SCRAM-SHA-512' \
--entity-type users --entity-name alice
Broker SCRAM 配置
# server.properties
listeners=SASL_SSL://0.0.0.0:9092
listener.security.protocol.map=SASL_SSL:SASL_SSL
# SCRAM 配置
sasl.mechanism.inter.broker.protocol=SCRAM-SHA-256
sasl.enabled.mechanisms=SCRAM-SHA-256,SCRAM-SHA-512
# JAAS 配置
java.security.auth.login.config=/etc/kafka/kafka_jaas.conf
客户端 SCRAM 配置
props.put("security.protocol", "SASL_SSL");
props.put("sasl.mechanism", "SCRAM-SHA-256");
props.put("sasl.jaas.config",
"org.apache.kafka.common.security.scram.ScramLoginModule required " +
"username=\"alice\" password=\"alice-secret\";");
6.5 ACL 授权
6.5.1 ACL 概念
ACL(Access Control List)定义谁可以对什么资源做什么操作。
授权模型:
- Principal:被授权的主体(如 User:alice)
- Resource:被授权的资源(主题、组、集群)
- Operation:允许的操作
- Permission:允许或拒绝
6.5.2 常用操作
| 操作 | 资源类型 | 说明 |
|---|---|---|
| Read | Topic, Group | 消费消息 |
| Write | Topic | 生产消息 |
| Create | Topic, Group | 创建资源 |
| Delete | Topic | 删除资源 |
| Alter | Topic, Group | 修改配置 |
| Describe | Topic, Group | 查看信息 |
| ClusterAction | Cluster | 集群操作 |
| All | 全部 | 所有操作 |
6.5.3 添加 ACL
# 允许用户读取主题
kafka-acls.sh --authorizer-properties zookeeper.connect=localhost:2181 \
--add --allow-principal User:alice \
--operation Read \
--topic orders
# 允许用户写入主题
kafka-acls.sh --authorizer-properties zookeeper.connect=localhost:2181 \
--add --allow-principal User:alice \
--operation Write \
--topic orders
# 允许用户消费组
kafka-acls.sh --authorizer-properties zookeeper.connect=localhost:2181 \
--add --allow-principal User:alice \
--operation Read \
--group orders-group
# 批量添加(读写)
kafka-acls.sh --authorizer-properties zookeeper.connect=localhost:2181 \
--add --allow-principal User:bob \
--operation Read --operation Write \
--topic orders
# 拒绝某个用户
kafka-acls.sh --authorizer-properties zookeeper.connect=localhost:2181 \
--add --deny-principal User:charlie \
--operation Read \
--topic sensitive-data
6.5.4 查看 ACL
# 查看主题的 ACL
kafka-acls.sh --authorizer-properties zookeeper.connect=localhost:2181 \
--list --topic orders
# 查看用户的 ACL
kafka-acls.sh --authorizer-properties zookeeper.connect=localhost:2181 \
--list --principal User:alice
# 查看组的 ACL
kafka-acls.sh --authorizer-properties zookeeper.connect=localhost:2181 \
--list --group orders-group
6.5.5 删除 ACL
# 删除主题的所有 ACL
kafka-acls.sh --authorizer-properties zookeeper.connect=localhost:2181 \
--remove --topic orders
# 删除特定 ACL
kafka-acls.sh --authorizer-properties zookeeper.connect=localhost:2181 \
--remove --allow-principal User:alice \
--operation Read \
--topic orders
6.5.6 预定义 ACL
# 允许任何人读取(公开主题)
kafka-acls.sh --authorizer-properties zookeeper.connect=localhost:2181 \
--add --allow-principal User:* \
--operation Read \
--topic public-topic
# 允许任何人生产
kafka-acls.sh --authorizer-properties zookeeper.connect=localhost:2181 \
--add --allow-principal User:* \
--operation Write \
--topic public-topic
6.6 ZooKeeper 安全
6.6.1 ZooKeeper 认证
# server.properties
zookeeper.set.acl=true
启用 ZooKeeper ACL
# 创建带 ACL 的 znode
zkCli.sh
> create /kafka-acls "data" ip:127.0.0.1:r
6.6.2 ZooKeeper 加密
# 配置 ZooKeeper TLS
secureClientPort=2182
ssl.keyStore.location=/etc/kafka/ssl/kafka.keystore.jks
ssl.keyStore.password=password
ssl.trustStore.location=/etc/kafka/ssl/kafka.truststore.jks
ssl.trustStore.password=password
附录:安全配置清单
生产环境安全配置
# ========== Broker 安全配置 ==========
# 监听器
listeners=SSL://0.0.0.0:9092,SASL_SSL://0.0.0.0:9093
listener.security.protocol.map=SSL:SSL,SASL_SSL:SASL_SSL
# SSL 配置
ssl.keystore.location=/etc/kafka/ssl/kafka.keystore.jks
ssl.keystore.password=password
ssl.key.password=password
ssl.truststore.location=/etc/kafka/ssl/kafka.truststore.jks
ssl.truststore.password=password
ssl.client.auth=required
ssl.protocol=TLSv1.3
# SASL 配置
sasl.enabled.mechanisms=SCRAM-SHA-256
sasl.mechanism.inter.broker.protocol=SCRAM-SHA-256
# ACL 配置
authorizer.class.name=kafka.security.authorizer.AclAuthorizer
allow.everyone.if.no.acl.found=false
# ZooKeeper ACL
zookeeper.set.acl=true
下一步学习
- 第七章:Kafka Connect - 数据集成框架
- 第八章:Kafka Streams - 流处理库
- 实践:安全配置实战