主题
第 15 章 安全与多租户:把 Kafka 关进笼子里
目标读者:之前所有章节都默认
localhost:9092裸跑,没碰过 SSL / SASL / ACL,对「多团队共享一个 Kafka 集群」「生产环境不能让所有应用都拥有 root 权限」这件事完全没有概念的同学。学完你会:能说清「鉴权 / 授权 / 限流 / 审计」四件事在 Kafka 里分别由谁负责;能用 SCRAM + ACL + Quota 把一个共享集群切成「每个团队只能读写自己 Topic」的多租户环境;能解释「SuperUser 为什么要慎用」「为什么 SASL_PLAINTEXT 在公网上等于裸奔」。
0. 一个生活类比:写字楼的安保系统
把 Kafka 集群想象成一栋写字楼,里面住着各个业务团队(订单、风控、推荐、数仓……)。这栋写字楼至少要有 4 道安保措施:
| Kafka 安全机制 | 写字楼对应物 | 干什么 |
|---|---|---|
| SSL(TLS) | 门口的金属探测门 + 防偷拍玻璃 | 加密通信、防中间人;只让门牌身份可信的访客进来 |
| SASL(鉴权 Authentication) | 前台的人脸识别 / 工卡刷卡 | 验证「你是谁」 |
| ACL(授权 Authorization) | 电梯权限:财务部只能去 12 层 | 验证「你能干什么」 |
| Quota(配额) | 每层楼的载重 / 电力上限 | 防止单个租户用力过猛挤死别人 |
| Audit Log(审计) | 大堂的监控录像 | 事后查谁来过、谁动了什么 |
很多人把「安全」简单理解成「上 SSL」,其实只解决了「门口探测门」一项;如果不配 SASL + ACL + Quota,等于「门口装了金属探测,但任何人进来都能直接刷开总裁办公室的门、还能把整栋楼的电用光」。
本章按这个顺序逐项铺开。
1. 监听器(Listener):一切的起点
任何安全配置都从「Broker 暴露什么端口、用什么协议」开始。server.properties 里的关键三行:
properties
listeners=PLAINTEXT://:9092,SASL_SSL://:9093,SSL://:9094
listener.security.protocol.map=PLAINTEXT:PLAINTEXT,SASL_SSL:SASL_SSL,SSL:SSL,CONTROLLER:PLAINTEXT
inter.broker.listener.name=SASL_SSL含义:
listeners:本 Broker 对外开几个端口、每个端口用什么 安全协议。listener.security.protocol.map:把 listener 名映射到 4 种内置安全协议之一。inter.broker.listener.name:Broker 之间通信走哪个 listener(即副本同步、Controller 心跳走哪个端口)。
四种内置安全协议的组合矩阵:
| 协议 | 加密 | 鉴权 | 适用场景 |
|---|---|---|---|
PLAINTEXT | ❌ | ❌ | 内网 dev、单机测试、彻底信任的环境 |
SSL | ✅ TLS | ✅ 客户端证书(mTLS) | 强加密 + 证书鉴权(Operator 偏爱) |
SASL_PLAINTEXT | ❌ | ✅ SASL(PLAIN/SCRAM/Kerberos…) | 不允许用在公网 / 跨机房!明文密码暴露 |
SASL_SSL | ✅ TLS | ✅ SASL | 生产环境标配 |
记住一条线:只要 listener 跨出 VPC / 上公网,必须 SASL_SSL 或 SSL,否则就是裸奔。
📌 与 RabbitMQ / RocketMQ / Pulsar 对比
- RabbitMQ:5671(AMQPS)+ 15671(管理 SSL),同样基于 TLS + SASL,机制几乎一致。
- RocketMQ:原生只支持 ACL(用户名 + 密码 + 权限),无内置 TLS(需要前置 Nginx)。
- Pulsar:支持 mTLS、JWT、OAuth2、Kerberos,授权基于 namespace,比 Kafka ACL 更易管理。
2. 鉴权(Authentication):你是谁?
「鉴权」回答的是「你这次请求由谁发起」。Kafka 给客户端发的最终结果是一个 KafkaPrincipal(如 User:alice),后面所有 ACL 检查都基于这个 Principal。
Kafka 支持的鉴权机制有 5 种:
2.1 SSL / TLS 双向认证(mTLS)
服务端和客户端互相验证证书,Broker 把客户端证书的 CN 当成 Principal。
关键配置(Broker):
properties
listeners=SSL://:9093
ssl.keystore.location=/etc/kafka/ssl/broker.keystore.jks
ssl.keystore.password=changeit
ssl.key.password=changeit
ssl.truststore.location=/etc/kafka/ssl/broker.truststore.jks
ssl.truststore.password=changeit
ssl.client.auth=required # ★ 必须 required 才是双向认证;none/requested 等于单向
ssl.endpoint.identification.algorithm=https # 防 MITM,校验 hostnamePrincipal 提取规则:默认从证书 Subject 的 CN 取,常自定义为:
properties
ssl.principal.mapping.rules=RULE:^CN=(.*?),.*$/$1/L,DEFAULT含义:把 CN=alice,OU=team,... 的证书映射成 User:alice(小写化)。
优点:无需密码、零依赖外部账号系统、加密 + 鉴权一次到位。 缺点:证书签发、轮转、吊销链路重,CA 跑路 / 证书过期就是大事故;客户端配置门槛高。 典型场景:内部基础设施互联(Broker ↔ Broker、Connect ↔ Broker)、已经有完善 PKI 体系的公司。
2.2 SASL/PLAIN:用户名 + 明文密码
最简单的 SASL 机制,用户名 + 密码以明文发送。只有跑在 SASL_SSL 上才能接受,否则密码就在网线上裸奔。
Broker JAAS 文件(/etc/kafka/kafka_jaas.conf):
text
KafkaServer {
org.apache.kafka.common.security.plain.PlainLoginModule required
username="admin"
password="admin-secret"
user_admin="admin-secret"
user_alice="alice-secret"
user_bob="bob-secret";
};启动 Broker 时通过 JVM 参数指定 JAAS:
bash
export KAFKA_OPTS="-Djava.security.auth.login.config=/etc/kafka/kafka_jaas.conf"
bin/kafka-server-start.sh config/server.properties优点:实现最简单、客户端只需用户名密码,对接老系统快。 缺点:用户和密码硬编码在配置文件里,改密码必须重启 Broker;不适合频繁加用户。 典型场景:测试环境、用户数极少的小集群。
2.3 SASL/SCRAM-SHA-256/SCRAM-SHA-512:带挑战的密码鉴权
SCRAM(Salted Challenge Response Authentication Mechanism)= 密码 + 盐 + 多轮哈希,密码不上线,每次会话用一次性挑战值证明拥有密码。Kafka 支持 SHA-256 / SHA-512 两种摘要长度。
最大优势:用户存在 ZooKeeper / KRaft 元数据里,可动态增删,不需要改 JAAS、不需要重启 Broker。
创建用户:
bash
kafka-configs.sh --bootstrap-server localhost:9092 \
--alter --add-config 'SCRAM-SHA-256=[iterations=4096,password=alice-secret]' \
--entity-type users --entity-name aliceBroker JAAS:
text
KafkaServer {
org.apache.kafka.common.security.scram.ScramLoginModule required
username="admin"
password="admin-secret";
};客户端 JAAS:
text
KafkaClient {
org.apache.kafka.common.security.scram.ScramLoginModule required
username="alice"
password="alice-secret";
};优点:动态增删用户、密码不传输、不需要 PKI;SHA-512 摘要长,安全性足够长期使用。 缺点:仍然需要管理密码(只是把痛苦从 PKI 换成密码库);用户量上千后管理成本上升。 典型场景:最常见的生产选型,是「不上 Kerberos / OAuth 的多租户集群」默认推荐。
2.4 SASL/GSSAPI(Kerberos)
企业级 SSO 方案:客户端先到 KDC 拿 TGT,再用 TGT 换 Kafka 服务的票(Service Ticket),用票完成鉴权。
Broker JAAS:
text
KafkaServer {
com.sun.security.auth.module.Krb5LoginModule required
useKeyTab=true
storeKey=true
keyTab="/etc/security/keytabs/kafka.service.keytab"
principal="kafka/broker1.example.com@EXAMPLE.COM";
};客户端配置:
properties
sasl.mechanism=GSSAPI
sasl.kerberos.service.name=kafka
security.protocol=SASL_SSL优点:和公司现有 AD / Hadoop 生态打通;票据有时效,自动过期;用户管理交给 KDC。 缺点:运维复杂度爆炸(KDC HA、Keytab 分发、principal 映射、时钟同步),新人上手痛苦;Python 客户端需要装 gssapi、pykrb5,跨平台兼容差。 典型场景:大型公司 / 已有 Kerberos 基础设施 / Hadoop 生态深度集成。
2.5 SASL/OAUTHBEARER:OAuth 2.0 Bearer Token
客户端从 OAuth Server 拿到 Bearer Token,通过 SASL 发给 Broker;Broker 用 JWKS 验证 token 签名 + 过期时间。
Broker 端(Kafka 3.x+):
properties
sasl.enabled.mechanisms=OAUTHBEARER
listener.name.sasl_ssl.oauthbearer.sasl.server.callback.handler.class=\
org.apache.kafka.common.security.oauthbearer.OAuthBearerValidatorCallbackHandler
listener.name.sasl_ssl.oauthbearer.sasl.jaas.config=\
org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required;
listener.name.sasl_ssl.oauthbearer.sasl.oauthbearer.jwks.endpoint.url=https://idp.example.com/jwks
listener.name.sasl_ssl.oauthbearer.sasl.oauthbearer.expected.audience=kafka优点:和企业 IDP 打通;细粒度 scope;token 短 TTL;不需要在 Kafka 端管密码。 缺点:需要 IDP(Keycloak / Okta / Auth0);token 刷新、JWKS 缓存、时钟漂移都要踩坑。 典型场景:云原生 / 微服务统一 IAM / Confluent Cloud(默认就是 OAUTHBEARER)。
2.6 五种鉴权机制横向对比
| 机制 | 协议层 | 凭证 | 动态增删用户 | 证书管理 | 复杂度 | 推荐场景 |
|---|---|---|---|---|---|---|
mTLS | TLS | 客户端证书 | 需吊销/轮转 | ★★★★★ | ★★★★ | 内部基础设施互联 |
SASL/PLAIN | SASL | 用户名 + 明文密码 | ❌(改 JAAS 重启) | ☆ | ☆ | 测试 / 小集群 |
SASL/SCRAM | SASL | 用户名 + 挑战哈希 | ✅(kafka-configs.sh) | ☆ | ★★ | 生产首选 |
SASL/GSSAPI | SASL | Kerberos Ticket | ✅(KDC 管理) | ☆ | ★★★★★ | 大企业 / Hadoop 生态 |
SASL/OAUTHBEARER | SASL | OAuth Token | ✅(IDP 管理) | ☆ | ★★★ | 云原生 / 统一 IAM |
记住选型口诀:「内网 mTLS、外网 SCRAM、企业 Kerberos、云原生 OAuth」。本章的代码示例统一用 SASL_SSL + SCRAM-SHA-256(生产最常见组合)。
3. 授权(Authorization):你能干什么?
鉴权回答了「你是谁」,授权回答「你能对哪些东西做什么操作」。Kafka 用 ACL(Access Control List) 表达授权规则。
3.1 启用 ACL
server.properties:
properties
authorizer.class.name=org.apache.kafka.metadata.authorizer.StandardAuthorizer # KRaft 模式
# authorizer.class.name=kafka.security.authorizer.AclAuthorizer # ZK 模式
allow.everyone.if.no.acl.found=false # ★ 关键:默认拒绝;未匹配 ACL 一律拒绝
super.users=User:admin;User:CN=broker # SuperUser 直接绕过 ACL 检查两条铁律:
allow.everyone.if.no.acl.found=false:默认拒绝,体现「白名单原则」。super.users是后门,只能给 Operator 账号,业务账号永远不准列入。
3.2 ACL 五元组
每条 ACL 都是一个五元组:
(Resource Type, Resource Name, Operation, Principal, Permission Type, Host)具体取值:
| 字段 | 可选值 | 说明 |
|---|---|---|
Resource Type | Topic / Group / Cluster / TransactionalId / DelegationToken | 资源种类 |
Resource Name | 字符串,可用 * 或前缀 learn.team-a.* | 资源名(支持 LITERAL / PREFIXED / WILDCARD 三种 Pattern) |
Operation | Read / Write / Create / Delete / Alter / Describe / ClusterAction / IdempotentWrite / All | 动作 |
Principal | User:alice / Group:devs | 鉴权出来的 Principal |
Permission Type | Allow / Deny | Deny 优先级高于 Allow |
Host | IP 字符串或 * | 来源 IP 限制 |
3.3 资源 × 操作 权限矩阵
记住这张表,授权命令就能信手拈来:
| Resource Type \ Operation | Read | Write | Create | Delete | Alter | Describe | ClusterAction | IdempotentWrite |
|---|---|---|---|---|---|---|---|---|
| Topic | 消费 | 生产 | 建 Topic | 删 Topic | 改配置 / 加分区 | 看元数据 | — | — |
| Group | 加入 group / 提交 offset | — | — | 删除 group | — | 看 group 信息 | — | — |
| Cluster | — | — | 建 Topic(与 Topic Create 二选一) | — | 改集群配置 | 看 broker 元数据 | broker 内部通信 | 幂等 Producer 注册 PID |
| TransactionalId | — | 写事务消息 | — | — | — | 查事务状态 | — | — |
记住一个最常见的「订单生产+订单消费」最小授权例子:
bash
# 生产者(写 learn.orders)
kafka-acls.sh --bootstrap-server localhost:9092 \
--add --allow-principal User:order-producer \
--operation Write --operation Describe \
--topic learn.orders
# 幂等 Producer 还需要 ClusterAction → IdempotentWrite
kafka-acls.sh --bootstrap-server localhost:9092 \
--add --allow-principal User:order-producer \
--operation IdempotentWrite \
--cluster
# 消费者组 order-svc(读 learn.orders)
kafka-acls.sh --bootstrap-server localhost:9092 \
--add --allow-principal User:order-consumer \
--operation Read --operation Describe \
--topic learn.orders \
--group order-svc3.4 PREFIXED Pattern:多租户的灵魂
如果给每个 Topic 单独写一条 ACL,运维要疯。Kafka 支持 前缀匹配:
bash
# team-a 团队对所有 learn.team-a.* Topic 都有读写权限
kafka-acls.sh --bootstrap-server localhost:9092 \
--add --allow-principal User:team-a \
--operation Read --operation Write --operation Describe --operation Create \
--topic 'learn.team-a.' --resource-pattern-type PREFIXED这条 ACL 涵盖未来所有 learn.team-a.xxx、learn.team-a.orders.v2,无需改 ACL。这是 Kafka 多租户的核心范式。
3.5 Allow 与 Deny 的优先级
ACL 评估顺序:
- 先看是不是 SuperUser → 是则直接通过;
- 找所有匹配的 Deny → 命中任意一条 → 拒绝;
- 找所有匹配的 Allow → 命中任意一条 → 通过;
- 都不命中 → 看
allow.everyone.if.no.acl.found:true 通过、false 拒绝。
典型用法:先给 User:team-a 整个前缀 Allow,再单独给某个敏感 Topic 加 Deny。
3.6 常用 ACL 命令速查
bash
# 列出某 topic 所有 ACL
kafka-acls.sh --bootstrap-server localhost:9092 --list --topic learn.orders
# 列出某 principal 所有 ACL
kafka-acls.sh --bootstrap-server localhost:9092 --list --principal User:alice
# 删除单条 ACL
kafka-acls.sh --bootstrap-server localhost:9092 \
--remove --allow-principal User:alice --operation Read \
--topic learn.orders
# 限制源 IP
kafka-acls.sh --bootstrap-server localhost:9092 \
--add --allow-principal User:alice --allow-host 10.0.0.10 \
--operation Read --topic learn.orders4. Quota 配额:把高速公路设个限速
问题场景:某个团队代码出 bug,开了 1000 个线程往 Kafka 灌数据,把整个集群的网卡占满,其它租户全部受牵连。这种「软隔离」就是 Quota 的活。
Kafka 提供三种 Quota:
| Quota 类型 | 单位 | 作用 |
|---|---|---|
| Producer 字节速率 | bytes/sec | 限制 Producer 写入速率 |
| Consumer 字节速率 | bytes/sec | 限制 Consumer 拉取速率 |
| Request 速率 | 每秒 Request 处理时间百分比 | 限制元数据 / 控制类请求(防 kafka-topics.sh --describe 类风暴) |
4.1 配额作用对象
Quota 可以挂在三种「实体」上,按优先级匹配:
<user, client-id>(最具体)<user>或<client-id><default>(兜底)
匹配规则:先找最具体的,找不到就用次具体,最后用默认。
4.2 配置 Quota 命令
bash
# 给 user=alice 全局限制:写 10MB/s,读 20MB/s
kafka-configs.sh --bootstrap-server localhost:9092 \
--alter --add-config 'producer_byte_rate=10485760,consumer_byte_rate=20971520' \
--entity-type users --entity-name alice
# 给 client-id=app-x 单独限速
kafka-configs.sh --bootstrap-server localhost:9092 \
--alter --add-config 'producer_byte_rate=5242880' \
--entity-type clients --entity-name app-x
# 给 user=alice + client-id=ingest-job 组合限速(最具体,最优先)
kafka-configs.sh --bootstrap-server localhost:9092 \
--alter --add-config 'producer_byte_rate=20971520' \
--entity-type users --entity-name alice \
--entity-type clients --entity-name ingest-job
# 兜底默认值(所有未单独配置的用户)
kafka-configs.sh --bootstrap-server localhost:9092 \
--alter --add-config 'producer_byte_rate=2097152' \
--entity-type users --entity-default
# Request quota:限制元数据请求占用 broker 时间不超过 200%(即 2 个 IO 线程)
kafka-configs.sh --bootstrap-server localhost:9092 \
--alter --add-config 'request_percentage=200' \
--entity-type users --entity-name alice4.3 超额会发生什么?
Broker 不会丢弃超额请求,而是 故意延后 给客户端响应(throttle)。响应里带个 throttle_time_ms,客户端 SDK 看到后会等待这个时间再发下一批,从而把速率压回阈值。
客户端 JMX 里可以看到
produce-throttle-time-avg/fetch-throttle-time-avg,这两个指标只要 > 0 就说明被限了。
4.4 Quota 与 ACL 的区别
| 维度 | ACL | Quota |
|---|---|---|
| 解决问题 | 「能不能做」 | 「能做多快」 |
| 拒绝行为 | 抛 TopicAuthorizationException | 故意延迟响应(throttle) |
| 目标 | 安全 | 隔离、SLA |
| 必要性 | 安全场景必须 | 多租户 / SLA 场景必须 |
多租户集群最佳实践 = ACL(你能不能用)+ Quota(你能用多少),缺一不可。
5. SuperUser 与审计日志
5.1 SuperUser
super.users=User:admin 配置的用户:
- 绕过所有 ACL(注意:仍然受 Quota 限制);
- 用于 Operator 排障、跨 Topic 拷贝、Connect 框架本身(避免给 Connect 集群配上百条 ACL)。
铁律:
- 只给运维 / 平台账号;
- 凭证存在保密柜(Vault / KMS),不放业务代码;
- 业务账号宁可多写几条 ACL,也不能放进 SuperUser。
5.2 审计日志(Audit Log)
Kafka 自带的 Authorizer 在启用 ACL 后会输出授权日志,但默认级别是 INFO,要打开 DEBUG 才看得到 deny / allow 的明细。
log4j.properties 加:
properties
log4j.logger.kafka.authorizer.logger=DEBUG, authorizerAppender
log4j.logger.kafka.request.logger=INFO, requestAppender
log4j.additivity.kafka.authorizer.logger=false
log4j.appender.authorizerAppender=org.apache.log4j.DailyRollingFileAppender
log4j.appender.authorizerAppender.DatePattern='.'yyyy-MM-dd-HH
log4j.appender.authorizerAppender.File=/var/log/kafka/kafka-authorizer.log
log4j.appender.authorizerAppender.layout=org.apache.log4j.PatternLayout
log4j.appender.authorizerAppender.layout.ConversionPattern=[%d] %p %m (%c)%n输出样例:
[2025-04-17 10:00:01] DEBUG Principal = User:alice is Allowed Operation = Read
from host = 10.0.0.10 on resource = Topic:LITERAL:learn.orders for request = Fetch
[2025-04-17 10:00:02] INFO Principal = User:bob is Denied Operation = Write
from host = 10.0.0.20 on resource = Topic:LITERAL:learn.orders for request = Produce进阶方案:
- Cruise Control / Confluent RBAC 提供更结构化的审计;
- OPA + Kafka Authorizer 插件:用 Open Policy Agent 写策略,决策审计统一上报;
- MirrorMaker / Connect 转储到 Splunk / ELK:收集
kafka-authorizer.log上报 SIEM。
6. 多租户最小权限模型范式
下面是一个真实可落地的 8 步范式,来自多家中大型公司的生产实践:
- 集群隔离 → 团队隔离 → 应用隔离 三级思路:跨业务域用集群隔离;同集群内不同团队用 Topic 前缀 + ACL 隔离;同团队内不同应用用 client-id + Quota 隔离。
- Topic 命名强制前缀:
<env>.<team>.<domain>.<event>,例如prod.team-order.payment.refunded。Schema Registry 也走相同前缀。 - 每个团队 1 个 SCRAM 用户(
team-order-prod/team-order-dev),每个应用 1 个 client-id(payment-svc-1/payment-svc-2)。禁用共享账号。 - PREFIXED ACL 一次到位:给
User:team-order在prod.team-order.前缀上的 Read / Write / Describe / Create。 - 跨团队消费走「申请-审批」工单:A 团队要消费 B 团队的 Topic,B Owner 审批后给 A 加一条 Read ACL;不要给写权限。
- Quota 默认兜底:所有用户 default 写 5MB/s、读 10MB/s;大数据 / 数仓单独申请提额。
- SuperUser 仅留 Operator:
super.users=User:kafka-ops;User:kafka-connect,禁用User:admin这种通用名。 - 审计日志强制保留 90 天:合规 + 事故追溯。
具体落地脚本见 15_security/code/scram_user_setup.sh 和 15_security/code/acl_examples.sh。
7. 生产环境完整配置示例
7.1 Broker server.properties
properties
# ====== 监听器:内部 SASL_SSL,控制面 PLAINTEXT(仅 KRaft 集群内部 IP) ======
listeners=SASL_SSL://0.0.0.0:9093,CONTROLLER://0.0.0.0:9094
advertised.listeners=SASL_SSL://broker1.example.com:9093
listener.security.protocol.map=SASL_SSL:SASL_SSL,CONTROLLER:PLAINTEXT
inter.broker.listener.name=SASL_SSL
controller.listener.names=CONTROLLER
# ====== TLS ======
ssl.keystore.location=/etc/kafka/ssl/broker.keystore.jks
ssl.keystore.password=changeit
ssl.key.password=changeit
ssl.truststore.location=/etc/kafka/ssl/broker.truststore.jks
ssl.truststore.password=changeit
ssl.endpoint.identification.algorithm=https
# ====== SASL ======
sasl.enabled.mechanisms=SCRAM-SHA-256
sasl.mechanism.inter.broker.protocol=SCRAM-SHA-256
listener.name.sasl_ssl.scram-sha-256.sasl.jaas.config=\
org.apache.kafka.common.security.scram.ScramLoginModule required \
username="kafka-broker" \
password="broker-secret";
# ====== ACL ======
authorizer.class.name=org.apache.kafka.metadata.authorizer.StandardAuthorizer
allow.everyone.if.no.acl.found=false
super.users=User:kafka-broker;User:kafka-ops
# ====== 默认 Quota ======
# (用 kafka-configs.sh 设置,不是 server.properties)7.2 客户端 client.properties
properties
bootstrap.servers=broker1.example.com:9093
security.protocol=SASL_SSL
sasl.mechanism=SCRAM-SHA-256
sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required \
username="alice" password="alice-secret";
ssl.truststore.location=/etc/kafka/ssl/client.truststore.jks
ssl.truststore.password=changeit
ssl.endpoint.identification.algorithm=httpsPython confluent-kafka 等价配置:
python
conf = {
"bootstrap.servers": "broker1.example.com:9093",
"security.protocol": "SASL_SSL",
"sasl.mechanism": "SCRAM-SHA-256",
"sasl.username": "alice",
"sasl.password": "alice-secret",
"ssl.ca.location": "/etc/kafka/ssl/ca.crt",
}8. 与其他 MQ 的安全模型对比
| 维度 | Kafka | RabbitMQ | RocketMQ | Pulsar |
|---|---|---|---|---|
| 鉴权 | SSL / SASL(多机制) | SASL PLAIN / EXTERNAL(mTLS) | ACL 用户名密码(4.x+) | mTLS / JWT / OAuth2 / Kerberos |
| 授权 | ACL 五元组 + Pattern | vhost + permission(regex) | Topic / Group ACL | Namespace 级 RBAC |
| 多租户 | Topic 前缀 + ACL + Quota | vhost 强隔离 | 不强 | tenant/namespace 原生 |
| 配额 | Producer / Consumer / Request | per-connection 限流 | 客户端限流 | namespace publish/dispatch rate |
| 审计 | log4j authorizer log | management plugin | 自带审计 | broker audit + Pulsar Manager |
结论:Pulsar 的多租户模型最原生(namespace、租户、隔离都是一等公民);Kafka 通过「前缀 ACL + Quota」也能做到工程级多租户,但要靠 运维 + 命名规范 双管齐下。
9. 常见踩坑
坑 1:SASL_PLAINTEXT 走公网
跨数据中心或公网通信用 SASL_PLAINTEXT = 密码(PLAIN)或 token(OAUTHBEARER)裸奔。只要出 VPC,必须 SASL_SSL。
坑 2:忘了开 ssl.client.auth=required
只配证书但 ssl.client.auth=none,只有加密、没有鉴权,所有人都能匿名访问。
坑 3:allow.everyone.if.no.acl.found=true
配了 ACL 但忘了改这一项 → 所有未匹配 ACL 的资源默认开放,等于没配。
坑 4:幂等 Producer 报 ClusterAuthorizationFailedException
打开 enable.idempotence=true 后,Broker 要求 IdempotentWrite 这个 Cluster 操作权限,常见忘配。
坑 5:消费者只配了 Topic Read,没配 Group Read
报错 GroupAuthorizationFailedException:消费者要 join group / 提交 offset,必须有 Group:Read。
坑 6:SuperUser 漂移
super.users=User:admin 这种通用名被业务方猜到并复用 → 整个 ACL 体系破防。用业务无关的命名 + 凭证保密。
坑 7:Quota 单位易错
producer_byte_rate=1048576 是 1MB/s,不是 1Mb/s(小写 b 是 bit);新人常按 Mbps 算,结果限速差 8 倍。
坑 8:JAAS 改动需要重启
SASL/PLAIN 的用户在 JAAS 文件里,改密码必须重启 broker;想动态加用户,必须 SCRAM。
10. 小结
- 「安全」= 加密(SSL)+ 鉴权(SASL)+ 授权(ACL)+ 限流(Quota)+ 审计,五件事缺一不可。
- 鉴权 5 选 1:生产首选 SCRAM-SHA-256;要 SSO 上 Kerberos 或 OAuth;纯内网零信任可上 mTLS。
- ACL 五元组:(Resource, Operation, Principal, Host, Allow/Deny),Pattern 支持
LITERAL/PREFIXED/WILDCARD,多租户必用 PREFIXED。 - Quota 三类:Producer 字节、Consumer 字节、Request 时间百分比;超额是 throttle 不是 reject。
- SuperUser 是后门,只给 Operator;审计日志默认 INFO,生产要打开 DEBUG。
- 多租户范式:每团队 1 SCRAM 用户 + 1 套 PREFIXED ACL + 兜底 Quota + 命名前缀 + 审计 90 天。
🎯 面试高频题
Q1:Kafka 支持哪些鉴权机制?怎么选?
考察点:安全机制对比、选型能力。
答案:
- 5 种机制:SSL(mTLS 双向证书)、SASL/PLAIN(用户名+明文密码)、SASL/SCRAM-SHA-256/512(盐+多轮哈希)、SASL/GSSAPI(Kerberos)、SASL/OAUTHBEARER(OAuth Token)。
- 选型口诀:
- 内网基础设施互联(broker↔broker、connect↔broker)→ mTLS;
- 跨业务团队多租户 → SCRAM-SHA-256(生产最常见);
- 公司有 Kerberos / Hadoop 生态 → GSSAPI;
- 云原生 / 统一 IAM → OAUTHBEARER;
- 测试环境 → SASL/PLAIN。
- 铁律:跨 VPC / 公网必须
SASL_SSL而不是SASL_PLAINTEXT,否则密码 / token 在网线上裸奔。 - 加分项:提到 SCRAM 用户存在 KRaft 元数据可动态增删、PLAIN 改密码必须重启;Kerberos 运维复杂度爆炸;OAUTHBEARER 是 Confluent Cloud 默认。
Q2:Kafka ACL 的五元组是什么?多租户怎么用 ACL 实现?
考察点:授权模型、多租户落地。
答案:
- 五元组:(Resource Type, Resource Name, Operation, Principal, Permission Type, Host)。其中 Permission Type 有 Allow / Deny,Deny 优先级更高。
- 资源类型:Topic / Group / Cluster / TransactionalId / DelegationToken。
- 关键操作:
- Topic:Read(消费)/ Write(生产)/ Create / Delete / Alter / Describe;
- Group:Read(join+commit)/ Describe / Delete;
- Cluster:ClusterAction(broker 内部)/ IdempotentWrite(幂等 producer)/ Alter / Describe / Create。
- Pattern:LITERAL(精确)/ PREFIXED(前缀)/ WILDCARD(通配)。多租户的核心范式是 PREFIXED:给
User:team-a在learn.team-a.前缀上一次性授权,未来新增 Topic 不用改 ACL。 - 配套:必须设
allow.everyone.if.no.acl.found=false实现白名单;super.users只给 Operator。 - 加分项:提到 Allow 需要业务账号要同时给 Topic + Group ACL;幂等 Producer 还要 Cluster IdempotentWrite;提到 OPA + Kafka Authorizer 插件实现策略化授权。
Q3:Quota 在 Kafka 里有哪几种?Producer 限速触发后会发生什么?
考察点:多租户隔离、限流模型。
答案:
- 三种 Quota:
producer_byte_rate(生产字节速率,bytes/sec);consumer_byte_rate(消费字节速率,bytes/sec);request_percentage(请求占用 broker 时间百分比,控制元数据风暴)。
- 作用对象优先级:
<user, client-id>><user>或<client-id>><default>。 - 超额行为:Broker 不丢消息也不报错,而是 故意延迟 响应,响应里带
throttle_time_ms,SDK 看到后 sleep 这段时间,自然把速率压回去。 - 客户端可观测:JMX
produce-throttle-time-avg/fetch-throttle-time-avg,> 0 即说明被限。 - 与 ACL 的关系:ACL 解决「能不能用」,Quota 解决「能用多少」,二者互补。
- 加分项:提到 Quota 单位是 byte 不是 bit;提到 Confluent 还提供 storage quota(按 GB 限存储);提到「Quota 限的是 broker 端处理速率,client 端实际感知到的是网络往返延迟变高」。
Q4:mTLS 和 SCRAM 各自的优缺点?为什么大多数公司选 SCRAM?
考察点:安全机制权衡。
答案:
- mTLS 优点:无需密码、加密 + 鉴权一次到位、零密码泄漏面;缺点:证书签发 / 轮转 / 吊销链路重,CA 跑路或证书过期 = 大事故;客户端配置门槛高(每应用一份证书)。
- SCRAM 优点:密码不在线上传输(每次会话用挑战哈希);用户存元数据,可通过
kafka-configs.sh动态增删,无需重启 broker;客户端只需用户名 + 密码,门槛低;缺点:仍要管密码(密码库 / 定期轮转);用户量上千后管理成本上升。 - 为什么大多数公司选 SCRAM:
- 已经有现成的密码库 / 工单系统;
- 不用搭 PKI 体系;
- 动态增删用户对多租户场景特别友好;
- 性能足够(SHA-256 / 512 在高吞吐下开销可忽略)。
- mTLS 的典型留场:Broker 之间的内部通信、Connect Cluster ↔ Broker、不希望任何密码出现的合规场景。
- 加分项:提到 SCRAM-SHA-512 安全性更强、SHA-256 兼容性更好;提到双方也可以叠加(SASL_SSL = SCRAM 鉴权 + TLS 加密 + 服务端证书校验)。
Q5:你们生产环境多租户 Kafka 集群的安全设计是怎样的?
考察点:架构设计、最佳实践。
答案(按 8 步范式回答):
- 分层隔离:跨业务域用集群隔离(线上 / 离线 / 数仓),同集群用 Topic 前缀 + ACL 隔离团队,团队内用 client-id + Quota 隔离应用。
- Topic 命名强约束:
<env>.<team>.<domain>.<event>(如prod.team-order.payment.refunded),Schema Registry 也走相同前缀。 - 账号体系:每团队 1 个 SCRAM 用户(
team-order-prod/team-order-dev),每应用 1 个 client-id;禁用共享账号。 - ACL 用 PREFIXED:给
User:team-order在prod.team-order.前缀的 Read/Write/Describe/Create 一次到位。 - 跨团队消费走工单:A 团队要消费 B 团队的 Topic,B Owner 审批后加一条 Read ACL,不给写。
- Quota 默认兜底:default 写 5MB/s、读 10MB/s;大数据 / 数仓单独申请。
- SuperUser 极少:仅
User:kafka-ops+ Connect 框架账号。 - 审计 90 天:
kafka-authorizer.log收集到 ELK / Splunk,合规 + 事故追溯。 - 加分项:提到自动化 ——「申请 Topic」走平台工单 → 自动建 Topic + 自动加 ACL + 自动加监控;提到密码用 Vault 管理 + 定期轮转。
Q6:消费者报 GroupAuthorizationFailedException 是为什么?怎么修?
考察点:ACL 实操、故障排查。
答案:
- 报错原因:消费者 join consumer group / commit offset 时,需要对该 group 有
Read权限,但 ACL 里没给。 - 常见诱因:
- 给消费者只配了
Topic:Read,忘了Group:Read; - group 名字写错了(带了环境前缀 / 大小写不一致);
- 用了 PREFIXED ACL 但 group 名不在前缀范围内。
- 给消费者只配了
- 修复:bash
kafka-acls.sh --bootstrap-server xxx \ --add --allow-principal User:order-consumer \ --operation Read --operation Describe \ --group order-svc - 类比:
Topic:Read是「能进图书馆」,Group:Read是「能在图书馆借阅记录里登记你借了什么书」,没有后者就借不走书。 - 加分项:提到生产中应该把 group 名也加前缀如
prod.team-order.order-svc,统一走 PREFIXED ACL;提到事务 Producer 还要TransactionalId:Write;提到幂等 Producer 要Cluster:IdempotentWrite。
本章配套:
15_security/demo.html:鉴权机制选择决策树 + ACL 五元组建模工具 + Quota 限速动画。15_security/code/scram_user_setup.sh:创建 SCRAM 用户脚本。15_security/code/sasl_producer.py:SASL_SSL + SCRAM 的 Python Producer 示例。15_security/code/acl_examples.sh:常用 ACL 命令集合。15_security/code/jaas.conf.example:客户端 JAAS 配置示例。
🔗 延伸阅读
- 第 4 章 Producer 深入 —— 幂等 Producer 需要的
IdempotentWriteACL。 - 第 13 章 幂等与事务 —— 事务 Producer 需要的
TransactionalId:WriteACL。 - 第 16 章 Kafka Connect —— Connect Worker 的 SuperUser 配置范式。
- 第 19 章 可观测性与运维 —— 审计日志的收集与告警。
🎬 可视化演示
演示加载缓慢或样式异常?点此在新标签页打开 ↗
💻 示例代码
bash
#!/usr/bin/env bash
# ====================================================================
# 第 15 章 - 安全与多租户
# acl_examples.sh - 常用 ACL / Quota 操作集合(生产可直接复用)
# --------------------------------------------------------------------
# 前提:authorizer.class.name=...StandardAuthorizer 已开启
# allow.everyone.if.no.acl.found=false
# ====================================================================
set -euo pipefail
BS=${BOOTSTRAP:-localhost:9092}
CC=${COMMAND_CONFIG:-/etc/kafka/admin.properties}
acl() {
kafka-acls.sh --bootstrap-server "$BS" --command-config "$CC" "$@"
}
cfg() {
kafka-configs.sh --bootstrap-server "$BS" --command-config "$CC" "$@"
}
# --------------------------------------------------------------------
# 0) 列出当前所有 ACL(一上来先看清现状)
# --------------------------------------------------------------------
list_all_acls() { acl --list; }
list_acls_of_user() { acl --list --principal "$1"; }
list_acls_of_topic() { acl --list --topic "$1"; }
# --------------------------------------------------------------------
# 1) 单 Topic 的最小 Producer 权限(含幂等 Producer)
# --------------------------------------------------------------------
grant_producer() {
local user=$1 topic=$2
acl --add --allow-principal "User:${user}" \
--operation Write --operation Describe \
--topic "$topic"
acl --add --allow-principal "User:${user}" \
--operation IdempotentWrite --cluster
}
# --------------------------------------------------------------------
# 2) 单 Topic + 单 Group 的最小 Consumer 权限
# --------------------------------------------------------------------
grant_consumer() {
local user=$1 topic=$2 group=$3
acl --add --allow-principal "User:${user}" \
--operation Read --operation Describe --topic "$topic"
acl --add --allow-principal "User:${user}" \
--operation Read --operation Describe --group "$group"
}
# --------------------------------------------------------------------
# 3) 事务 Producer 额外 TransactionalId 权限
# --------------------------------------------------------------------
grant_tx_producer() {
local user=$1 tx_id_prefix=$2
acl --add --allow-principal "User:${user}" \
--operation Write --operation Describe \
--transactional-id "$tx_id_prefix" \
--resource-pattern-type PREFIXED
}
# --------------------------------------------------------------------
# 4) 多租户:PREFIXED ACL(最常用)
# --------------------------------------------------------------------
grant_team_prefix() {
local team=$1 # team-order
local env=$2 # prod
local prefix="${env}.${team}."
acl --add --allow-principal "User:${team}" \
--operation Read --operation Write --operation Describe \
--operation Create --operation Alter \
--topic "$prefix" --resource-pattern-type PREFIXED
acl --add --allow-principal "User:${team}" \
--operation Read --operation Describe \
--group "$prefix" --resource-pattern-type PREFIXED
acl --add --allow-principal "User:${team}" \
--operation IdempotentWrite --cluster
acl --add --allow-principal "User:${team}" \
--operation Write --operation Describe \
--transactional-id "$prefix" --resource-pattern-type PREFIXED
}
# --------------------------------------------------------------------
# 5) 跨团队消费
# --------------------------------------------------------------------
grant_cross_team_read() {
local consumer_user=$1 source_topic=$2 consumer_group=$3
acl --add --allow-principal "User:${consumer_user}" \
--operation Read --operation Describe --topic "$source_topic"
acl --add --allow-principal "User:${consumer_user}" \
--operation Read --operation Describe --group "$consumer_group"
}
# --------------------------------------------------------------------
# 6) Deny(优先级 > Allow)
# --------------------------------------------------------------------
deny_topic_for_user() {
local user=$1 topic=$2
acl --add --deny-principal "User:${user}" \
--operation All --topic "$topic"
}
# --------------------------------------------------------------------
# 7) Quota
# --------------------------------------------------------------------
set_user_quota() {
local user=$1 producer_mb=$2 consumer_mb=$3
cfg --alter \
--add-config "producer_byte_rate=$((producer_mb*1024*1024)),consumer_byte_rate=$((consumer_mb*1024*1024))" \
--entity-type users --entity-name "$user"
}
set_default_quota() {
local producer_mb=$1 consumer_mb=$2
cfg --alter \
--add-config "producer_byte_rate=$((producer_mb*1024*1024)),consumer_byte_rate=$((consumer_mb*1024*1024))" \
--entity-type users --entity-default
}
set_request_quota() {
local user=$1 percentage=$2
cfg --alter --add-config "request_percentage=${percentage}" \
--entity-type users --entity-name "$user"
}
list_quotas() {
cfg --describe --entity-type users
}
# --------------------------------------------------------------------
# 8) 收回权限
# --------------------------------------------------------------------
revoke_topic_write() {
local user=$1 topic=$2
acl --remove --allow-principal "User:${user}" \
--operation Write --topic "$topic"
}
# ====================================================================
# Demo:组装一个完整的多租户 onboarding 流程
# ====================================================================
onboard_team_demo() {
local team=team-order
local env=prod
echo "[1] 创建团队 SCRAM 用户"
./scram_user_setup.sh add "$team" "$(openssl rand -base64 18)"
echo "[2] 给团队 PREFIXED ACL"
grant_team_prefix "$team" "$env"
echo "[3] 给团队 Quota(写 10MB/s 读 20MB/s)"
set_user_quota "$team" 10 20
echo "[4] 设置默认兜底 Quota(写 5MB/s 读 10MB/s)"
set_default_quota 5 10
echo "[5] 验收:列出 ACL + Quota"
list_acls_of_user "User:${team}"
cfg --describe --entity-type users --entity-name "$team"
}
# ====================================================================
# 命令分发
# ====================================================================
case "${1:-help}" in
list) list_all_acls ;;
list-user) list_acls_of_user "${2:?need User:xxx}" ;;
list-topic) list_acls_of_topic "${2:?need topic}" ;;
grant-producer) grant_producer "$2" "$3" ;;
grant-consumer) grant_consumer "$2" "$3" "$4" ;;
grant-tx) grant_tx_producer "$2" "$3" ;;
grant-team) grant_team_prefix "$2" "$3" ;;
grant-cross) grant_cross_team_read "$2" "$3" "$4" ;;
deny) deny_topic_for_user "$2" "$3" ;;
quota) set_user_quota "$2" "$3" "$4" ;;
quota-default) set_default_quota "$2" "$3" ;;
quota-request) set_request_quota "$2" "$3" ;;
list-quotas) list_quotas ;;
revoke-write) revoke_topic_write "$2" "$3" ;;
onboard-demo) onboard_team_demo ;;
*)
cat <<EOF
Usage: $0 <subcommand> [args]
list 列出所有 ACL
list-user User:alice 列出某用户 ACL
list-topic learn.orders 列出某 Topic ACL
grant-producer alice learn.orders 授权 Producer
grant-consumer alice learn.orders order-svc
grant-tx alice order-tx- 授权事务 Producer(前缀)
grant-team team-order prod 团队 PREFIXED 一键授权
grant-cross team-b prod.team-a.payment.refunded prod.team-b.watch
deny bob learn.secret Deny 全部操作
quota alice 10 20 单用户 Quota(producer/consumer MB/s)
quota-default 5 10 默认兜底 Quota
quota-request alice 200 Request 时间百分比
list-quotas 列出所有用户 Quota
revoke-write alice learn.orders 收回写权限
onboard-demo 演示多租户 onboarding 全流程
EOF
;;
esactxt
// ====================================================================
// 第 15 章 - 客户端 JAAS 配置示例
// 适用:Java 客户端、kafka-*.sh 命令行、Kafka Connect、ksqlDB Server
//
// 使用方式(任选其一):
// 1) JVM 参数:
// export KAFKA_OPTS="-Djava.security.auth.login.config=/path/to/jaas.conf"
// bin/kafka-console-producer.sh ...
//
// 2) 直接写在 producer.properties(推荐):
// sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule \
// required username="alice" password="alice-secret";
//
// Python confluent-kafka 不读 JAAS 文件,直接用 sasl.username / sasl.password
// ====================================================================
// ===== (A) Broker 端:服务端配置 + inter-broker =====
KafkaServer {
org.apache.kafka.common.security.scram.ScramLoginModule required
username="kafka-broker"
password="broker-secret-CHANGE-ME";
};
// ===== (B) 普通客户端:SASL/SCRAM-SHA-256(生产首选) =====
KafkaClient {
org.apache.kafka.common.security.scram.ScramLoginModule required
username="alice"
password="alice-secret-CHANGE-ME";
};
// ===== (C) SASL/PLAIN(仅测试) =====
// KafkaClient {
// org.apache.kafka.common.security.plain.PlainLoginModule required
// username="alice"
// password="alice-secret";
// };
// ===== (D) SASL/GSSAPI(Kerberos) =====
// 需要先 kinit 拿 TGT,或 useKeyTab=true + keyTab 路径
// KafkaClient {
// com.sun.security.auth.module.Krb5LoginModule required
// useKeyTab=true
// storeKey=true
// keyTab="/etc/security/keytabs/alice.keytab"
// principal="alice@EXAMPLE.COM";
// };
// ===== (E) SASL/OAUTHBEARER =====
// KafkaClient {
// org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required
// oauth.token.endpoint.uri="https://idp.example.com/protocol/openid-connect/token"
// oauth.client.id="alice-svc"
// oauth.client.secret="xxxxxxxx";
// };
// ====================================================================
// 配套 client.properties 示例:
//
// security.protocol=SASL_SSL
// sasl.mechanism=SCRAM-SHA-256
// sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule \
// required username="alice" password="alice-secret";
// ssl.truststore.location=/etc/kafka/ssl/client.truststore.jks
// ssl.truststore.password=changeit
// ssl.endpoint.identification.algorithm=https
//
// 命令行调用:
// kafka-console-producer.sh --bootstrap-server localhost:9093 \
// --topic learn.security.demo \
// --producer.config client.properties
// ====================================================================python
"""
第 15 章 · 安全与多租户
sasl_producer.py - 用 SASL_SSL + SCRAM-SHA-256 连接 Kafka 的 Python Producer
依赖:confluent-kafka >= 2.5
pip install confluent-kafka
前提:
1. Broker 已开启 SASL_SSL + SCRAM-SHA-256
2. 已用 scram_user_setup.sh 创建用户 alice,密码 alice-secret
3. 已用 acl_examples.sh 给 User:alice 加上 learn.security.demo Topic 的 Write 权限
用法:
python sasl_producer.py
"""
from __future__ import annotations
import json
import os
import time
from confluent_kafka import Producer
TOPIC = "learn.security.demo"
def build_producer() -> Producer:
conf = {
# ------- 基础 -------
"bootstrap.servers": os.getenv("KAFKA_BOOTSTRAP", "localhost:9093"),
"client.id": "sasl-producer-demo",
# ------- 安全协议 -------
# SASL_SSL = TLS 加密 + SASL 鉴权
"security.protocol": "SASL_SSL",
# ------- SASL 机制 -------
"sasl.mechanism": "SCRAM-SHA-256",
"sasl.username": os.getenv("KAFKA_USERNAME", "alice"),
"sasl.password": os.getenv("KAFKA_PASSWORD", "alice-secret"),
# ------- TLS(验证 broker 证书;mTLS 时还需 keystore) -------
# CA 证书:用于校验 broker 证书是否由可信 CA 签发
"ssl.ca.location": os.getenv("KAFKA_CA_CERT", "/etc/kafka/ssl/ca.crt"),
# 强烈建议开启 hostname 校验,防中间人;与 broker 端
# ssl.endpoint.identification.algorithm=https 配套
"ssl.endpoint.identification.algorithm": "https",
# ------- 可靠性 -------
"acks": "all",
"enable.idempotence": True, # ★ 注意:需要给账号 Cluster:IdempotentWrite ACL
"compression.type": "zstd",
"linger.ms": 20,
"batch.size": 64 * 1024,
"retries": 10,
"delivery.timeout.ms": 120_000,
}
return Producer(conf)
def delivery_report(err, msg):
if err is not None:
# 鉴权失败时,常见错误码:
# _AUTHENTICATION(broker 拒绝凭证)
# _ALL_BROKERS_DOWN(TLS 握手失败时也可能报这个)
print(f"❌ Delivery failed: {err}")
else:
print(f"✅ Delivered to {msg.topic()}[{msg.partition()}]@{msg.offset()}: "
f"key={msg.key()!r} val={msg.value()!r}")
def main() -> None:
producer = build_producer()
print(f"🔐 Connecting as user={os.getenv('KAFKA_USERNAME', 'alice')} "
f"to {os.getenv('KAFKA_BOOTSTRAP', 'localhost:9093')} via SASL_SSL+SCRAM-SHA-256")
for i in range(10):
payload = {"order_id": i, "ts": time.time(), "msg": f"hello-secure-{i}"}
producer.produce(
topic=TOPIC,
key=str(i).encode(),
value=json.dumps(payload).encode(),
on_delivery=delivery_report,
)
# 触发回调(不是 flush)
producer.poll(0)
time.sleep(0.2)
# 退出前 flush,确保所有 in-flight 消息全部 ack 或 error
remaining = producer.flush(10)
if remaining:
print(f"⚠️ Still {remaining} messages in queue after flush timeout")
else:
print("🎉 All messages delivered")
if __name__ == "__main__":
main()
# ============================================================================
# 故障排查指南
# ----------------------------------------------------------------------------
# 1) AUTHENTICATION_FAILED / SASL authentication failed
# → 用户名 / 密码错;或者 broker 端没启用对应的 SASL 机制
#
# 2) UNKNOWN_SERVER_ERROR + broker log: SSL handshake failed
# → ca.location 不对 / broker 证书过期 / 时间不同步
#
# 3) TOPIC_AUTHORIZATION_FAILED
# → 用户没有 Topic:Write 的 ACL
# kafka-acls.sh ... --add --allow-principal User:alice \
# --operation Write --operation Describe --topic learn.security.demo
#
# 4) CLUSTER_AUTHORIZATION_FAILED(仅 enable.idempotence=True 时出现)
# → 用户没有 Cluster:IdempotentWrite 的 ACL
# kafka-acls.sh ... --add --allow-principal User:alice \
# --operation IdempotentWrite --cluster
#
# 5) GROUP_AUTHORIZATION_FAILED(消费者端,本脚本不涉及)
# → 用户没有 Group:Read
# ============================================================================bash
#!/usr/bin/env bash
# ====================================================================
# 第 15 章 · 安全与多租户
# scram_user_setup.sh - 创建 / 修改 / 删除 SCRAM-SHA-256 用户
# --------------------------------------------------------------------
# 前提:
# 1. Broker 已开启 SASL_SSL + SCRAM-SHA-256(参见 server.properties)
# 2. 当前 shell 有一个有 SuperUser 权限的客户端配置 admin.properties
#
# 用法:
# ./scram_user_setup.sh add alice alice-secret
# ./scram_user_setup.sh update bob new-secret
# ./scram_user_setup.sh list alice
# ./scram_user_setup.sh delete bob
# ====================================================================
set -euo pipefail
BOOTSTRAP=${BOOTSTRAP:-localhost:9092}
COMMAND_CONFIG=${COMMAND_CONFIG:-/etc/kafka/admin.properties}
# admin.properties 示例(SuperUser 凭证):
# bootstrap.servers=localhost:9092
# security.protocol=SASL_SSL
# sasl.mechanism=SCRAM-SHA-256
# sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required \
# username="admin" password="admin-secret";
# ssl.truststore.location=/etc/kafka/ssl/client.truststore.jks
# ssl.truststore.password=changeit
action=${1:-help}
user=${2:-}
pass=${3:-}
usage() {
cat <<EOF
Usage:
$0 add <user> <password> # 创建用户(SCRAM-SHA-256, iterations=4096)
$0 update <user> <password> # 修改密码
$0 list [user] # 列出(指定用户或全部)
$0 delete <user> # 删除
EOF
exit 1
}
case "$action" in
add|update)
[[ -z "$user" || -z "$pass" ]] && usage
echo "🔧 ${action} SCRAM-SHA-256 user: $user"
kafka-configs.sh --bootstrap-server "$BOOTSTRAP" --command-config "$COMMAND_CONFIG" \
--alter --add-config "SCRAM-SHA-256=[iterations=4096,password=${pass}]" \
--entity-type users --entity-name "$user"
echo "✅ done"
;;
list)
if [[ -z "$user" ]]; then
echo "📋 列出全部 SCRAM 用户"
kafka-configs.sh --bootstrap-server "$BOOTSTRAP" --command-config "$COMMAND_CONFIG" \
--describe --entity-type users
else
echo "📋 用户 $user 的 SCRAM 配置"
kafka-configs.sh --bootstrap-server "$BOOTSTRAP" --command-config "$COMMAND_CONFIG" \
--describe --entity-type users --entity-name "$user"
fi
;;
delete)
[[ -z "$user" ]] && usage
echo "❌ 删除 SCRAM-SHA-256 user: $user"
kafka-configs.sh --bootstrap-server "$BOOTSTRAP" --command-config "$COMMAND_CONFIG" \
--alter --delete-config 'SCRAM-SHA-256' \
--entity-type users --entity-name "$user"
echo "✅ done"
;;
*)
usage
;;
esac
# ====================================================================
# 一些常用片段(生产中常组合使用)
# --------------------------------------------------------------------
#
# 1) 创建一个团队主账号 + 一组 client-id
# ./scram_user_setup.sh add team-order $(openssl rand -base64 18)
#
# 2) 配套 PREFIXED ACL(参见 acl_examples.sh)
# kafka-acls.sh ... --add --allow-principal User:team-order \
# --operation Read --operation Write --operation Describe --operation Create \
# --topic 'prod.team-order.' --resource-pattern-type PREFIXED
#
# 3) 配套 Quota(参见 acl_examples.sh)
# kafka-configs.sh ... --alter \
# --add-config 'producer_byte_rate=10485760,consumer_byte_rate=20971520' \
# --entity-type users --entity-name team-order
# ====================================================================acl_examples.sh ↗ · jaas.conf.example ↗ · sasl_producer.py ↗ · scram_user_setup.sh ↗