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

如何使用纯Java Kafka客户端连接MSK集群(无需SpringBoot)

纯Java Kafka客户端连接Amazon MSK集群(非SpringBoot实现)

前置要求

  • EC2实例已绑定具备MSK访问权限的IAM角色(例如允许kafka:DescribeCluster、kafka:GetBootstrapBrokers及消息生产/消费权限)
  • EC2与MSK集群处于同一VPC,安全组规则允许EC2访问MSK的对应端口(明文通信用9092,IAM加密通信用9094)
  • 项目中引入Kafka客户端及MSK IAM认证依赖(Maven示例):
<dependencies>
    <dependency>
        <groupId>org.apache.kafka</groupId>
        <artifactId>kafka-clients</artifactId>
        <version>3.6.1</version>
    </dependency>
    <dependency>
        <groupId>software.amazon.msk</groupId>
        <artifactId>aws-msk-iam-auth</artifactId>
        <version>1.1.5</version>
    </dependency>
</dependencies>

生产者实现(纯Java)

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.serialization.StringSerializer;

import java.util.Properties;

public class MSKPlainJavaProducer {
    public static void main(String[] args) {
        // MSK集群引导Broker地址,从MSK控制台获取
        String bootstrapServers = "b-1.xxx.xx.kafka.us-east-1.amazonaws.com:9094,b-2.xxx.xx.kafka.us-east-1.amazonaws.com:9094";
        String topicName = "test-topic";

        Properties props = new Properties();
        // 基础配置
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        props.put(ProducerConfig.ACKS_CONFIG, "all");

        // IAM认证配置
        props.put("security.protocol", "SASL_SSL");
        props.put("sasl.mechanism", "AWS_MSK_IAM");
        props.put("sasl.jaas.config", "software.amazon.msk.auth.iam.IAMLoginModule required;");
        props.put("sasl.client.callback.handler.class", "software.amazon.msk.auth.iam.IAMClientCallbackHandler");

        try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
            for (int i = 0; i < 10; i++) {
                String message = "MSK test message - " + i;
                ProducerRecord<String, String> record = new ProducerRecord<>(topicName, message);
                producer.send(record, (metadata, exception) -> {
                    if (exception == null) {
                        System.out.printf("Message sent to partition %d, offset %d%n", metadata.partition(), metadata.offset());
                    } else {
                        exception.printStackTrace();
                    }
                });
            }
            producer.flush();
            System.out.println("All messages sent successfully");
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

消费者实现(纯Java)

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.serialization.StringDeserializer;

import java.time.Duration;
import java.util.Collections;
import java.util.Properties;

public class MSKPlainJavaConsumer {
    public static void main(String[] args) {
        String bootstrapServers = "b-1.xxx.xx.kafka.us-east-1.amazonaws.com:9094,b-2.xxx.xx.kafka.us-east-1.amazonaws.com:9094";
        String topicName = "test-topic";
        String groupId = "test-consumer-group";

        Properties props = new Properties();
        // 基础配置
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");

        // IAM认证配置
        props.put("security.protocol", "SASL_SSL");
        props.put("sasl.mechanism", "AWS_MSK_IAM");
        props.put("sasl.jaas.config", "software.amazon.msk.auth.iam.IAMLoginModule required;");
        props.put("sasl.client.callback.handler.class", "software.amazon.msk.auth.iam.IAMClientCallbackHandler");

        try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) {
            consumer.subscribe(Collections.singletonList(topicName));
            System.out.println("Subscribed to topic: " + topicName);

            while (true) {
                ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
                for (ConsumerRecord<String, String> record : records) {
                    System.out.printf("Received message: key=%s, value=%s, partition=%d, offset=%d%n",
                            record.key(), record.value(), record.partition(), record.offset());
                }
            }
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

关键配置说明

  • bootstrap.servers:替换为你的MSK集群实际引导Broker地址,可从AWS控制台的MSK集群详情页获取
  • 若无需IAM认证(仅限VPC内部明文访问场景),移除sasl.*及security.protocol配置,将security.protocol设为PLAINTEXT,端口改用9092
  • EC2实例的IAM角色会自动被AWSDefaultCredentialsProvider识别,无需手动配置密钥

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 19:38:09