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

如何实现基于用户名密码认证的Camel Kafka生产者与消费者(无SSL证书)

使用Apache Camel实现带SASL用户名密码认证的Kafka生产者与消费者

没问题,我来给你详细讲讲怎么配置完全匹配你需求的Camel Kafka生产者和消费者——用用户名密码做SASL认证,不需要SSL证书,刚好适配你提供的Broker信息。

首先先明确核心配置逻辑:Camel Kafka组件支持直接传递原生Kafka客户端的参数,所以我们只需要指定SASL相关的认证参数即可,重点是securityProtocol、saslMechanism和saslJaasConfig这三个关键项。

一、先整理你的Broker配置

把你提供的多个Broker地址合并成逗号分隔的字符串:

kafka04-orgName.services.orgsj001.us-west.bluemix.net:9093,kafka02-orgName.services.orgsj001.us-west.bluemix.net:9093,kafka01-orgName.services.orgsj001.us-west.bluemix.net:9093,kafka03-orgName.services.orgsj001.us-west.bluemix.net:9093

二、生产者配置示例

1. Java DSL 实现

import org.apache.camel.builder.RouteBuilder;

public class KafkaProducerRoute extends RouteBuilder {
    @Override
    public void configure() throws Exception {
        // 从direct端点接收消息,发送到Kafka
        from("direct:kafka-producer")
            .to("kafka:你的主题名称?" +
                "brokers=kafka04-orgName.services.orgsj001.us-west.bluemix.net:9093,kafka02-orgName.services.orgsj001.us-west.bluemix.net:9093,kafka01-orgName.services.orgsj001.us-west.bluemix.net:9093,kafka03-orgName.services.orgsj001.us-west.bluemix.net:9093&" +
                "securityProtocol=SASL_PLAINTEXT&" +
                "saslMechanism=PLAIN&" +
                "saslJaasConfig=org.apache.kafka.common.security.plain.PlainLoginModule required username=\"**********\" password=\"********\";")
            .log("消息已成功发送到Kafka主题: ${headers[kafka.TOPIC]}");
    }
}

参数说明:

  • securityProtocol=SASL_PLAINTEXT:指定使用SASL明文认证(无需SSL加密,符合你的需求)
  • saslMechanism=PLAIN:采用最常用的用户名密码认证机制
  • saslJaasConfig:配置Kafka的JAAS登录模块,填入你的用户名和密码

2. Spring XML 实现

如果用Spring XML配置,注意双引号需要转义为%22:

<route id="kafkaProducerRoute">
    <from uri="direct:kafka-producer"/>
    <to uri="kafka:你的主题名称?
        brokers=kafka04-orgName.services.orgsj001.us-west.bluemix.net:9093,kafka02-orgName.services.orgsj001.us-west.bluemix.net:9093,kafka01-orgName.services.orgsj001.us-west.bluemix.net:9093,kafka03-orgName.services.orgsj001.us-west.bluemix.net:9093
        &securityProtocol=SASL_PLAINTEXT
        &saslMechanism=PLAIN
        &saslJaasConfig=org.apache.kafka.common.security.plain.PlainLoginModule required username=%22**********%22 password=%22********%22;"/>
    <log message="消息已成功发送到Kafka主题: ${headers[kafka.TOPIC]}"/>
</route>

三、消费者配置示例

消费者的SASL认证配置和生产者一致,额外需要指定消费者组ID(Kafka必须参数):

1. Java DSL 实现

import org.apache.camel.builder.RouteBuilder;

public class KafkaConsumerRoute extends RouteBuilder {
    @Override
    public void configure() throws Exception {
        // 从Kafka主题消费消息
        from("kafka:你的主题名称?" +
                "brokers=kafka04-orgName.services.orgsj001.us-west.bluemix.net:9093,kafka02-orgName.services.orgsj001.us-west.bluemix.net:9093,kafka01-orgName.services.orgsj001.us-west.bluemix.net:9093,kafka03-orgName.services.orgsj001.us-west.bluemix.net:9093&" +
                "groupId=你的消费者组ID&" +
                "securityProtocol=SASL_PLAINTEXT&" +
                "saslMechanism=PLAIN&" +
                "saslJaasConfig=org.apache.kafka.common.security.plain.PlainLoginModule required username=\"**********\" password=\"********\";")
            .log("从Kafka收到消息: ${body}")
            // 这里添加你的业务处理逻辑
            .to("direct:process-message");
    }
}

2. Spring XML 实现

<route id="kafkaConsumerRoute">
    <from uri="kafka:你的主题名称?
        brokers=kafka04-orgName.services.orgsj001.us-west.bluemix.net:9093,kafka02-orgName.services.orgsj001.us-west.bluemix.net:9093,kafka01-orgName.services.orgsj001.us-west.bluemix.net:9093,kafka03-orgName.services.orgsj001.us-west.bluemix.net:9093
        &groupId=你的消费者组ID
        &securityProtocol=SASL_PLAINTEXT
        &saslMechanism=PLAIN
        &saslJaasConfig=org.apache.kafka.common.security.plain.PlainLoginModule required username=%22**********%22 password=%22********%22;"/>
    <log message="从Kafka收到消息: ${body}"/>
    <to uri="direct:process-message"/>
</route>

四、重要注意事项

  • 安全最佳实践:不要把用户名密码硬编码在代码里!建议用环境变量、Spring配置文件(如application.properties)或密钥管理工具注入。比如Spring Boot中可以这样配置:
    camel.kafka.producer.sasl-jaas-config=org.apache.kafka.common.security.plain.PlainLoginModule required username="${kafka.user}" password="${kafka.password}";
    
    然后通过环境变量或配置中心传入kafka.user和kafka.password。
  • 版本兼容性:确保Camel Kafka组件版本与Kafka客户端版本匹配,比如Camel 3.x对应Kafka 2.x,Camel 4.x对应Kafka 3.x/4.x,避免版本冲突导致认证失败。
  • Broker端验证:确认你的Kafka Broker已经开启了SASL_PLAINTEXT监听,并且配置了PLAIN认证机制,否则客户端会连接失败。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 09:36:25