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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 07:24:12