代码

sparkstreaming本地读取数据代码块

package com.mydemo

import org.apache.spark.SparkConf
import org.apache.spark.streaming.dstream.DStream
import org.apache.spark.streaming.{Seconds, StreamingContext}

object strDemo {

  def main(args: Array[String]): Unit = {

    //1.初始化Spark配置信息
    val sparkConf = new SparkConf().setMaster("local[*]")
      .setAppName("StreamWordCount")

    //2.初始化SparkStreamingContext,间隔时间
    val ssc = new StreamingContext(sparkConf, Seconds(5))

    //3.监控文件夹
    val dirStream = ssc.textFileStream("file:///C:/User/IdeaProjects/spark01/")

    //4.将数据进行切分
    val wordStreams: DStream[String] = dirStream.flatMap(_.split(","))

    //5.将切割后的数据转换为(数据,1)的二元组格式
    val wordAndOneStreams = wordStreams.map((_, 1))

    //6.通过二元组中相同key的数据求和
    val wordAndCountStreams = wordAndOneStreams.reduceByKey(_ + _)

    //7.打印
    wordAndCountStreams.print()

    //8.启动SparkStreamingContext
    ssc.start()
    ssc.awaitTermination()
  }
}

原因

sparkstreaming读取数据是以流的形式读取的,在HDFS中由于HDFS的上传、下载、移动等操作都是以文件流的形式操作,而在本地中的相同文件夹中复制移动不是以文件流的形式操作的,若是想在本地使用sparkstreaming读取改变的文件夹或者文件可以以以下操作写入数据
写入文件代码块

package com.mydemo

import java.io.FileWriter

import scala.util.Random

object DataCreater {
	//初始化地址值,循环次数,数据
  private val datapath = "C:\\User\\IdeaProjects\\spark01\\score.txt"
  private val max_records = 20
  private val brand = Array("手机", "笔记本", "小龙虾", "卫生纸", "吸尘器",  "苹果", "洗面奶", "保温杯")

  def Creater(): Unit ={
  	
    val rand = new Random()
    val writer: FileWriter = new FileWriter(datapath,true)

    // create age of data
    for(i <- 1 to max_records){
      //电器名称
      var phonePlus = brand(rand.nextInt(10))
      //电器价格
      var price = rand.nextInt(999)+1000
      //数据拼接
      writer.write( phonePlus + "," + price)
      writer.write(System.getProperty("line.separator"))
    }
    writer.flush()
    writer.close()
  }
  def main(args: Array[String]): Unit = {
    Creater()
    System.exit(1)
  }
}

这样的话就好啦

更多推荐