You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Kafka Broker与客户端侧加密安全配置密码?

Kafka 服务端与客户端密码加密存储实现示例

一、Kafka服务端(server.properties)实现方案

Kafka支持通过自定义密码解密器解析加密后的配置项,核心是实现org.apache.kafka.common.config.Configurable接口,启动时自动解密加密密码值。

1. 实现自定义解密类

以AES对称加密为例,编写解密逻辑:

import org.apache.kafka.common.config.Configurable;
import org.apache.kafka.common.config.ConfigException;
import javax.crypto.Cipher;
import javax.crypto.spec.SecretKeySpec;
import java.util.Base64;
import java.util.Map;

public class AESPasswordDecryptor implements Configurable {
    private static final String ALGORITHM = "AES";
    private String secretKey; // 加密密钥需提前约定并妥善存储

    @Override
    public void configure(Map<String, ?> configs) {
        // 从Kafka配置读取解密密钥(也可通过环境变量/外部密钥管理系统获取)
        this.secretKey = (String) configs.get("password.decryptor.secret");
        if (secretKey == null || secretKey.length() != 16) { // AES-128要求密钥16位
            throw new ConfigException("Missing or invalid password decryptor secret key");
        }
    }

    // Kafka自动调用此方法处理配置值
    public String decrypt(String encryptedValue) {
        try {
            SecretKeySpec keySpec = new SecretKeySpec(secretKey.getBytes(), ALGORITHM);
            Cipher cipher = Cipher.getInstance(ALGORITHM);
            cipher.init(Cipher.DECRYPT_MODE, keySpec);
            byte[] decodedBytes = Base64.getDecoder().decode(encryptedValue);
            return new String(cipher.doFinal(decodedBytes));
        } catch (Exception e) {
            throw new RuntimeException("Failed to decrypt password", e);
        }
    }
}

2. 打包部署解密类

  • 将上述代码打包为jar文件(如kafka-password-decryptor.jar)
  • 把jar放入Kafka安装目录的libs文件夹

3. 修改server.properties配置

替换明文密码为加密字符串,并指定解密器:

# 配置解密器密钥(建议通过环境变量传入,避免写在配置文件)
password.decryptor.secret=your_aes_secret_key_16
ssl.truststore.password=${decrypt:encrypted_truststore_pass_base64}
ssl.keystore.password=${decrypt:encrypted_keystore_pass_base64}
ssl.key.password=${decrypt:encrypted_key_pass_base64}
listener.name.sasl.ssl.scram-sha-256.sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required \
    username="admin" \
    password="${decrypt:encrypted_scram_pass_base64}";

# 指定解密器类
config.providers=decrypt
config.provider.decrypt.class=AESPasswordDecryptor

注:加密密码时,用同一份AES密钥对明文加密,生成Base64编码字符串替换到配置中。


二、Spring Cloud Stream客户端(application.properties)实现方案

Spring生态可通过自定义EnvironmentPostProcessor解密配置项,以下是具体实现:

1. 实现解密处理器

import org.springframework.boot.SpringApplication;
import org.springframework.boot.env.EnvironmentPostProcessor;
import org.springframework.core.env.ConfigurableEnvironment;
import org.springframework.core.env.MapPropertySource;
import javax.crypto.Cipher;
import javax.crypto.spec.SecretKeySpec;
import java.util.Base64;
import java.util.HashMap;
import java.util.Map;

public class KafkaPasswordDecryptor implements EnvironmentPostProcessor {
    private static final String ALGORITHM = "AES";
    private static final String SECRET_KEY = "your_aes_secret_key_16"; // 可与服务端密钥一致或单独约定
    private static final String CIPHER_PREFIX = "{cipher}";

    @Override
    public void postProcessEnvironment(ConfigurableEnvironment environment, SpringApplication application) {
        Map<String, Object> decryptedProps = new HashMap<>();
        
        // 定义需要解密的配置项
        String[] targetProps = {
            "spring.cloud.stream.kafka.binder.jaas.options.password",
            "spring.cloud.stream.kafka.binder.configuration.ssl.truststore.password"
        };

        for (String propKey : targetProps) {
            String encryptedValue = environment.getProperty(propKey);
            if (encryptedValue != null && encryptedValue.startsWith(CIPHER_PREFIX)) {
                String encryptedPart = encryptedValue.substring(CIPHER_PREFIX.length());
                decryptedProps.put(propKey, decrypt(encryptedPart));
            }
        }

        if (!decryptedProps.isEmpty()) {
            environment.getPropertySources().addFirst(new MapPropertySource("decrypted-kafka-props", decryptedProps));
        }
    }

    private String decrypt(String encryptedValue) {
        try {
            SecretKeySpec keySpec = new SecretKeySpec(SECRET_KEY.getBytes(), ALGORITHM);
            Cipher cipher = Cipher.getInstance(ALGORITHM);
            cipher.init(Cipher.DECRYPT_MODE, keySpec);
            byte[] decodedBytes = Base64.getDecoder().decode(encryptedValue);
            return new String(cipher.doFinal(decodedBytes));
        } catch (Exception e) {
            throw new RuntimeException("Failed to decrypt Kafka password", e);
        }
    }
}

2. 注册处理器

在src/main/resources/META-INF/spring.factories中添加配置,让Spring自动加载处理器:

org.springframework.boot.env.EnvironmentPostProcessor=com.your.package.KafkaPasswordDecryptor

3. 修改application.properties配置

将明文密码替换为带{cipher}前缀的加密字符串:

spring.cloud.stream.kafka.binder.jaas.options.password={cipher}encrypted_scram_pass_base64
spring.cloud.stream.kafka.binder.configuration.ssl.truststore.password={cipher}encrypted_truststore_pass_base64

额外注意事项

  • 加密密钥禁止存储在配置文件中,建议通过环境变量、密钥管理系统(如HashiCorp Vault)或服务器本地安全文件读取。
  • 生产环境推荐使用非对称加密(RSA),私钥仅在解密端持有,公钥用于加密密码,提升安全性。
  • 替换生产配置前,务必先测试加密/解密流程是否正常运行。

内容的提问来源于stack exchange,提问作者Sudarsh1

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.19 13:52:05