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

Scalding单元测试:如何实现TypedPipe本地文件写入及验证?

为Scalding TypedPipe.write编写单元测试:本地写入与MiniDFSCluster集成方案

问题背景

公司通过自定义API增强Scalding的写入功能,用于追踪数据集元数据。在将普通写入切换为该特殊写入时,需要对比转换前后Key/Value、TSV/CSV、Thrift等各类数据集的二进制文件一致性。当前需要解决的核心问题是:如何为TypedPipe的.write方法编写单元测试,且无法使用Scalding的JobTest框架(因实际写入的数据会被bijection codec包装或通过case class封装后,再由元数据写入API执行写入)。

本地写入尝试(未生成文件)

以下代码在OSX本地磁盘未写入任何内容:

implicit val timeZone: TimeZone = DateOps.UTC
implicit val dateParser: DateParser = DateParser.default
implicit def flowDef: FlowDef = new FlowDef()
implicit def mode: Mode = Local(true)

val fileStrPath = root + "/test"

println("writing data to " + fileStrPath)

TypedPipe
  .from(Seq[Long](1, 2, 3, 4, 5))
  // .map((x: Long) => { println(x.toString); System.out.flush(); x })
  .write(TypedTsv[Long](fileStrPath))
  .forceToDisk

MiniDFSCluster尝试(无法关联Scalding写入)

尝试搭建MiniDFSCluster,但不知道如何将其与Scalding的写入逻辑关联,代码如下:

def setUpTempFolder: String = {
  val tempFolder = new TemporaryFolder
  tempFolder.create()
  tempFolder.getRoot.getAbsolutePath
}
val root: String = setUpTempFolder
println(s"root = $root")
val tempDir = Files.createTempDirectory(setUpTempFolder).toFile
val hdfsCluster: MiniDFSCluster = {
  val configuration = new Configuration()
  configuration.set(MiniDFSCluster.HDFS_MINIDFS_BASEDIR, tempDir.getAbsolutePath)
  configuration.set("io.compression.codecs", classOf[LzopCodec].getName)
  new MiniDFSCluster.Builder(configuration)
    .manageNameDfsDirs(true)
    .manageDataDfsDirs(true)
    .format(true)
    .build()
}
hdfsCluster.waitClusterUp()
val fs: DistributedFileSystem = hdfsCluster.getFileSystem
val rootPath = new Path(root)
fs.mkdirs(rootPath)

需求

  • 不使用Scalding REPL,实现Scalding本地文件写入的可行方法
  • 若采用MiniDFSCluster,需明确:
    • 如何将其与Scalding写入逻辑关联
    • 写入完成后读取文件的方法

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 12:24:19