Kafka跑上生产环境后,最逃不掉的一件事就是给客户端加上认证和传输加密。很多团队初期图方便,直接裸用PLAINTEXT,等安全审计或者多团队共用集群时,才开始补救。我最近正好把一套三节点的Kafka集群完整配置了SASL_SSL,从证书体系搭建到broker端改造,再到Java客户端和命令行工具接入,整个过程踩了不少坑,也积累了一些经验。这篇东西不是官方文档复读,而是把配置SASL_SSL认证和传输加密的完整决策过程、操作细节、验证方法,以及几个最容易翻车的地方都讲清楚。适合正在做Kafka安全改造、或者准备从零搭建带认证加密集群的工程师参考。
1. 项目背景与需求拆解
1.1 为什么Kafka生产环境必须考虑SASL_SSL
Kafka默认的PLAINTEXT监听模式,所有消息都是裸奔的。如果你把Kafka部署在不可信网络里,任何一个能连通9092端口的人,都能直接拉取数据、篡改topic内容。就算只在公司内网用,一旦有多团队共享集群,你也没法确认写入数据的人到底是谁。这就是两个独立的问题:一个是数据在网络上传输时的机密性问题,另一个是访问者身份的可信性问题。
传输加密解决的是前者,用SSL/TLS把broker和客户端之间的流量包加密;认证解决的是后者,用SASL机制确认对方是不是有权限的用户。两者组合在一起,就是SASL_SSL这个listener协议。它同时实现了“你是谁”和“传输内容不被窃听”这两个目标,也是目前Kafka生产部署里比较推荐的安全基线。
实际项目中我还遇到过更麻烦的场景:公司安全合规要求所有敏感数据的传输必须加密,否则不能上生产。这时候SASL_SSL就不是可选项,而是硬性门槛。所以如果你现在还在规划Kafka集群,建议从一开始就把SASL_SSL设计进去,省得后面再对整个集群做安全改造,那会牵扯到所有客户端同时升级,相当折腾。
1.2 认证和传输加密各自到底管什么
很多人容易混淆这两个概念。SASL(Simple Authentication and Security Layer)是一套认证框架,Kafka用它来验证客户端的身份。认证成功后,broker才知道“你是你”,然后根据ACL判断你能干什么。SSL/TLS则负责保护通道,防止中途被嗅探或者篡改。只有加密没有认证,相当于你有个保险箱,但谁都能打开往里放东西;只有认证没有加密,相当于门卫查了身份证,但院内说的话隔墙都能听到。
在Kafka的listener里,你可以单独配置SASL_PLAINTEXT,那就是只认证不加密;也可以单独配置SSL,那是只加密不认证(服务端证书验证不等同于客户端身份认证,除非用mTLS)。SASL_SSL则是把两者都串起来。还有一点,broker和broker之间的通信也可以走SASL_SSL,集群内部的复制流量同样需要保护。很多人在配置客户端时很上心,却忘了内部broker之间的通道,这是安全盲区。
1.3 SASL机制选型:为什么我强烈推荐SCRAM
Kafka支持的SASL机制有好几种:PLAIN、SCRAM-SHA-256、SCRAM-SHA-512、GSSAPI(Kerberos)、OAUTHBEARER等。PLAIN最简单,本质就是明文传输用户名密码,但注意它只是认证时用明文用户名密码,不代表传输不加密,配上SSL之后整体仍然是安全的。不过PLAIN有个硬伤:用户信息得写在JAAS文件或ZooKeeper(旧版)里,管理起来非常笨重,改个密码要动broker。
SCRAM-SHA-256/512是目前我用下来最合适的机制。它同样是用户名密码模型,但密码不以明文存储,服务端保存的是盐和迭代哈希,认证时通过质询-响应流程完成,且Kafka提供了动态创建和修改用户的工具。这就解决了运维痛点:我可以随时用命令新增一个用户,给某个应用单独分配账号,不需要重启broker。GSSAPI(Kerberos)功能强大,但部署复杂度高,要是公司没有现成的Kerberos域,不建议轻易碰。OAUTHBEARER在Kafka 2.x以后也逐步成熟,适合和统一认证中心对接,但需要自己写登录回调,重一点。
所以这篇博文我以SCRAM-SHA-256为例展开。版本上我用的Kafka 3.3.1,Java 11,三节点集群,认证机制选的是SCRAM-SHA-256,传输加密用自建内部CA签发的服务端证书。
2. 证书体系与信任关系搭建
2.1 准备工具链和目录结构
搭建SASL_SSL第一步不是改Kafka配置,而是先把证书体系搞定。我们需要给每一个broker单独签发一张证书,这张证书必须包含broker的hostname或IP,最好用SAN(Subject Alternative Name)而不是Common Name,因为新版本Java已经强制校验SAN了。我习惯在每台broker机器上建一个专门的目录,比如/etc/kafka/security,里面放keystore和truststore。
这里需要决策一件事:是用Java自带的keytool生成JKS,还是用openssl生成PEM再导入。个人推荐用openssl生成CA和私钥,再用keytool生成keystore并导入证书,因为这样CA管理更自由,也方便后续给其他组件签证书。如果你公司有内部CA服务,可以省略手动生成CA的步骤,直接从CA签发broker证书即可。
还要准备一个统一的生成脚本,放在ansible或者shell脚本里,保证每台broker的证书流程一致。脚本里需要指定的关键信息包括:broker的hostname(必须和客户端实际访问的hostname一致)、组织名、证书有效期。有效期我建议签一年,太长了有风险,太短了运维烦。
2.2 生成内部CA和Broker证书
先生成CA私钥和自签名CA证书:
# 生成CA私钥 openssl genrsa -out ca.key 2048 # 生成CA证书,这里CN是内部CA的名称 openssl req -x509 -new -nodes -key ca.key -days 3650 -out ca.crt \ -subj "/C=CN/ST=Beijing/L=Beijing/O=Example Inc/CN=Example Internal CA"然后把ca.crt分发给每台broker,以及所有需要连接Kafka的客户端机器。接下来为第一台broker生成私钥和证书签名请求(CSR):
# 生成broker私钥 openssl genrsa -out broker1.key 2048 # 生成CSR,注意CN和SAN openssl req -new -key broker1.key -out broker1.csr \ -subj "/C=CN/ST=Beijing/L=Beijing/O=Example Inc/CN=kafka1.example.com" # 创建SAN配置文件,里面包含broker的hostname和IP cat > broker1.ext <<EOF subjectAltName = DNS:kafka1.example.com,IP:192.168.1.11 EOF然后用CA签发出broker证书:
openssl x509 -req -in broker1.csr -CA ca.crt -CAkey ca.key -CAcreateserial \ -out broker1.crt -days 365 -extfile broker1.ext这里有个非常关键的细节:客户端连接broker时,会用broker的hostname或者IP去校验证书里的SAN。如果你在server.properties里配的是kafka1.example.com,但客户端访问的是IP,那SSL握手会直接失败。所以SAN里必须把你所有可能被客户端访问的方式都列进去,包括内网域名、外网域名、IP地址。如果不确定,宁可多写几个SAN。
其他broker重复同样的操作,生成各自独立的私钥和证书。然后需要把私钥和证书导入到Java keystore中,同时把CA证书导入到truststore:
# 创建一个PKCS12 keystore,导入broker私钥和证书 openssl pkcs12 -export -in broker1.crt -inkey broker1.key \ -out broker1.p12 -name kafka1 \ -CAfile ca.crt -caname root -password pass:changeit # 用keytool将PKCS12转为JKS(也可以直接用PKCS12格式,Kafka支持) keytool -importkeystore -deststorepass changeit -destkeypass changeit \ -destkeystore broker1.keystore.jks -srckeystore broker1.p12 \ -srcstoretype PKCS12 -srcstorepass changeit -alias kafka1 # 创建truststore,导入CA证书 keytool -import -trustcacerts -alias caroot -file ca.crt \ -keystore broker1.truststore.jks -storepass changeit -noprompt注意这里我把keystore和truststore都设成了同一个密码changeit,实际环境请换成强密码,而且每个broker的密码最好独立。很多老教程还在用JKS,但Java 9以后官方更推荐PKCS12,你可以直接用.p12作为keystore文件,Kafka完全支持。我个人后来改用了PKCS12,省去JKS格式转换这一层麻烦。
2.3 客户端信任关系怎么建立
客户端需要什么证书?取决于你是否开启了双向认证(mTLS)。我们用的是SASL_SSL + SCRAM,客户端通过SASL的用户名密码认证,所以只需要让客户端信任我们的内部CA即可,不需要给每个客户端签发证书。这样运维负担小很多。
所以交付给开发同学的东西包括:ca.crt文件、broker的hostname和端口、SASL用户名密码、SASL机制类型。开发需要在客户端所在机器的信任库中导入ca.crt,或者在Java应用启动参数里指定truststore。这一步看似简单,实际上在跨部门协作时经常出幺蛾子,我后面会在问题排查章节细讲。
3. Broker端SASL_SSL完整配置
3.1 先配置JAAS文件
Kafka的SASL认证是通过JAAS(Java Authentication and Authorization Service)配置加载的。我们需要创建一个JAAS文件,里面定义了两个section:一个是KafkaServer,用于接收客户端连接;另一个是KafkaClient,用于broker之间互相通信。如果想用Kafka的SCRAM动态用户管理,KafkaServer的配置里不能直接写死密码,而是要用下面的方式:
KafkaServer { org.apache.kafka.common.security.scram.ScramLoginModule required; }; KafkaClient { org.apache.kafka.common.security.scram.ScramLoginModule required username="kafka-admin" password="admin-secret"; };把这段内容保存为/etc/kafka/kafka_jaas.conf。注意KafkaClient里的用户名密码是给broker之间通信用的,这个用户必须存在,而且要有足够的权限。如果集群里只有三节点,通常用同一个超级管理员账号即可。
然后修改Kafka启动脚本,或者用systemd管理时添加JVM参数。如果是传统bin/kafka-server-start.sh,可以修改脚本顶部的KAFKA_OPTS,或者启动前设置环境变量:
export KAFKA_OPTS="-Djava.security.auth.login.config=/etc/kafka/kafka_jaas.conf"在systemd里就需要把Environment配置写上,这个别忘,忘了会报Failed to load login module之类的错误。
3.2 server.properties核心参数
接下来是config/server.properties。我建议为SASL_SSL单独设置一个listener,而不是直接改原来的9092,这样可以平滑过渡。比如继续保留9092作为PLAINTEXT供测试?但生产环境不建议同时开明文listener,除非你确信只是临时迁移。稳妥做法是直接新增一个9093端口:
# 监听器定义 listeners=SASL_SSL://0.0.0.0:9093 advertised.listeners=SASL_SSL://kafka1.example.com:9093 # 安全协议 security.inter.broker.protocol=SASL_SSL sasl.enabled.mechanisms=SCRAM-SHA-256 sasl.mechanism.inter.broker.protocol=SCRAM-SHA-256 # SSL配置 ssl.keystore.location=/etc/kafka/security/broker1.keystore.jks ssl.keystore.password=changeit ssl.key.password=changeit ssl.truststore.location=/etc/kafka/security/broker1.truststore.jks ssl.truststore.password=changeit ssl.client.auth=none这里我特意把ssl.client.auth设成none,因为我们用的是SASL认证,不需要客户端证书。如果你对安全性要求更高,可以改成required做双向TLS,但那样客户端配置复杂度会上升一大截,一般没必要。
还有一个容易被忽略的老参数:listener.name.sasl_ssl.scram-sha-256.sasl.jaas.config。在新版本Kafka中,你可以直接在server.properties里配置JAAS,而不一定用外部文件。不过既然我们已经在JVM参数里指定了JAAS文件,就不需要再重复配置。但如果你看到某些新教程这么写,也不奇怪:
listener.name.sasl_ssl.scram-sha-256.sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required;两种方式选一种,别混着写,否则会报配置冲突。我个人倾向用外部JAAS文件,因为统一管理比较直观。
3.3 Broker间通信是一条单独的认证链路
集群内broker之间也需要认证,否则一个恶意客户端伪装成broker,就能复制整个topic数据。Kafka的security.inter.broker.protocol专门控制这条链路。上面配置里我已设置为SASL_SSL,并且ssl.keystore和ssl.truststore都是每个broker自己的证书和信任库。这里有个小细节:broker之间连接时,目标broker的主机名校验同样会做,所以SAN配置一定要包含broker之间的真实hostname。
如果你有三台机器,分别叫kafka1、kafka2、kafka3,那每台broker的证书里都应该在SAN中包含自身hostname,同时信任同一个CA。这样A连B时,B的证书里必须有kafka2这个SAN,否则握手失败。我见过不少人在单节点测试时一切正常,一扩展到集群就报SSLHandshakeException,基本都是SAN漏了集群内其他broker的hostname。
配置好之后,在重启集群前,最好先在每台机器上用openssl验证一下证书:
openssl s_client -connect kafka2:9093 -CAfile ca.crt如果输出里有Verify return code: 0 (ok),说明证书链没问题。
3.4 滚动重启和初始用户创建
集群配置改完后,需要重启每个broker。强烈建议滚动重启:先重启一台,确认日志里没有报错,再重启下一台。Kafka在滚动重启时会导致分区leader切换,但只要unclean.leader.election没乱开,短暂抖动是可以接受的。日志里如果看到类似:
INFO [SocketServer listenerType=BROKER] ... SSL listener started on port 9093说明SSL监听已经正常。接着创建SCRAM用户。Kafka提供了一个命令行工具:kafka-configs.sh。注意是在3.x版本里直接用它,老版本需要用sasl.scram相关的工具。创建一个超级管理员用户,用于broker间通信:
bin/kafka-configs.sh --bootstrap-server kafka1:9093 \ --alter --add-config 'SCRAM-SHA-256=[password=admin-secret]' \ --entity-type users --entity-name kafka-admin执行这个命令时,因为集群还没完全准备好,可能需要先通过PLAINTEXT或者本机连接。我当时的做法是保留一个临时内网PLAINTEXT端口,把SCRAM用户创建完后,再删掉那个PLAINTEXT listener。当然如果你用外部JAAS文件写死了KafkaServer,其实不需要动态用户,但动态用户的管理方式更好,推荐大家学会用这个工具。
给应用创建用户也同理:
bin/kafka-configs.sh --bootstrap-server kafka1:9093 \ --alter --add-config 'SCRAM-SHA-256=[password=app-secret]' \ --entity-type users --entity-name app-user之后这个用户就可以通过SASL证书认证连上来。
4. 客户端接入配置与实操验证
4.1 命令行工具快速验证认证链路
服务端配好之后,第一件事就是拿命令行自带的工具做冒烟测试。不要直接上Java代码,先用kafka-console-producer.sh和kafka-console-consumer.sh验证链路。命令里需要带上SSL和SASL相关参数:
# 创建topic bin/kafka-topics.sh --bootstrap-server kafka1:9093 \ --create --topic test-topic \ --command-config /tmp/client.properties # 启动生产者 bin/kafka-console-producer.sh \ --bootstrap-server kafka1:9093 \ --topic test-topic \ --producer.config /tmp/client.properties这里的/tmp/client.properties内容如下:
security.protocol=SASL_SSL sasl.mechanism=SCRAM-SHA-256 sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required \ username="app-user" password="app-secret"; ssl.truststore.location=/etc/kafka/security/client.truststore.jks ssl.truststore.password=changeit注意sasl.jaas.config必须在一行内,如果跨行需要加反斜杠。另外,命令行里不要把密码写在进程参数中,容易泄露到history里,最好用配置文件。这里我只是为了演示。能正常发送和消费,说明认证和加密链路已经通了。
4.2 Java客户端属性配置要点
Java应用中,KafkaProducer和KafkaConsumer的配置并不复杂,关键在于Properties要设置完整。下面是一个典型的配置:
Properties props = new Properties(); props.put("bootstrap.servers", "kafka1.example.com:9093,kafka2.example.com:9093,kafka3.example.com:9093"); props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer"); // 安全配置 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=\"app-user\" password=\"app-secret\";"); props.put("ssl.truststore.location", "/etc/kafka/security/client.truststore.jks"); props.put("ssl.truststore.password", "changeit");这里有几个坑:
第一,sasl.jaas.config这个属性在旧版本Kafka客户端里不支持,如果你还在用很久以前的老版本,就需要在JVM参数里指定JAAS文件。但现在主流客户端版本都支持在Properties里直接设置,灵活度更高。
第二,ssl.truststore的路径和密码在所有客户端机器上要保持一致,或者通过配置中心分发。如果应用以容器方式运行,千万别把truststore放在镜像里,应该用挂载或注入的方式,避免证书泄露。
第三,Java客户端默认会校验服务端主机名,这个必须开着。要是你看到Caused by: java.security.cert.CertificateException: No name matching kafka1.example.com found,说明SAN里没有这个域名,别想着关校验,而是回去修证书。
为了便于团队协作,我通常会把安全配置封装成一个公共类,避免每个业务方都去翻文档自己拼。最好把连接串、用户名密码等做成环境变量,不要提交到代码仓库。
4.3 其他语言客户端如何接入
如果你的应用是Python,使用confluent-kafka-python的话,配置项和Java类似,只是参数名略有差异:
conf = { 'bootstrap.servers': 'kafka1.example.com:9093', 'security.protocol': 'SASL_SSL', 'sasl.mechanism': 'SCRAM-SHA-256', 'sasl.username': 'app-user', 'sasl.password': 'app-secret', 'ssl.ca.location': '/etc/kafka/security/ca.crt', }Go语言使用segmentio/kafka-go的话,稍微繁琐一点,需要在Transport配置TLS和SASL。这里不多展开,但核心思路一致:protocol设成SASL_SSL,提供credential,配置好TLS信任链。
4.4 图形化工具怎么配置
很多同学喜欢用可视化工具管理Kafka,比如Kafka Tool(现在叫Offset Explorer)或者Kafka UI。这些工具都支持SASL_SSL配置。以Offset Explorer为例,新建连接时选择Security Protocol为SASL_SSL,然后填写SASL Mechanism为SCRAM-SHA-256,Username/Password填上,还需要在Advanced SSL config里指定truststore路径。这里也是我经常被同事问的地方:明明配了证书怎么连不上?大多数情况是工具的JVM用的默认信任库,没有把我们内部CA导进去,需要在工具启动脚本里加上-Djavax.net.ssl.trustStore=xxx,或者在界面里指定信任库。
所以在交付文档里,我一般会附带一个“客户端连接配置样例”,覆盖Java、命令行、图形工具三种,这样大家在接入时不用自己猜。
5. 常见问题与排查技巧实录
5.1 SSL握手失败,证书链不完整
这是我在配置过程中遇到最多的一类问题。报错通常长这样:
Caused by: javax.net.ssl.SSLHandshakeException: PKIX path building failed: sun.security.provider.certpath.SunCertPathBuilderException: unable to find valid certification path to requested target原因无非是客户端不信任broker的证书,要么ca.crt没导入到truststore,要么导入的不是签发broker证书的那个CA。排查时可以先用keytool -list -keystore client.truststore.jks确认CA别名存在,再用openssl verify -CAfile ca.crt broker1.crt检查证书链。如果是自签名broker证书(而不是CA签发),那就必须把broker证书本身导入客户端truststore,但这样的话每个broker证书都要导一遍,太麻烦。所以我才坚持用内部CA统一签发。另外,时间不同步也可能导致证书校验失败,注意客户端和服务器的系统时间要同步到误差范围内。
5.2 SASL认证失败,机制不匹配
另一个高频问题:
org.apache.kafka.common.errors.SaslAuthenticationException: Authentication failed during authentication due to invalid credentials with SASL mechanism SCRAM-SHA-256如果用户名密码没错,最先怀疑两点:一是服务端server.properties里sasl.enabled.mechanisms没包含SCRAM-SHA-256,或者客户端设置的mechanism和服务端不匹配;二是SCRAM用户并没有真正创建。用kafka-configs.sh --describe --entity-type users --entity-name app-user看看是否返回了SCRAM-SHA-256配置。如果发现用户不存在,就回到3.4节创建用户那一步。
还有种隐蔽情况,broker间通信和外部客户端共用了同一个JAAS里的KafkaClient账号,但那个账号已经被删掉或密码改了,导致broker启动后正常,但互相连接时认证失败。日志里会频繁出现Failed authentication with kafka2之类的提示。所以集群中的所有broker的KafkaClient账号要一直保持一致,严格管理密码变更流程。
5.3 性能影响和延迟问题
启用SASL_SSL后,很多团队会担心性能下降。说实话,加密确实有CPU开销,但影响程度取决于流量大小和CPU能力。我实测过一个中等负载集群,开启SASL_SSL后吞吐量下降大概5-10%,延迟增加也在毫秒级别。对大多数业务来说可以接受。但如果你的消息量极大,比如每秒几十万条,就要注意SSL握手不是大头,因为连接会复用,真正消耗CPU的是持续加解密。
如果发现延迟明显变高,可以从几个方向排查:
- SSL算法套件是否选了很慢的算法,比如用默认配置和用
TLS_ECDHE_RSA_WITH_AES_128_GCM_SHA256性能差很多。 - 是否所有连接都走了一个broker,没有均衡。
- 客户端是否频繁重连,导致SSL握手过多。对比一下
socket.connection.setup.timeout.ms配置是否过小。
针对这个话题,热搜里提到的“kafka消息延迟高”,其实很多时候不是SASL_SSL造成的,而是acks配置、批量大小、分区数设计的问题。但安全改造后,如果延迟异常,建议先对照改造前后的基线数据,再动手排查证书和密钥套件。我习惯在压测时开jmx指标里的ssl相关计数器,观察握手次数和失败率。
5.4 安全加固的进阶建议
SASL_SSL配通只是第一步,生产环境里还有不少需要加固的细节。
- 最小化super.users:不要把所有应用账号都设为超级用户,否则ACL形同虚设。通过
--add-config给每个应用只授权需要的topic。 - 定期轮换证书:内部CA签发的证书设为一年有效期,在快到期前用openssl重新签发批次证书,然后滚动更新broker。
- 开启审计日志:Kafka的authorizer日志会记录用户的访问请求,这些对追查安全问题很有用。
- 管理好信任库:客户端机器上的truststore密码不要明文写在配置文件里,可以用密钥管理服务或者环境变量注入。
还有一点非常重要的“后门”:千万不要把PLAINTEXT listener长期留着。很多事故都是因为某人图方便,临时开放了9092明文端口,忘记关闭,结果安全策略形同虚设。建议在server.properties里只保留SASL_SSL这一个listener,不要贪多。
6. 个人经验总结
最开始我其实是抱着“配置一下不复杂”的心态去做的,结果前前后后折腾了两天,主要在证书SAN和SCAM用户上栽了跟头。后来我把整个流程固定成了一个标准操作手册,每次新加broker或者新接入应用都按这个流程走,基本很顺畅。
有一个小技巧分享给你:调试阶段,可以用openssl s_client -connect broker:9093 -tlsextdebug代替实际客户端,能快速看到服务端证书和信任链情况。另外,Kafka客户端的日志级别调到DEBUG后,能看到非常详细的认证握手过程,定位问题时不要怕刷屏,安全改造本来就值得花时间看清每一步。
最后建议所有正在实施Kafka安全的团队,先把证书和用户管理的基础设施想好,不要为了赶进度省略信任库和SAN校验,后面欠下的债都要还的。希望这篇SASL_SSL配置过程复盘能让你少走几条弯路,把安全改造做得更利索。