Kafka配置SASL_SSL认证传输加密

发布时间:2026/7/28 0:29:33
Kafka配置SASL_SSL认证传输加密 Kafka配置SASL_SSL认证传输加密在大数据与消息队列领域Apache Kafka 凭借其高吞吐、低延迟、高可扩展性等特性成为实时数据流处理的核心组件。然而随着数据安全与合规性要求日益严格Kafka 需要同时满足传输加密SSL/TLS与身份认证SASL的需求。本文将循序渐进地讲解如何在 Kafka 中配置 SASL_SSL确保客户端与 Broker 之间的通信既加密又经过认证。## 一、基础概念### 1.1 为什么需要 SASL_SSL默认情况下Kafka 使用明文传输这意味着- 任何能监听网络的人都可以看到消息内容。- 任何人都可以冒充合法客户端连接 Broker。SASL_SSL 结合了两层安全-SSL/TLS为数据传输提供加密防止窃听与篡改。-SASLSimple Authentication and Security Layer提供身份认证机制如 PLAIN、SCRAM、GSSAPI 等。因此启用 SASL_SSL 后Kafka 集群变为“加密信道 认证入口”只有持有合法证书和凭证的客户端才能通信。### 1.2 核心组件-证书由 CA 签发用于 SSL 握手验证双方身份。-密钥库Keystore存储 Broker 或客户端的私钥与证书。-信任库Truststore存储受信任的 CA 证书。-JAAS 配置文件定义 SASL 认证的具体实现如用户名密码。## 二、环境准备### 2.1 准备证书我们使用 Java 自带的keytool生成自签名证书生产环境应使用 CA 签发。bash# 生成 Broker 的密钥库包含私钥与自签名证书keytool -genkey -alias kafka-broker -keyalg RSA -keystore broker.keystore.jks -dname CNlocalhost, OUdev, Oexample, LBeijing, SBeijing, CCN -storepass changeme -keypass changeme# 导出证书keytool -export -alias kafka-broker -keystore broker.keystore.jks -file broker.cer -storepass changeme# 创建客户端的信任库信任 Broker 的证书keytool -import -alias kafka-broker -keystore client.truststore.jks -file broker.cer -storepass changeme -noprompt这里我们只创建了 Broker 的证书客户端信任它。如果需要双向认证mTLS还需为客户生成证书但 SASL_SSL 通常只要求 SSL 单向认证认证由 SASL 层完成。### 2.2 配置 JAAS 文件SASL 使用 PLAIN 机制明文密码生产环境建议使用 SCRAM。创建kafka_server_jaas.confplaintextKafkaServer { org.apache.kafka.common.security.plain.PlainLoginModule required usernameadmin passwordadmin-secret user_adminadmin-secret user_alicealice-secret;};-username/passwordBroker 间通信的凭证。-user_xxx定义客户端用户的密码。## 三、Kafka Broker 配置编辑server.properties加入以下配置properties# 开启 SSLlistenersSASL_SSL://localhost:9093advertised.listenersSASL_SSL://localhost:9093# SSL 配置ssl.keystore.location/path/to/broker.keystore.jksssl.keystore.passwordchangemessl.key.passwordchangemessl.truststore.location/path/to/client.truststore.jksssl.truststore.passwordchangeme# SASL 配置sasl.enabled.mechanismsPLAINsasl.mechanism.inter.broker.protocolPLAINsecurity.inter.broker.protocolSASL_SSL# JAAS 文件listener.name.sasl_ssl.plain.sasl.jaas.configorg.apache.kafka.common.security.plain.PlainLoginModule required \ usernameadmin \ passwordadmin-secret \ user_adminadmin-secret \ user_alicealice-secret;注意listener.name.sasl_ssl.plain.sasl.jaas.config的格式与 JAAS 文件相同但直接写在配置中更方便。启动 Kafka Broker 时需指定 JAAS 文件如果使用文件方式bashexport KAFKA_OPTS-Djava.security.auth.login.config/path/to/kafka_server_jaas.confbin/kafka-server-start.sh config/server.properties## 四、Java 客户端代码示例### 4.1 生产者示例javaimport org.apache.kafka.clients.producer.*;import java.util.Properties;public class SecureProducer { public static void main(String[] args) { Properties props new Properties(); // 1. 必须指向 SASL_SSL 端口 props.put(bootstrap.servers, localhost:9093); // 2. SSL 配置信任 Broker 证书 props.put(ssl.truststore.location, /path/to/client.truststore.jks); props.put(ssl.truststore.password, changeme); // 3. SASL 配置 props.put(sasl.mechanism, PLAIN); props.put(security.protocol, SASL_SSL); // 4. JAAS 配置客户端凭证 props.put(sasl.jaas.config, org.apache.kafka.common.security.plain.PlainLoginModule required username\alice\ password\alice-secret\;); // 5. 序列化与主题 props.put(key.serializer, org.apache.kafka.common.serialization.StringSerializer); props.put(value.serializer, org.apache.kafka.common.serialization.StringSerializer); KafkaProducerString, String producer new KafkaProducer(props); // 发送消息 ProducerRecordString, String record new ProducerRecord(secure-topic, Hello SASL_SSL!); producer.send(record, (metadata, exception) - { if (exception null) { System.out.println(消息发送成功分区 metadata.partition() 偏移量 metadata.offset()); } else { exception.printStackTrace(); } }); producer.close(); }}### 4.2 消费者示例javaimport org.apache.kafka.clients.consumer.*;import java.time.Duration;import java.util.Collections;import java.util.Properties;public class SecureConsumer { public static void main(String[] args) { Properties props new Properties(); // 1. 连接 SASL_SSL Broker props.put(bootstrap.servers, localhost:9093); // 2. SSL 信任库 props.put(ssl.truststore.location, /path/to/client.truststore.jks); props.put(ssl.truststore.password, changeme); // 3. SASL 认证 props.put(sasl.mechanism, PLAIN); props.put(security.protocol, SASL_SSL); props.put(sasl.jaas.config, org.apache.kafka.common.security.plain.PlainLoginModule required username\alice\ password\alice-secret\;); // 4. 反序列化与消费者组 props.put(key.deserializer, org.apache.kafka.common.serialization.StringDeserializer); props.put(value.deserializer, org.apache.kafka.common.serialization.StringDeserializer); props.put(group.id, secure-group); props.put(auto.offset.reset, earliest); KafkaConsumerString, String consumer new KafkaConsumer(props); consumer.subscribe(Collections.singletonList(secure-topic)); try { while (true) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecordString, String record : records) { System.out.printf(收到消息key%s, value%s, 分区%d, 偏移量%d%n, record.key(), record.value(), record.partition(), record.offset()); } } } finally { consumer.close(); } }}## 五、高级用法SCRAM 认证PLAIN 机制将密码以明文传输不安全。建议使用SCRAMSalted Challenge Response Authentication Mechanism它通过哈希与盐值保护密码。### 5.1 配置 SCRAM首先在 Kafka 中创建 SCRAM 用户bashbin/kafka-configs.sh --zookeeper localhost:2181 \ --alter --add-config SCRAM-SHA-256[iterations8192,passwordalice-secret] \ --entity-type users --entity-name aliceBroker 配置改为propertiessasl.enabled.mechanismsSCRAM-SHA-256sasl.mechanism.inter.broker.protocolSCRAM-SHA-256listener.name.sasl_ssl.scram-sha-256.sasl.jaas.configorg.apache.kafka.common.security.scram.ScramLoginModule required;客户端 JAAS 对应改为javaprops.put(sasl.mechanism, SCRAM-SHA-256);props.put(sasl.jaas.config, org.apache.kafka.common.security.scram.ScramLoginModule required username\alice\ password\alice-secret\;);## 六、总结本文从基础概念出发逐步讲解了 Kafka SASL_SSL 配置的核心步骤1.生成证书为 Broker 创建密钥库将证书导入客户端信任库。2.配置 Broker修改server.properties同时启用 SSL 与 SASL。3.编写客户端Java 生产者、消费者需指定 SSL 信任库与 SASL 凭证。4.安全增强推荐使用 SCRAM 替代 PLAIN 机制防止密码泄露。SASL_SSL 是保护 Kafka 生产环境的基石它确保了数据在传输过程中的机密性、完整性并提供了强身份认证。在生产部署中还应当注意- 使用正式 CA 签发证书而非自签名。- 定期轮换证书与密码。- 结合网络隔离如防火墙规则形成纵深防御。通过本文的指导你应该能够自主搭建一个安全的 Kafka 集群并编写对应的加密认证客户端。