【发布时间】:2017-12-20 23:25:20
【问题描述】:
我有一个 Scala 列表:fileNames,其中包含本地目录中存在的文件名。 例如:
fileNames(2)
res0: String = file:///tmp/audits/xx_user.log
我正在尝试使用 Scala 将列表中的文件:fileNames 从本地移动到 HDFS。为此,我按照以下步骤操作:
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();
hadoopconf.addResource(new Path("/etc/hadoop/conf/core-site.xml"));
val fs = FileSystem.get(hadoopconf);
val outFileStream = fs.create(new Path("hdfs://mydev/user/devusr/testfolder"))
代码在此处运行良好。当我尝试添加 inputStream 时,我收到如下错误消息:
val inStream = fs.open(new Path(fileNames(2)))
java.lang.IllegalArgumentException: Wrong FS: file:/tmp/audits/xx_user.log, expected: hdfs://mergedev
我也试过直接指定文件名,结果是一样的:
val inStream = fs.open(new Path("file:///tmp/audits/xx_user.log"))
java.lang.IllegalArgumentException: Wrong FS: file:/tmp/audits/xx_user.log, expected: hdfs://mergedev
但是当我尝试将文件直接加载到 spark 中时,它工作正常:
val localToSpark = spark.read.text(fileNames(2))
localToSpark: org.apache.spark.sql.DataFrame = [value: string]
localToSpark.collect
res1: Array[org.apache.spark.sql.Row] = Array([[Wed Dec 20 06:18:02 UTC 2017] INFO: ], [*********************************************************************************************************], [ ], [[Wed Dec 20 06:18:02 UTC 2017] INFO: Diagnostic log for xx_user.]
谁能告诉我此时我做错了什么:
val inStream = fs.open(new Path(fileNames(2))) 我得到了错误。
【问题讨论】:
-
如果你有 Spark,为什么不直接将文件从那里写入 HDFS?