如何使用纯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
相关产品推荐
相关产品推荐

