如何为Kafka服务端设置访问账号密码?Spring Boot项目该如何适配配置
1 服务端Kafka身份验证配置
Kafka默认未开启身份验证,可通过SASL/PLAIN机制快速实现账号密码访问限制,操作步骤如下:
- 第一步:创建服务端JAAS认证配置文件
kafka_server_jaas.conf,内容如下,可自定义多组账号密码:
KafkaServer { org.apache.kafka.common.security.plain.PlainLoginModule required username="admin" password="Admin@123456" user_admin="Admin@123456" user_client="Client@654321"; };
说明:username和password是broker内部通信使用的账号,user_xxx开头的配置是开放给客户端使用的账号,等号右侧为对应密码。
- 第二步:修改Kafka配置文件
config/server.properties,新增/修改以下配置项:
# 开启SASL监听 listeners=SASL_PLAINTEXT://0.0.0.0:9092 # 对外暴露的访问地址,替换为你的服务器实际IP advertised.listeners=SASL_PLAINTEXT://你的服务器公网IP:9092 # broker间通信协议 security.inter.broker.protocol=SASL_PLAINTEXT # 启用PLAIN认证机制 sasl.enabled.mechanisms=PLAIN sasl.mechanism.inter.broker.protocol=PLAIN
- 第三步:重启Kafka服务,启动时指定加载JAAS配置:
# 先设置环境变量,替换为你实际的配置文件路径 export KAFKA_OPTS="-Djava.security.auth.login.config=/opt/kafka/config/kafka_server_jaas.conf" # 再执行启动命令 bin/kafka-server-start.sh -daemon config/server.properties
2 Spring Boot 配置修改
你现有配置类需要在生产者、消费者的配置中新增SASL认证相关参数,修改后的完整配置如下:
@EnableKafka @Configuration @SuppressWarnings("SpringFacetCodeInspection") public class KafkaConfig { @Bean ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory(ConsumerFactory<String, String> consumerFactory) { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); return factory; } @Bean public ConsumerFactory<String, String> consumerFactory() { final Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "你的服务器IP:9092"); props.put(ConsumerConfig.GROUP_ID_CONFIG, UUID.randomUUID().toString()); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest"); // 新增认证配置 props.put("security.protocol", "SASL_PLAINTEXT"); props.put("sasl.mechanism", "PLAIN"); // 替换为你在服务端配置的客户端账号密码 props.put("sasl.jaas.config", "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"client\" password=\"Client@654321\";"); return new DefaultKafkaConsumerFactory<>(props); } @Bean public ProducerFactory<String, String> producerFactory() { final Map<String, Object> props = new HashMap<>(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "你的服务器IP:9092"); props.put(ProducerConfig.LINGER_MS_CONFIG, 10); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); // 新增认证配置,和消费者保持一致即可 props.put("security.protocol", "SASL_PLAINTEXT"); props.put("sasl.mechanism", "PLAIN"); props.put("sasl.jaas.config", "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"client\" password=\"Client@654321\";"); return new DefaultKafkaProducerFactory<>(props); } @Bean public KafkaTemplate<String, String> kafkaTemplate(ProducerFactory<String, String> producerFactory) { return new KafkaTemplate<>(producerFactory); } }
注意事项
生产环境建议使用SASL_SSL协议配合SSL证书加密传输,避免账号密码在网络中明文传输,提升安全性。
内容的提问来源于stack exchange,提问作者My Name
相关产品推荐
相关产品推荐

