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

如何用Scala在Databricks Workspace存储文件并供Executor访问?

解决方案:Databricks Spark作业通过Workspace文件连接Kafka

核心问题分析

Executor节点无法直接访问Databricks Workspace的文件系统——Workspace仅挂载在Driver节点,Executor属于独立计算节点,没有Workspace的访问权限。Scala直接读写/Workspace路径报错,是因为原生文件API没有适配Workspace的权限映射,而Python通过Databricks内部封装工具绕过了这个限制。


步骤1:将Workspace中的JKS文件分发到所有Executor

利用Spark的分布式缓存,把Driver可访问的Workspace文件复制到每个Executor的本地临时目录:

// 替换为你的JKS文件在Workspace的实际路径
val workspaceJksPath = "/Workspace/Shared/credentials/your-keystore.jks"

// 将文件添加到Spark分布式缓存,Executor启动时会自动下载该文件
spark.sparkContext.addFile(workspaceJksPath)

// 在Executor中获取本地缓存的JKS文件路径
val localJksPath = org.apache.spark.SparkFiles.get("your-keystore.jks")

步骤2:配置Kafka SSL参数

将Kafka连接的SSL配置指向Executor本地的JKS文件路径,而非Workspace路径:

val kafkaConfig = Map(
  "bootstrap.servers" -> "kafka-broker-host:9093",
  "security.protocol" -> "SSL",
  "ssl.truststore.location" -> localJksPath,
  "ssl.truststore.password" -> "your-truststore-pass",
  "ssl.keystore.location" -> localJksPath, // 若需客户端认证则添加
  "ssl.keystore.password" -> "your-keystore-pass",
  "subscribe" -> "target-topic"
)

// 示例:读取Kafka流
val kafkaStream = spark.readStream
  .format("kafka")
  .options(kafkaConfig)
  .load()

步骤3:解决Scala访问Workspace的权限问题

如果需要在Scala中读写Workspace文件,不要使用Java/Scala原生的File API,改用dbutils.fs工具(Databricks封装的文件操作工具,已适配Workspace权限):

// 读取Workspace文件内容
val jksContent = dbutils.fs.head(workspaceJksPath)

// 写入Workspace文件(需确保你有该目录的写入权限)
dbutils.fs.put("/Workspace/Shared/credentials/new-keystore.jks", jksContent, overwrite = true)

关键注意事项

  • 永远不要在Executor代码中直接引用/Workspace路径,Executor无法识别该文件系统。
  • 确保Workspace目标目录的权限正确:Shared目录默认开放读权限,写入权限需要管理员配置。
  • spark.sparkContext.addFile支持从Driver本地路径(包括Workspace挂载的/Workspace)添加文件,Spark会自动处理文件分发。

内容的提问来源于stack exchange,提问作者John Doe

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 07:58:15