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

基于IAM认证连接AWS MSK的Spring Boot生产消费应用开发求助

Spring Boot 集成 AWS MSK 并通过 IAM 认证实现消息生产消费

1. 依赖配置

在pom.xml中添加必要依赖:

<!-- Spring Kafka 核心依赖 -->
<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka</artifactId>
</dependency>
<!-- AWS MSK IAM 认证库 -->
<dependency>
    <groupId>software.amazon.msk</groupId>
    <artifactId>aws-msk-iam-auth</artifactId>
    <version>1.1.5</version>
</dependency>
<!-- AWS SDK 核心依赖(用于身份认证) -->
<dependency>
    <groupId>software.amazon.awssdk</groupId>
    <artifactId>sts</artifactId>
    <version>2.20.100</version>
</dependency>

2. IAM 权限配置

核心权限要求

确保你的IAM实体(用户/角色)拥有以下MSK相关权限:

  • kafka:DescribeCluster
  • kafka:WriteData(生产者需要)
  • kafka:ReadData(消费者需要)
  • 若需自动创建主题,添加kafka:CreateTopic

EKS 环境适配(IRSA 方式)

因为你已在EKS部署应用,推荐使用IAM Roles for Service Accounts (IRSA) 实现权限绑定:

  1. 创建IAM角色,关联包含上述权限的策略
  2. 为EKS中的Service Account绑定该IAM角色(通过eksctl或AWS控制台配置)
  3. 在Deployment的yaml中指定serviceAccountName为绑定好的SA

3. 生产者实现

配置文件(application.yml)

spring:
  kafka:
    bootstrap-servers: <MSK_BROKER_SASL_ENDPOINTS> # 格式:broker1:9098,broker2:9098
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.apache.kafka.common.serialization.StringSerializer
      properties:
        security.protocol: SASL_SSL
        sasl.mechanism: AWS_MSK_IAM
        sasl.jaas.config: software.amazon.msk.auth.iam.IAMLoginModule required;
        sasl.client.callback.handler.class: software.amazon.msk.auth.iam.IAMClientCallbackHandler

生产者代码

import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Component;

@Component
public class MskEventProducer {
    private final KafkaTemplate<String, String> kafkaTemplate;
    private static final String TOPIC = "test-topic";

    public MskEventProducer(KafkaTemplate<String, String> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }

    public void sendEvent(String message) {
        kafkaTemplate.send(TOPIC, message);
    }
}

4. 消费者实现

配置文件(application.yml)

spring:
  kafka:
    consumer:
      group-id: msk-consumer-group
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      properties:
        security.protocol: SASL_SSL
        sasl.mechanism: AWS_MSK_IAM
        sasl.jaas.config: software.amazon.msk.auth.iam.IAMLoginModule required;
        sasl.client.callback.handler.class: software.amazon.msk.auth.iam.IAMClientCallbackHandler

消费者代码

import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;

@Component
public class MskEventConsumer {

    @KafkaListener(topics = "test-topic", groupId = "msk-consumer-group")
    public void consumeEvent(String message) {
        // 处理消息逻辑
        System.out.println("Received message from MSK: " + message);
    }
}

5. EKS 部署注意事项

  • 确保MSK集群与EKS集群处于同一VPC,或已配置VPC对等连接
  • MSK的安全组需开放9098端口给EKS节点的安全组
  • 若使用IRSA,需确保OIDC提供商已配置到EKS集群

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 08:30:37