Quarkus中如何从AWS Secrets Manager加载Kafka配置密钥?
问题描述
我们的项目当前使用Smallrye Kafka(通过quarkus-smallrye-reactive-messaging-kafka)处理各类消息任务,其配置在application.yaml中如下:
kafka.sasl.jaas.config: org.apache.kafka.common.security.plain.PlainLoginModule required username="${kafka.user}" password="${kafka.pswd}";
安全团队指出明文存储用户名/密码存在风险,要求将其迁移至AWS Secrets Manager。我们目前已使用Secrets Manager存储下游Web服务的凭据,通过注入SecretsManagerClient(来自quarkus-amazon-secretsmanager)的服务进行访问。现咨询:
- 能否在应用配置中直接从Secrets Manager加载
kafka.user和kafka.pswd? - 是否可修改Smallrye-Kafka的默认实例化行为,通过自定义服务获取用户名和密码?
解决方案
针对你的问题,分两种场景给出具体实现方式:
1. 直接在配置中从AWS Secrets Manager加载凭据
Quarkus的AWS Secrets Manager扩展支持直接在配置文件中引用存储的密钥,无需额外代码,这是最简洁的实现方式。
操作步骤:
- 确保
quarkus-amazon-secretsmanager扩展版本不低于Quarkus 2.10.x(低版本可能不支持该特性) - 在
application.yaml中使用${sm://<secret-name>/<key>}格式直接引用Secrets Manager中的值:
注:kafka.user: ${sm://your-kafka-secret/username} kafka.pswd: ${sm://your-kafka-secret/password} kafka.sasl.jaas.config: org.apache.kafka.common.security.plain.PlainLoginModule required username="${kafka.user}" password="${kafka.pswd}";your-kafka-secret是你在AWS Secrets Manager中存储Kafka凭据的密钥名称,username和password是密钥JSON结构中的对应字段;如果密钥是纯文本格式,直接使用${sm://your-kafka-secret}即可,同时需要调整JAAS配置的拼接逻辑适配纯文本内容。
2. 自定义服务控制Kafka客户端实例化
如果需要更灵活的逻辑(比如动态刷新凭据、自定义密钥解析规则),可以通过Quarkus提供的KafkaClientConfigCustomizer接口实现自定义配置注入:
实现步骤:
- 创建自定义配置定制器类,注入
SecretsManagerClient获取凭据:import io.quarkus.kafka.client.runtime.KafkaClientConfigCustomizer; import software.amazon.awssdk.services.secretsmanager.SecretsManagerClient; import software.amazon.awssdk.services.secretsmanager.model.GetSecretValueRequest; import software.amazon.awssdk.services.secretsmanager.model.GetSecretValueResponse; import jakarta.enterprise.context.ApplicationScoped; import java.util.Map; import com.fasterxml.jackson.databind.ObjectMapper; @ApplicationScoped public class KafkaSecretConfigCustomizer implements KafkaClientConfigCustomizer { private final SecretsManagerClient secretsManagerClient; private final ObjectMapper objectMapper; // 构造函数注入依赖 public KafkaSecretConfigCustomizer(SecretsManagerClient secretsManagerClient, ObjectMapper objectMapper) { this.secretsManagerClient = secretsManagerClient; this.objectMapper = objectMapper; } @Override public void customize(String configName, Map<String, Object> config) { // 仅对Kafka相关配置生效 if (configName.startsWith("kafka.")) { // 从Secrets Manager获取密钥 GetSecretValueRequest request = GetSecretValueRequest.builder() .secretId("your-kafka-secret") .build(); GetSecretValueResponse response = secretsManagerClient.getSecretValue(request); // 解析JSON格式的密钥内容 Map<String, String> secret = objectMapper.readValue(response.secretString(), Map.class); // 构造JAAS配置并注入 String jaasConfig = String.format( "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"%s\" password=\"%s\";", secret.get("username"), secret.get("password") ); config.put("sasl.jaas.config", jaasConfig); } } } - 移除
application.yaml中原有的kafka.user、kafka.pswd及硬编码的JAAS配置项,确保Kafka客户端会使用自定义器提供的配置。
额外说明:
- 第一种方式适合大多数常规场景,无需编写代码;第二种方式适合需要自定义逻辑的场景,比如结合Quarkus定时任务实现凭据自动刷新。
- 确保运行应用的IAM角色拥有
secretsmanager:GetSecretValue权限,能够访问目标AWS Secrets Manager密钥。
内容的提问来源于stack exchange,提问作者Amoeba
相关产品推荐
相关产品推荐

