Вот основная часть настройки кода конечная точка Кафки:
Код: Выделить всё
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
Код: Выделить всё
...
.to(String.format("azure-eventhubs:?connectionString=RAW(%s)", connectionStringWithTopic));
Код: Выделить всё
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\";", connectioString));
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 producer = new KafkaProducer(properties);
ProducerRecord record = new ProducerRecord(eventHubName, "key", "This is a message from regular Java app using Confluent Kafka!");
producer.send(record);
Кто-нибудь сталкивался с подобными проблемами или мог бы дать представление о том, что может произойти? здесь не так? Могут ли быть в Apache Camel какие-то особые настройки или конфигурации, которые я мог упустить из виду?
Подробнее здесь: https://stackoverflow.com/questions/787 ... zure-event