【问题标题】:AWS MSK configuration issue with SpringSpring 的 AWS MSK 配置问题
【发布时间】:2023-01-03 20:34:51
【问题描述】:

我们最近从自我管理的 Kafka 实例迁移到完全托管的 AWS MSK 集群。我们仅启用了基于 IAM 的角色身份验证以从本地系统连接到 MSK 集群。

当我远程登录到集群的公共 url 时,我得到了成功的响应,但是当我尝试启动我的 java 应用程序时,它由于不同的错误而失败了。下面是我的 KafkaConfiguration

错误 :

Invalid login module control flag 'com.amazonaws.auth.AWSStaticCredentialsProvider' in JAAS config
@Configuration
public class KafkaConfiguration {

    @Value("${aws.kafka.bootstrap-servers}")
    private String bootstrapServers;

    @Value("${aws.kafka.accessKey}")
    private String accessKey;

    @Value("${aws.kafka.secret}")
    private String secret;

    @Bean
    public KafkaAdmin kafkaAdmin() {
        AWSCredentials awsCredentials = new BasicAWSCredentials(accessKey, secret);
        Map<String, Object> configs = new HashMap<>();
        configs.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        configs.put(AdminClientConfig.SECURITY_PROTOCOL_CONFIG, "SASL_SSL");
        configs.put(SaslConfigs.SASL_MECHANISM, "AWS_MSK_IAM");
        configs.put(SaslConfigs.SASL_JAAS_CONFIG, "com.amazonaws.auth.AWSCredentialsProvider com.amazonaws.auth.AWSStaticCredentialsProvider(" + awsCredentials + ")");
        return new KafkaAdmin(configs);
    }

    @Bean
    public ProducerFactory<String, String> producerFactory() {
        AWSCredentials awsCredentials = new BasicAWSCredentials(accessKey, secret);

        Map<String, Object> configProps = new HashMap<>();
        configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        configProps.put("security.protocol", "SASL_SSL");
        configProps.put(SaslConfigs.SASL_MECHANISM, "AWS_MSK_IAM");
        configProps.put(SaslConfigs.SASL_JAAS_CONFIG, "com.amazonaws.auth.AWSCredentialsProvider com.amazonaws.auth.AWSStaticCredentialsProvider(" + awsCredentials + ")");
        return new DefaultKafkaProducerFactory<>(configProps);
    }

    @Bean
    public KafkaTemplate<String, String> kafkaTemplate() {
        return new KafkaTemplate<>(producerFactory());
    }
}

消费者配置:

@EnableKafka
@Configuration
public class KafkaConsumerConfig {

    @Value("${aws.kafka.bootstrap-servers}")
    private String bootstrapServers;

    @Value("${aws.kafka.accessKey}")
    private String accessKey;

    @Value("${aws.kafka.secret}")
    private String secret;

    public ConsumerFactory<String, String> consumerFactory() {
        AWSCredentials awsCredentials = new BasicAWSCredentials(accessKey, secret);

        Map<String, Object> configProps = new HashMap<>();
        configProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        configProps.put("security.protocol", "SASL_SSL");
        configProps.put(SaslConfigs.SASL_MECHANISM, "AWS_MSK_IAM");
        configProps.put(SaslConfigs.SASL_JAAS_CONFIG, "com.amazonaws.auth.AWSCredentialsProvider com.amazonaws.auth.AWSStaticCredentialsProvider(" + awsCredentials + ")");
        configProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        configProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        configProps.put(ConsumerConfig.GROUP_ID_CONFIG, "iTopLight");
        return new DefaultKafkaConsumerFactory<>(configProps);
    }

    @Bean
    public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> rawKafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        return factory;
    }
}

【问题讨论】:

    标签: spring-kafka aws-msk


    【解决方案1】:

    将 MSK 与 IAM 身份验证连接的选项不止一种。

    首先你需要在你的项目中使用这个库。

    <dependency>
        <groupId>software.amazon.msk</groupId>
        <artifactId>aws-msk-iam-auth</artifactId>
        <version>1.0.0</version>        
    </dependency>
    

    然后,您需要提供 AWS 访问凭证提供商。您可以使用环境变量或使用系统属性的第一个选项。

    系统属性解决方案看起来像。

    @EnableKafka
    @Configuration
    public class KafkaConsumerConfig {
    
        @Value("${aws.kafka.bootstrap-servers}")
        private String bootstrapServers;
    
        @Value("${aws.kafka.accessKey}")
        private String accessKey;
    
        @Value("${aws.kafka.secret}")
        private String secret;
    
        public ConsumerFactory<String, String> consumerFactory() {
            System.setProperty("aws.accessKeyId", accessKey);
            System.setProperty("aws.secretKey", secret);
    
            Map<String, Object> configProps = new HashMap<>();
            configProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
            configProps.put("security.protocol", "SASL_SSL");
            configProps.put(SaslConfigs.SASL_MECHANISM, "AWS_MSK_IAM");
            configProps.put(SaslConfigs.SASL_JAAS_CONFIG, "software.amazon.msk.auth.iam.IAMLoginModule required");
            configProps.put("sasl.client.callback.handler.class", "software.amazon.msk.auth.iam.IAMClientCallbackHandler");
            configProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
            configProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
            configProps.put(ConsumerConfig.GROUP_ID_CONFIG, "iTopLight");
            return new DefaultKafkaConsumerFactory<>(configProps);
        }
    
        @Bean
        public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> rawKafkaListenerContainerFactory() {
            ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
            factory.setConsumerFactory(consumerFactory());
            return factory;
        }
    }
    

    您可以查看aws-msiam-auth project for providers。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2022-08-04
      • 1970-01-01
      • 1970-01-01
      • 2019-10-01
      • 1970-01-01
      • 1970-01-01
      • 2011-07-02
      • 2017-06-14
      相关资源
      最近更新 更多