【发布时间】:2015-10-28 06:27:47
【问题描述】:
我正在运行一个简单的 WordCount 程序。 Spark Streaming 正在 HDFS 中的目录中查看新文件,并应在它们进入时对其进行处理。
我开始我的流媒体工作,我将一堆小文件添加到一个 tmp HDFS 目录,然后我将这些文件移动到监视的 HDFS 目录(全部使用简单的 shell 命令,- MV)。但我的流式传输作业没有将这些文件识别为新文件,因此没有处理它们(我检查了文件是否移动良好)。
目前我正在使用 textFileStream,但我愿意使用 fileStream。我使用的是 1.3.1 或 1.4.0 Spark 版本。 我想提一下,使用 1.0.x 版本的 spark,一切都很好(它检测到新的 -moved- 文件)!
代码是:
//files are moved from /user/share/jobs-data/gstream/tmp to /user/share/jobs-data/gstream/streams, both directories are on HDFS.
val sparkConf = new SparkConf().setAppName(this.getClass().getName())
sparkConf.setMaster(master)
val sc = new SparkContext(sparkConf)
val ssc = new StreamingContext(sc, Milliseconds(1000))
val data = ssc.textFileStream(args(1)) //args(1) == /user/share/jobs-data/gstream/streams
val words = data.flatMap(.split(" "))
val wordCounts = words.map(word => (word, 1)).reduceByKey( + _)
wordCounts.print()
ssc.start()
ssc.awaitTermination()
谁能给点意见,谢谢?
【问题讨论】:
-
请发布您的代码。否则我们只是在猜测你做了什么
-
其实代码就是一个简单的字数统计:val data = ssc.textFileStream(args(1)); val words = data.flatMap(.split(" ")); val wordCounts = words.map(word => (word, 1)).reduceByKey( + _); wordCounts.print();
-
请将代码编辑到您的问题中。您的监视目录和 tmp 目录是否在同一个文件系统上? “必须通过将文件从同一文件系统中的另一个位置“移动”它们来将文件写入受监视的目录”。也许您的 tmp 有所不同。
-
好的,我编辑了我的问题并添加了一些其他细节
-
我认为那里有一个
ssc.start?你能显示设置ssc的代码吗?
标签: scala apache-spark spark-streaming