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
相关产品推荐
相关产品推荐

