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

Databricks读取S3中Kafka证书文件遇路径错误求助

Databricks读取Kafka证书时S3路径解析异常的解决思路

问题场景

在Databricks Notebook中读取Kafka证书时,先触发以下错误:

DbxDlTransferError: Terminated with exception: Kafka store file location only supports external location or UC Volume path on Shared cluster. Use external location or UC Volume Path to provide it.: None

改用S3路径(s3://my-bucket/tmp/k2/truststore.jks)后,又出现文件找不到的异常,核心错误信息显示传入的双斜杠//被Java解析为单斜杠/:

Caused by: java.nio.file.NoSuchFileException: s3:/my-bucket/tmp/k2/truststore.jks
at sun.nio.fs.UnixException.translateToIOException(UnixException.java:86)
at sun.nio.fs.UnixException.rethrowAsIOException(UnixException.java:102)
at sun.nio.fs.UnixException.rethrowAsIOException(UnixException.java:107)
at sun.nio.fs.UnixFileSystemProvider.newByteChannel(UnixFileSystemProvider.java:214)
at java.nio.file.Files.newByteChannel(Files.java:361)
at java.nio.file.Files.newByteChannel(Files.java:407)
at java.nio.file.spi.FileSystemProvider.newInputStream(FileSystemProvider.java:384)
at java.nio.file.Files.newInputStream(Files.java:152)
at kafkashaded.org.apache.kafka.common.security.ssl.DefaultSslEngineFactory$FileBasedStore.load(DefaultSslEngineFactory.java:368)
... 86 more

涉及代码:

truststore_location = "s3://my-bucket/tmp/k2/truststore.jks"
cluster_ca_certificate_location = "s3://my-bucket/tmp/k2/cluster-ca-certificate.pem"


kafka_server = "server_ip:9093"

kafka_topic = "kafka_topic"
kafka_group_id = "group_id"
scram_login_module = 'org.apache.kafka.common.security.scram.ScramLoginModule required username="" password=""'

input_df = (
  spark
  .read
  .format("kafka")
  .option("kafka.bootstrap.servers", kafka_server)
  .option("kafka.ssl.truststore.location", truststore_location)
  .option("kafka.ssl.truststore.password", truststore_pass)
  .option("kafka.ssl.ca.location", cluster_ca_certificate_location)
  .option("kafka.security.protocol", "SASL_SSL")
  .option("kafka.sasl.mechanism", "SCRAM-SHA-256")
  .option("kafka.sasl.jaas.config", "kafkashaded.org.apache.kafka.common.security.scram.ScramLoginModule required username='{}' password='{}';".format(consumer_username, consumer_password))
  .option("subscribe", kafka_topic)
  .option("kafka.group.id", kafka_group_id)
  .option("failOnDataLoss", "false")
  .load()
)

解决思路

核心原因

Kafka SSL配置项底层依赖Java NIO读取文件,而Java原生文件系统不支持S3路径格式,会自动将s3://解析为s3:/,导致路径无效。需通过Databricks支持的方式让Kafka客户端访问证书文件。

方案1:下载证书到集群本地临时目录

通过Databricks的文件操作API将S3上的证书下载到集群本地临时路径,再配置Kafka使用本地路径:

# 从S3下载证书到本地临时目录
dbutils.fs.cp("s3://my-bucket/tmp/k2/truststore.jks", "file:/tmp/truststore.jks")
dbutils.fs.cp("s3://my-bucket/tmp/k2/cluster-ca-certificate.pem", "file:/tmp/cluster-ca-certificate.pem")

# 更新路径为本地路径
truststore_location = "/tmp/truststore.jks"
cluster_ca_certificate_location = "/tmp/cluster-ca-certificate.pem"

# 后续Kafka读取代码保持不变

注意:临时目录在集群重启后会被清空,若需持久化,可改用集群挂载的云存储本地路径,或在Notebook启动阶段自动执行下载逻辑。

方案2:使用Unity Catalog Volume路径

根据初始错误提示,Shared集群支持UC Volume路径,可将证书上传到UC Volume后直接引用:

  1. 在Unity Catalog中创建Volume,上传证书文件;
  2. 替换Kafka配置中的路径为Volume路径:
truststore_location = "/Volumes/[catalog_name]/[schema_name]/[volume_name]/truststore.jks"
cluster_ca_certificate_location = "/Volumes/[catalog_name]/[schema_name]/[volume_name]/cluster-ca-certificate.pem"

此方式无需手动下载文件,Databricks会自动处理路径映射,适合长期稳定使用场景。

方案3:挂载S3桶到DBFS后引用

将S3桶挂载到DBFS,使用DBFS本地路径配置Kafka:

  1. 挂载S3桶到DBFS(需确保集群有权限访问S3):
dbutils.fs.mount(
  source="s3://my-bucket/tmp/k2/",
  mount_point="/mnt/kafka_certs",
  extra_configs={"fs.s3a.access.key": "你的访问密钥", "fs.s3a.secret.key": "你的密钥"}
)
  1. 使用挂载后的DBFS路径:
truststore_location = "/dbfs/mnt/kafka_certs/truststore.jks"
cluster_ca_certificate_location = "/dbfs/mnt/kafka_certs/cluster-ca-certificate.pem"

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 09:52:02