Skip to content

第 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_SSLSSL,否则就是裸奔。

📌 与 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,校验 hostname

Principal 提取规则:默认从证书 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 alice

Broker 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 客户端需要装 gssapipykrb5,跨平台兼容差。 典型场景:大型公司 / 已有 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 五种鉴权机制横向对比

机制协议层凭证动态增删用户证书管理复杂度推荐场景
mTLSTLS客户端证书需吊销/轮转★★★★★★★★★内部基础设施互联
SASL/PLAINSASL用户名 + 明文密码❌(改 JAAS 重启)测试 / 小集群
SASL/SCRAMSASL用户名 + 挑战哈希✅(kafka-configs.sh★★生产首选
SASL/GSSAPISASLKerberos Ticket✅(KDC 管理)★★★★★大企业 / Hadoop 生态
SASL/OAUTHBEARERSASLOAuth 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 检查

两条铁律

  1. allow.everyone.if.no.acl.found=false:默认拒绝,体现「白名单原则」。
  2. super.users 是后门,只能给 Operator 账号,业务账号永远不准列入。

3.2 ACL 五元组

每条 ACL 都是一个五元组:

(Resource Type, Resource Name, Operation, Principal, Permission Type, Host)

具体取值:

字段可选值说明
Resource TypeTopic / Group / Cluster / TransactionalId / DelegationToken资源种类
Resource Name字符串,可用 * 或前缀 learn.team-a.*资源名(支持 LITERAL / PREFIXED / WILDCARD 三种 Pattern)
OperationRead / Write / Create / Delete / Alter / Describe / ClusterAction / IdempotentWrite / All动作
PrincipalUser:alice / Group:devs鉴权出来的 Principal
Permission TypeAllow / DenyDeny 优先级高于 Allow
HostIP 字符串或 *来源 IP 限制

3.3 资源 × 操作 权限矩阵

记住这张表,授权命令就能信手拈来:

Resource Type \ OperationReadWriteCreateDeleteAlterDescribeClusterActionIdempotentWrite
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-svc

3.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.xxxlearn.team-a.orders.v2无需改 ACL。这是 Kafka 多租户的核心范式。

3.5 Allow 与 Deny 的优先级

ACL 评估顺序:

  1. 先看是不是 SuperUser → 是则直接通过;
  2. 找所有匹配的 Deny → 命中任意一条 → 拒绝;
  3. 找所有匹配的 Allow → 命中任意一条 → 通过;
  4. 都不命中 → 看 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.orders

4. 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 可以挂在三种「实体」上,按优先级匹配:

  1. <user, client-id>(最具体)
  2. <user><client-id>
  3. <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 alice

4.3 超额会发生什么?

Broker 不会丢弃超额请求,而是 故意延后 给客户端响应(throttle)。响应里带个 throttle_time_ms,客户端 SDK 看到后会等待这个时间再发下一批,从而把速率压回阈值。

客户端 JMX 里可以看到 produce-throttle-time-avg / fetch-throttle-time-avg,这两个指标只要 > 0 就说明被限了。

4.4 Quota 与 ACL 的区别

维度ACLQuota
解决问题「能不能做」「能做多快」
拒绝行为TopicAuthorizationException故意延迟响应(throttle)
目标安全隔离、SLA
必要性安全场景必须多租户 / SLA 场景必须

多租户集群最佳实践 = ACL(你能不能用)+ Quota(你能用多少),缺一不可。


5. SuperUser 与审计日志

5.1 SuperUser

super.users=User:admin 配置的用户:

  • 绕过所有 ACL(注意:仍然受 Quota 限制);
  • 用于 Operator 排障、跨 Topic 拷贝、Connect 框架本身(避免给 Connect 集群配上百条 ACL)。

铁律

  1. 只给运维 / 平台账号;
  2. 凭证存在保密柜(Vault / KMS),不放业务代码;
  3. 业务账号宁可多写几条 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 步范式,来自多家中大型公司的生产实践:

  1. 集群隔离 → 团队隔离 → 应用隔离 三级思路:跨业务域用集群隔离;同集群内不同团队用 Topic 前缀 + ACL 隔离;同团队内不同应用用 client-id + Quota 隔离。
  2. Topic 命名强制前缀<env>.<team>.<domain>.<event>,例如 prod.team-order.payment.refunded。Schema Registry 也走相同前缀。
  3. 每个团队 1 个 SCRAM 用户team-order-prod / team-order-dev),每个应用 1 个 client-idpayment-svc-1 / payment-svc-2)。禁用共享账号
  4. PREFIXED ACL 一次到位:给 User:team-orderprod.team-order. 前缀上的 Read / Write / Describe / Create。
  5. 跨团队消费走「申请-审批」工单:A 团队要消费 B 团队的 Topic,B Owner 审批后给 A 加一条 Read ACL;不要给写权限。
  6. Quota 默认兜底:所有用户 default 写 5MB/s、读 10MB/s;大数据 / 数仓单独申请提额。
  7. SuperUser 仅留 Operatorsuper.users=User:kafka-ops;User:kafka-connect,禁用 User:admin 这种通用名。
  8. 审计日志强制保留 90 天:合规 + 事故追溯。

具体落地脚本见 15_security/code/scram_user_setup.sh15_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=https

Python 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 的安全模型对比

维度KafkaRabbitMQRocketMQPulsar
鉴权SSL / SASL(多机制)SASL PLAIN / EXTERNAL(mTLS)ACL 用户名密码(4.x+)mTLS / JWT / OAuth2 / Kerberos
授权ACL 五元组 + Patternvhost + permission(regex)Topic / Group ACLNamespace 级 RBAC
多租户Topic 前缀 + ACL + Quotavhost 强隔离不强tenant/namespace 原生
配额Producer / Consumer / Requestper-connection 限流客户端限流namespace publish/dispatch rate
审计log4j authorizer logmanagement 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 支持哪些鉴权机制?怎么选?

考察点:安全机制对比、选型能力。

答案

  1. 5 种机制:SSL(mTLS 双向证书)、SASL/PLAIN(用户名+明文密码)、SASL/SCRAM-SHA-256/512(盐+多轮哈希)、SASL/GSSAPI(Kerberos)、SASL/OAUTHBEARER(OAuth Token)。
  2. 选型口诀
    • 内网基础设施互联(broker↔broker、connect↔broker)→ mTLS
    • 跨业务团队多租户 → SCRAM-SHA-256(生产最常见);
    • 公司有 Kerberos / Hadoop 生态 → GSSAPI
    • 云原生 / 统一 IAM → OAUTHBEARER
    • 测试环境 → SASL/PLAIN。
  3. 铁律:跨 VPC / 公网必须 SASL_SSL 而不是 SASL_PLAINTEXT,否则密码 / token 在网线上裸奔。
  4. 加分项:提到 SCRAM 用户存在 KRaft 元数据可动态增删、PLAIN 改密码必须重启;Kerberos 运维复杂度爆炸;OAUTHBEARER 是 Confluent Cloud 默认。

Q2:Kafka ACL 的五元组是什么?多租户怎么用 ACL 实现?

考察点:授权模型、多租户落地。

答案

  1. 五元组:(Resource Type, Resource Name, Operation, Principal, Permission Type, Host)。其中 Permission Type 有 Allow / Deny,Deny 优先级更高
  2. 资源类型:Topic / Group / Cluster / TransactionalId / DelegationToken。
  3. 关键操作
    • Topic:Read(消费)/ Write(生产)/ Create / Delete / Alter / Describe;
    • Group:Read(join+commit)/ Describe / Delete;
    • Cluster:ClusterAction(broker 内部)/ IdempotentWrite(幂等 producer)/ Alter / Describe / Create。
  4. Pattern:LITERAL(精确)/ PREFIXED(前缀)/ WILDCARD(通配)。多租户的核心范式是 PREFIXED:给 User:team-alearn.team-a. 前缀上一次性授权,未来新增 Topic 不用改 ACL。
  5. 配套:必须设 allow.everyone.if.no.acl.found=false 实现白名单;super.users 只给 Operator。
  6. 加分项:提到 Allow 需要业务账号要同时给 Topic + Group ACL;幂等 Producer 还要 Cluster IdempotentWrite;提到 OPA + Kafka Authorizer 插件实现策略化授权。

Q3:Quota 在 Kafka 里有哪几种?Producer 限速触发后会发生什么?

考察点:多租户隔离、限流模型。

答案

  1. 三种 Quota
    • producer_byte_rate(生产字节速率,bytes/sec);
    • consumer_byte_rate(消费字节速率,bytes/sec);
    • request_percentage(请求占用 broker 时间百分比,控制元数据风暴)。
  2. 作用对象优先级<user, client-id> > <user><client-id> > <default>
  3. 超额行为:Broker 不丢消息也不报错,而是 故意延迟 响应,响应里带 throttle_time_ms,SDK 看到后 sleep 这段时间,自然把速率压回去。
  4. 客户端可观测:JMX produce-throttle-time-avg / fetch-throttle-time-avg,> 0 即说明被限。
  5. 与 ACL 的关系:ACL 解决「能不能用」,Quota 解决「能用多少」,二者互补。
  6. 加分项:提到 Quota 单位是 byte 不是 bit;提到 Confluent 还提供 storage quota(按 GB 限存储);提到「Quota 限的是 broker 端处理速率,client 端实际感知到的是网络往返延迟变高」。

Q4:mTLS 和 SCRAM 各自的优缺点?为什么大多数公司选 SCRAM?

考察点:安全机制权衡。

答案

  1. mTLS 优点:无需密码、加密 + 鉴权一次到位、零密码泄漏面;缺点:证书签发 / 轮转 / 吊销链路重,CA 跑路或证书过期 = 大事故;客户端配置门槛高(每应用一份证书)。
  2. SCRAM 优点:密码不在线上传输(每次会话用挑战哈希);用户存元数据,可通过 kafka-configs.sh 动态增删,无需重启 broker;客户端只需用户名 + 密码,门槛低;缺点:仍要管密码(密码库 / 定期轮转);用户量上千后管理成本上升。
  3. 为什么大多数公司选 SCRAM
    • 已经有现成的密码库 / 工单系统;
    • 不用搭 PKI 体系;
    • 动态增删用户对多租户场景特别友好;
    • 性能足够(SHA-256 / 512 在高吞吐下开销可忽略)。
  4. mTLS 的典型留场:Broker 之间的内部通信、Connect Cluster ↔ Broker、不希望任何密码出现的合规场景。
  5. 加分项:提到 SCRAM-SHA-512 安全性更强、SHA-256 兼容性更好;提到双方也可以叠加(SASL_SSL = SCRAM 鉴权 + TLS 加密 + 服务端证书校验)。

Q5:你们生产环境多租户 Kafka 集群的安全设计是怎样的?

考察点:架构设计、最佳实践。

答案(按 8 步范式回答)

  1. 分层隔离:跨业务域用集群隔离(线上 / 离线 / 数仓),同集群用 Topic 前缀 + ACL 隔离团队,团队内用 client-id + Quota 隔离应用。
  2. Topic 命名强约束<env>.<team>.<domain>.<event>(如 prod.team-order.payment.refunded),Schema Registry 也走相同前缀。
  3. 账号体系:每团队 1 个 SCRAM 用户(team-order-prod / team-order-dev),每应用 1 个 client-id;禁用共享账号。
  4. ACL 用 PREFIXED:给 User:team-orderprod.team-order. 前缀的 Read/Write/Describe/Create 一次到位。
  5. 跨团队消费走工单:A 团队要消费 B 团队的 Topic,B Owner 审批后加一条 Read ACL,不给写
  6. Quota 默认兜底:default 写 5MB/s、读 10MB/s;大数据 / 数仓单独申请。
  7. SuperUser 极少:仅 User:kafka-ops + Connect 框架账号。
  8. 审计 90 天kafka-authorizer.log 收集到 ELK / Splunk,合规 + 事故追溯。
  9. 加分项:提到自动化 ——「申请 Topic」走平台工单 → 自动建 Topic + 自动加 ACL + 自动加监控;提到密码用 Vault 管理 + 定期轮转。

Q6:消费者报 GroupAuthorizationFailedException 是为什么?怎么修?

考察点:ACL 实操、故障排查。

答案

  1. 报错原因:消费者 join consumer group / commit offset 时,需要对该 group 有 Read 权限,但 ACL 里没给。
  2. 常见诱因
    • 给消费者只配了 Topic:Read,忘了 Group:Read
    • group 名字写错了(带了环境前缀 / 大小写不一致);
    • 用了 PREFIXED ACL 但 group 名不在前缀范围内。
  3. 修复
    bash
    kafka-acls.sh --bootstrap-server xxx \
      --add --allow-principal User:order-consumer \
      --operation Read --operation Describe \
      --group order-svc
  4. 类比Topic:Read 是「能进图书馆」,Group:Read 是「能在图书馆借阅记录里登记你借了什么书」,没有后者就借不走书。
  5. 加分项:提到生产中应该把 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 配置示例。

🔗 延伸阅读

🎬 可视化演示

演示加载缓慢或样式异常?点此在新标签页打开 ↗

💻 示例代码

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
    ;;
esac
txt
// ====================================================================
// 第 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 ↗