Databricks读取Kafka流遇SSL异常:密钥库文件无法加载
问题描述
在Databricks中通过SSL密钥库认证读取Kafka流时,出现连接失败问题,提示无法找到密钥库文件,但文件实际存在于DBFS中,尝试挂载ADLS Gen2读取后问题依旧。
错误信息
Caused by: kafkashaded.org.apache.kafka.common.KafkaException: Failed to load SSL keystore /dbfs/FileStore/Certs/client.keystore.jks Caused by: java.nio.file.NoSuchFileException: /dbfs/FileStore/Certs/client.keyst
驱动日志额外错误:
22/11/04 12:18:07 ERROR DefaultSslEngineFactory: Modification time of key store could not be obtained
使用的代码
df = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers","kafka server with port") \ .option("kafka.security.protocol", "SSL") \ .option("kafka.ssl.truststore.location",'/dbfs/FileStore/Certs/client.truststore.jks' ) \ .option("kafka.ssl.keystore.location", '/dbfs/FileStore/Certs/client.keystore.jks') \ .option("kafka.ssl.keystore.password", keystore_pass) \ .option("kafka.ssl.truststore.password", truststore_pass) \ .option("kafka.ssl.keystore.type", "JKS") \ .option("kafka.ssl.truststore.type", "JKS") \ .option("subscribe","sports") \ .option("startingOffsets", "earliest") \ .load()
解决方案
替换DBFS路径格式
Kafka客户端基于JVM运行,无法直接识别/dbfs本地模拟路径,需改用dbfs:/格式的路径,修改代码中的密钥库配置:.option("kafka.ssl.truststore.location",'dbfs:/FileStore/Certs/client.truststore.jks' ) \ .option("kafka.ssl.keystore.location", 'dbfs:/FileStore/Certs/client.keystore.jks') \复制密钥库到集群本地临时目录
如果路径格式修改无效,可将DBFS中的密钥库文件复制到集群节点的本地临时目录,再引用本地路径:# 复制文件到本地临时目录 dbutils.fs.cp("dbfs:/FileStore/Certs/client.truststore.jks", "file:/tmp/client.truststore.jks") dbutils.fs.cp("dbfs:/FileStore/Certs/client.keystore.jks", "file:/tmp/client.keystore.jks") # 更新读取流的路径配置 df = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers","kafka server with port") \ .option("kafka.security.protocol", "SSL") \ .option("kafka.ssl.truststore.location",'/tmp/client.truststore.jks' ) \ .option("kafka.ssl.keystore.location", '/tmp/client.keystore.jks') \ .option("kafka.ssl.keystore.password", keystore_pass) \ .option("kafka.ssl.truststore.password", truststore_pass) \ .option("kafka.ssl.keystore.type", "JKS") \ .option("kafka.ssl.truststore.type", "JKS") \ .option("subscribe","sports") \ .option("startingOffsets", "earliest") \ .load()检查密钥库文件完整性
日志中的修改时间获取失败提示,可能是文件损坏或权限不足导致。可通过以下命令验证:# 查看DBFS目录下的文件信息 dbutils.fs.ls("dbfs:/FileStore/Certs/") # 用keytool验证密钥库有效性 %sh keytool -list -v -keystore /dbfs/FileStore/Certs/client.keystore.jks -storepass <你的密钥库密码>若验证失败,重新上传密钥库文件到DBFS。
确认集群权限配置
检查集群实例是否有DBFS对应目录的访问权限;若使用ADLS Gen2挂载,确认挂载点的权限设置正确,服务主体拥有存储容器的读取权限。
内容的提问来源于stack exchange,提问作者JayanthGoulla
相关产品推荐
相关产品推荐

