如何在Java Kafka生产者中实现--producer.config配置文件加载
解决Kafka Java生产者加载SSL配置文件无效的问题
你在命令行里通过--producer.config指定SSL配置文件的方式是命令行工具的专属用法,Java Kafka客户端并没有producer.config这个配置参数,所以直接在代码里用properties.put("producer.config", "...")是无效的。下面给你两种正确的实现方式:
方式一:直接在代码中设置SSL相关属性
把ssl.properties里的配置项逐个添加到Properties对象中,替换掉原来无效的那行代码:
public class MyProducer { public static void main(String[] args) throws Exception { Properties properties = new Properties(); // 基础配置 properties.put("bootstrap.servers", "<My bootstrap server>"); properties.put("key.serializer", StringSerializer.class.getName()); properties.put("value.serializer", StringSerializer.class.getName()); // 添加SSL相关配置 properties.put("security.protocol", "SSL"); properties.put("ssl.truststore.location", "<My custom value>"); properties.put("ssl.truststore.password", "<My custom value>"); properties.put("ssl.key.password", "<My custom value>"); properties.put("ssl.keystore.location", "<My custom value>"); properties.put("ssl.keystore.password", "<My custom value>"); KafkaProducer<String, String> kafkaProducer = new KafkaProducer<>(properties); // 注意:ProducerRecord第一个参数是topic名称,不是bootstrap server,你原代码此处有误 ProducerRecord<String, String> producerRecord = new ProducerRecord<>( "<My Kafka topic name>", "Hello World from program"); Future<RecordMetadata> future = kafkaProducer.send( producerRecord, (metadata, exception) -> { if(exception != null){ System.out.println("Something went wrong"); exception.printStackTrace(); } else { System.out.println("Successfully transmitted"); } }); future.get(); kafkaProducer.close(); } }
方式二:加载外部的SSL配置文件
如果不想在代码里硬编码SSL配置,可以直接把外部的ssl.properties文件加载到Properties对象中:
public class MyProducer { public static void main(String[] args) throws Exception { Properties properties = new Properties(); // 加载外部SSL配置文件 try (FileInputStream fis = new FileInputStream("/Users/DY/SSL/ssl.properties")) { properties.load(fis); } // 添加基础配置 properties.put("bootstrap.servers", "<My bootstrap server>"); properties.put("key.serializer", StringSerializer.class.getName()); properties.put("value.serializer", StringSerializer.class.getName()); KafkaProducer<String, String> kafkaProducer = new KafkaProducer<>(properties); ProducerRecord<String, String> producerRecord = new ProducerRecord<>( "<My Kafka topic name>", "Hello World from program"); Future<RecordMetadata> future = kafkaProducer.send( producerRecord, (metadata, exception) -> { if(exception != null){ System.out.println("Something went wrong"); exception.printStackTrace(); } else { System.out.println("Successfully transmitted"); } }); future.get(); kafkaProducer.close(); } }
内容的提问来源于stack exchange,提问作者DockYard
相关产品推荐
相关产品推荐

