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

Apache Camel Kafka组件连接Azure Event Hub时的SASL认证错误

Apache Camel Kafka组件连接Azure Event Hub SASL认证失败问题排查与解决

问题背景

在Windows 10系统、Java 17环境下,使用Apache Camel 4.4.3的Kafka组件向Azure Event Hub(Premium)发送消息时,SASL认证阶段持续失败。但其他Kafka客户端(Kafka UI、Confluent Java客户端、Apache Camel Azure Eventhubs组件)均能成功连接并发送消息。

核心配置代码

camelContext.addRoutes(new RouteBuilder() {
    @Override
    public void configure() {
        String kafkaEndpointUri = "kafka:{{eventhub.name}}?brokers={{eventhub.namespace}}.servicebus.windows.net:9093"
            + "&securityProtocol=SASL_SSL"
            + "&saslMechanism=PLAIN"
            + "&saslJaasConfig=org.apache.kafka.common.security.plain.PlainLoginModule required username=\"$ConnectionString\" password=\"{{eventhub.connectionstring}}\";";

        from("timer://foo?repeatCount=1")
        .setBody(constant("Hello from Camel to Azure EventHub!"))
        .to(kafkaEndpointUri);
    }
});

认证阶段异常日志

[kafka-producer-network-thread | producer-1] DEBUG org.apache.kafka.common.security.authenticator.SaslClientAuthenticator - [Producer clientId=producer-1] Set SASL client state to SEND_APIVERSIONS_REQUEST
[main] DEBUG org.apache.camel.impl.DefaultCamelContext - start() took 536 millis
Camel application is running. Press Ctrl + C to terminate.
[kafka-producer-network-thread | producer-1] DEBUG org.apache.kafka.common.security.authenticator.SaslClientAuthenticator - [Producer clientId=producer-1] Creating SaslClient: client=null;service=kafka;serviceHostname=EVENTHUB_NAMESPACE.servicebus.windows.net;mechs=[PLAIN]
[kafka-producer-network-thread | producer-1] DEBUG org.apache.kafka.common.network.Selector - [Producer clientId=producer-1] Created socket with SO_RCVBUF = 65536, SO_SNDBUF = 131072, SO_TIMEOUT = 0 to node -1
[kafka-producer-network-thread | producer-1] DEBUG org.apache.kafka.clients.NetworkClient - [Producer clientId=producer-1] Completed connection to node -1. Fetching API versions.
[kafka-producer-network-thread | producer-1] DEBUG org.apache.kafka.common.network.SslTransportLayer - [SslTransportLayer channelId=-1 key=channel=java.nio.channels.SocketChannel[connection-pending remote=EVENTHUB_NAMESPACE.servicebus.windows.net/XXX.XXX.XXX.XXX:9093], selector=sun.nio.ch.WEPollSelectorImpl@2643665b, interestOps=8, readyOps=0] SSL handshake completed successfully with peerHost 'EVENTHUB_NAMESPACE.servicebus.windows.net' peerPort 9093 peerPrincipal 'CN=servicebus.windows.net, O=Microsoft Corporation, L=Redmond, ST=WA, C=US' protocol 'TLSv1.3' cipherSuite 'TLS_AES_256_GCM_SHA384'
[kafka-producer-network-thread | producer-1] DEBUG org.apache.kafka.common.security.authenticator.SaslClientAuthenticator - [Producer clientId=producer-1] Set SASL client state to RECEIVE_APIVERSIONS_RESPONSE
[kafka-producer-network-thread | producer-1] DEBUG org.apache.kafka.common.security.authenticator.SaslClientAuthenticator - [Producer clientId=producer-1] Set SASL client state to SEND_HANDSHAKE_REQUEST
[kafka-producer-network-thread | producer-1] DEBUG org.apache.kafka.common.security.authenticator.SaslClientAuthenticator - [Producer clientId=producer-1] Set SASL client state to RECEIVE_HANDSHAKE_RESPONSE
[kafka-producer-network-thread | producer-1] DEBUG org.apache.kafka.common.security.authenticator.SaslClientAuthenticator - [Producer clientId=producer-1] Set SASL client state to INITIAL
[kafka-producer-network-thread | producer-1] DEBUG org.apache.kafka.common.security.authenticator.SaslClientAuthenticator - [Producer clientId=producer-1] Set SASL client state to INTERMEDIATE
[Camel (camel-1) thread #1 - timer://foo] DEBUG org.apache.camel.processor.SendProcessor - >>>> kafka://EVENTHUB_NAME?brokers=EVENTHUB_NAMESPACE.servicebus.windows.net%3A9093&saslJaasConfig=xxxxxx&saslMechanism=PLAIN&securityProtocol=SASL_SSL Exchange[]
[Camel (camel-1) thread #1 - timer://foo] DEBUG org.apache.camel.component.kafka.KafkaProducer - Sending message to topic: EVENTHUB_NAME, partition: null, key: null
[kafka-producer-network-thread | producer-1] WARN  org.apache.kafka.common.network.Selector - [Producer clientId=producer-1] Unexpected error from EVENTHUB_NAMESPACE.servicebus.windows.net/XXX.XXX.XXX.XXX (channelId=-1); closing connection
java.lang.RuntimeException: non-nullable field authBytes was serialized as null

其他可行客户端代码

Apache Camel Event Hubs组件

...
.to(String.format("azure-eventhubs:?connectionString=RAW(%s)", connectionStringWithTopic));

Confluent Java客户端

Properties properties = new Properties();
properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, brokerList);
properties.put("security.protocol", "SASL_SSL");
properties.put("sasl.mechanism", "PLAIN");
properties.put("sasl.jaas.config", String.format("org.apache.kafka.common.security.plain.PlainLoginModule required username=\"$ConnectionString\" password=\"%s\";", connectionString));
properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
properties.put(ProducerConfig.CLIENT_ID_CONFIG, "KafkaProducer");
Producer<String, String> producer = new KafkaProducer<>(properties);
ProducerRecord<String, String> record = new ProducerRecord<>(eventHubName, "key", "This is a message from regular Java app using Confluent Kafka!");
producer.send(record);

问题原因与解决方案

原因分析

异常java.lang.RuntimeException: non-nullable field authBytes was serialized as null本质是Kafka客户端收到的SASL认证内容不完整,根源在于Camel的URI参数解析机制:Azure Event Hub的连接字符串包含大量特殊字符(如=、;),Camel默认会将这些字符视为URI参数的分隔符,导致saslJaasConfig的配置被截断或错误解析,无法完整传递给Kafka客户端。

解决方案1:使用RAW()包裹saslJaasConfig参数

修改Kafka端点URI,将saslJaasConfig的内容用RAW()包裹,告诉Camel不要解析该参数内的特殊字符,直接完整传递给Kafka组件:

camelContext.addRoutes(new RouteBuilder() {
    @Override
    public void configure() {
        String kafkaEndpointUri = "kafka:{{eventhub.name}}?brokers={{eventhub.namespace}}.servicebus.windows.net:9093"
            + "&securityProtocol=SASL_SSL"
            + "&saslMechanism=PLAIN"
            + "&saslJaasConfig=RAW(org.apache.kafka.common.security.plain.PlainLoginModule required username=\"$ConnectionString\" password=\"{{eventhub.connectionstring}}\";)";

        from("timer://foo?repeatCount=1")
        .setBody(constant("Hello from Camel to Azure EventHub!"))
        .to(kafkaEndpointUri);
    }
});

解决方案2:直接设置Kafka原生Properties

绕过Camel的URI参数解析,通过KafkaConstants.KAFKA_PRODUCER_CONFIGS直接传递Kafka原生配置,避免特殊字符解析问题:

camelContext.addRoutes(new RouteBuilder() {
    @Override
    public void configure() {
        from("timer://foo?repeatCount=1")
        .setBody(constant("Hello from Camel to Azure EventHub!"))
        .to("kafka:{{eventhub.name}}")
        .setProperty(KafkaConstants.KAFKA_PRODUCER_CONFIGS, () -> {
            Properties props = new Properties();
            props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "{{eventhub.namespace}}.servicebus.windows.net:9093");
            props.put("security.protocol", "SASL_SSL");
            props.put("sasl.mechanism", "PLAIN");
            props.put("sasl.jaas.config", "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"$ConnectionString\" password=\"{{eventhub.connectionstring}}\";");
            return props;
        });
    }
});

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 16:12:06