Spark Databricks(Scala)下二进制数据写入、读取与删除方法咨询
Databricks Scala 单字节数组二进制文件读写删方案
核心说明
你需要操作的是单节点内存中的Array[Byte]变量,属于单文件读写场景,直接使用Hadoop FileSystem API(操作挂载ADLS)或Java标准IO(操作本地路径)即可,不需要调用Spark分布式读写API,避免不必要的格式转码。
场景1:操作挂载的Azure Data Lake路径(/mnt/somewhere/tmp/data.bin)
写入字节数组到目标路径
import org.apache.hadoop.fs.{FileSystem, Path} // 替换为你实际的字节数组变量 val byteContent: Array[Byte] = ??? // 初始化FileSystem,自动继承挂载存储的认证权限 val hadoopConf = spark.sparkContext.hadoopConfiguration val fs = FileSystem.get(hadoopConf) val adlsPath = new Path("/mnt/somewhere/tmp/data.bin") // 写入数据,第二个参数true表示覆盖已存在的同名文件 val outputStream = fs.create(adlsPath, true) outputStream.write(byteContent) outputStream.close()
读取文件为InputStream
val inputStream = fs.open(adlsPath) // 此处编写你的InputStream处理逻辑 // 处理完成后关闭流 inputStream.close()
使用完成后删除文件
// 第二个参数false表示非递归删除,仅删除单个文件 fs.delete(adlsPath, false) fs.close()
场景2:操作Driver节点本地/tmp路径(file:/tmp/data.bin)
写入字节数组到本地路径
import java.io.{FileOutputStream, FileInputStream, File} val byteContent: Array[Byte] = ??? val localFile = new File("/tmp/data.bin") val outputStream = new FileOutputStream(localFile) outputStream.write(byteContent) outputStream.close()
读取为InputStream
val inputStream = new FileInputStream(localFile) // 此处编写InputStream处理逻辑 inputStream.close()
删除本地文件
localFile.delete()
注意事项
- 以上代码可直接在Databricks Scala Notebook中运行,无需额外引入依赖
- 本地路径为Driver节点的本地存储,仅适合仅在Driver侧处理的场景,Executor节点无法访问该路径
- 操作ADLS挂载路径前,请确认当前集群账号对目标路径有读写删除权限
- 不要使用Spark原生的
writeAPI处理单二进制文件,该API为分布式数据集设计,会生成多分区文件且会引入不必要的格式转码,破坏原始二进制内容
内容的提问来源于stack exchange,提问作者Karzyfox
相关产品推荐
相关产品推荐

