【问题标题】:Move file from one folder to another on HDFS in Scala / Spark在 Scala / Spark 中的 HDFS 上将文件从一个文件夹移动到另一个文件夹
【发布时间】:2018-06-21 22:51:46
【问题描述】:

我有两条路径,一条用于文件,一条用于文件夹。我想将文件移动到 HDFS 上的那个文件夹中。我怎么能在 Scala 中做到这一点?我也在用 Spark

如果相同的代码也适用于 Windows 路径,就像在 HDFS 上读取/写入文件一样,但不是必需的。

我尝试了以下方法:

val fs = FileSystem.get(sc.hadoopConfiguration)
fs.moveFromLocalFile(something, something2)

我收到以下错误:

线程“主”java.lang.IllegalArgumentException 中的异常:错误 FS:hdfs:/user/o/datasets/data.txt,预期:file:///

moveToLocalFile() 也是如此,因为它们旨在在文件系统之间而不是在文件系统内传输文件。我也尝试过fs.rename(),但这根本没有做任何事情(没有错误或任何事情)。

我基本上在一个目录中创建文件(使用流写入它们),一旦完成,他们需要移动到另一个目录。这个不同的目录由 Spark 流监控,当 Spark 流尝试处理未完成的文件时我遇到了一些问题

【问题讨论】:

  • Spark 流尝试处理未完成的文件。您需要明确忽略以句点或下划线开头的任何文件
  • 当我创建文件时,它们的临时形式仍然具有相同的文件名,但它们的大小为 0(字节),直到它们完成,然后它们具有最终大小和相同的名称。跨度>
  • 是的,除非您忽略它们,否则 Spark Streaming 会引发错误
  • 如何检测程序中的大小?由于文件名没有改变
  • 问题没看懂,不过好像和原帖无关

标签: scala hadoop apache-spark hdfs


【解决方案1】:

试试下面的 Scala 代码。

import org.apache.hadoop.conf.Configuration
import org.apache.hadoop.fs.FileSystem
import org.apache.hadoop.fs.Path

val hadoopConf = new Configuration()
val hdfs = FileSystem.get(hadoopConf)

val srcPath = new Path(srcFilePath)
val destPath = new Path(destFilePath)

hdfs.copyFromLocalFile(srcPath, destPath)

您还应该检查 Spark 是否在 conf/spark-env.sh 文件中设置了 HADOOP_CONF_DIR 变量。这将确保 Spark 能够找到 Hadoop 配置设置。

build.sbt 文件的依赖关系:

libraryDependencies += "org.apache.hadoop" % "hadoop-common" % "2.6.0"
libraryDependencies += "org.apache.commons" % "commons-io" % "1.3.2"
libraryDependencies += "org.apache.hadoop" % "hadoop-hdfs" % "2.6.0"

您可以使用 apache commons 中的 IOUtils 将数据从 InputStream 复制到 OutputStream

import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;

import org.apache.commons.io.IOUtils;



val hadoopconf = new Configuration();
val fs = FileSystem.get(hadoopconf);

//Create output stream to HDFS file
val outFileStream = fs.create(new Path("hdfs://<namenode>:<port>/output_path"))

//Create input stream from local file
val inStream = fs.open(new Path("hdfs://<namenode>:<port>/input_path"))

IOUtils.copy(inStream, outFileStream)

//Close both files
inStream.close()
outFileStream.close()

【讨论】:

  • 不幸的是第一个解决方案不起作用,我如何检查是否设置了HADOOP_CONF_DIR?第二种解决方案也不适用于我的系统。我基本上在一个目录中创建文件(用流写入它们),一旦完成,他们需要移动到另一个目录。这个不同的目录由 Spark 流监控,当 Spark 流尝试处理未完成的文件时,我遇到了一些问题。
  • @osk 你的问题没有提到 Spark... 而HADOOP_CONF_DIR 是一个环境变量,所以搜索你如何为各自的操作系统寻找它们,或者如果你正在使用 Spark,然后打开spark-env.sh 文件,并将其设置在那里
  • @Sahil,我正在研究相同的解决方案,并试图找到一种以分布式方式复制大型数据集的方法,因为我看到 IOUtils 是一个非 Hadoop 包 org.apache.commons。 io.IOUtils 它可能无法以分布式方式工作。您能否确认 IOUtis 可以在分布式文件副本中工作。我正在尝试将 HDFS 中的文件文件复制到同一集群上的另一个 HDFS 目录
【解决方案2】:
import org.apache.hadoop.fs.{FileAlreadyExistsException, FileSystem, FileUtil, Path}

val srcFileSystem: FileSystem = FileSystemUtil
  .apply(spark.sparkContext.hadoopConfiguration)
  .getFileSystem(sourceFile)
val dstFileSystem: FileSystem = FileSystemUtil
  .apply(spark.sparkContext.hadoopConfiguration)
  .getFileSystem(sourceFile)
FileUtil.copy(
  srcFileSystem,
  new Path(new URI(sourceFile)),
  dstFileSystem,
  new Path(new URI(targetFile)),
  true,
  spark.sparkContext.hadoopConfiguration)

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多