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
相关产品推荐
相关产品推荐

