I'm trying to connect to a Kafka topic which has no security enabled and is part of our enterprise network. So I do not have any issues connecting to this. I tried to connect using spark structured streaming. Below is the piece of code.
val df = spark
.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "KAFKA-BROKER:PORT")
.option("subscribe", "KAFKA-TOPIC")
.option("maxOffsetsPerTrigger", 20)
.option("startingOffsets", "earliest")
.load()
This gives me a dataframe with the below schema.
df:org.apache.spark.sql.DataFrame
key:binary
value:binary
topic:string
partition:integer
offset:long
timestamp:timestamp
timestampType:integer
key and value are in binary format. So I tried connecting to Schema Registry. This has SSL enabled on it. I got the .pk12 from the kafka team to connect to schema registry. I tested this p12 file by curl command.(I extracted client_cert.pem, ca_cert.pem and client_key.pem). The curl Command was successful.
%sh
cd /dbfs/FileStore/tables/
curl -s --cert ./client_cert.pem --key client_key.pem --tlsv1.2 --cacert ca_cert.pem https://sc-soh-dev01-schema-registry.px-npe2105.pks.t-mobile.com:30301/subjects/
I got the result from the above curl command which listed the available schemas.
["kafka-topic-value","kafka-topic1-value"]
I tried importing the p12 file into databricks Keystore
keytool -importkeystore -srckeystore /dbfs/FileStore/tables/security.p12 -srcstoretype pkcs12 -destkeystore /usr/lib/jvm/java-8-openjdk-amd64/jre/lib/security/cacerts -deststoretype JKS
I've configured the below spark configurations in Spark config at cluster level while starting the databricks cluster.
spark.databricks.userInfoFunctions.enabled true
spark.driver.extraJavaOptions -Djavax.net.ssl.trustStore=/usr/lib/jvm/zulu8-ca-amd64/jre/lib/security/cacerts -Djava.security.properties=/databricks/spark/dbconf/java/extra_v1.security
spark.databricks.delta.preview.enabled true
spark.executor.extraJavaOptions -Djavax.net.ssl.trustStore=/usr/lib/jvm/zulu8-ca-amd64/jre/lib/security/cacerts -Djava.security.properties=/databricks/spark/dbconf/java/extra_v1.security
and then tried to stream using from_avro.
val df = spark
.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "KAFKA-BROKER:PORT")
.option("subscribe", "KAFKA-TOPIC")
.option("maxOffsetsPerTrigger", 20)
.option("startingOffsets", "earliest")
.load()
.select(from_avro($"value", "KAFKA-TOPIC-VALUE", "https://SCHEMA-REGISTRY:30301/").as("value"))
I get the below error.
SSLHandshakeException: Received fatal alert: bad_certificate
at sun.security.ssl.Alert.createSSLException(Alert.java:131)
at sun.security.ssl.Alert.createSSLException(Alert.java:117)
at sun.security.ssl.TransportContext.fatal(TransportContext.java:335)
at sun.security.ssl.Alert$AlertConsumer.consume(Alert.java:293)
at sun.security.ssl.TransportContext.dispatch(TransportContext.java:185)
at sun.security.ssl.SSLTransport.decode(SSLTransport.java:156)
at sun.security.ssl.SSLSocketImpl.decode(SSLSocketImpl.java:1418)
I used the below options with the spark structured streaming as well.
.option("schema.registry.url", "https://SCHEMA-REGISTRY:30301/")
.option("schema.registry.ssl.truststore.location", "/tmp/sohkeystore")
.option("schema.registry.ssl.truststore.password", "123456")
.option("schema.registry.ssl.keystore.location", "/tmp/sohkeystore")
.option("schema.registry.ssl.keystore.password", "123456")
.option("schema.registry.ssl.key.password", "123456")
I tried using init scripts following the below link. https://docs.microsoft.com/en-us/azure/databricks/kb/python/import-custom-ca-cert
No luck with any of those above options. Can anyone help me with where I'm going wrong?