直接上代碼吧
說下測試思路:
該代碼監(jiān)控的/tmp/sparkStream/目錄;
首先,創(chuàng)建該目錄mkdir -p /tmp/sparkStream;
然后,運(yùn)行spark程序;
最后,向監(jiān)控目錄/tmp/sparkStream/添加數(shù)據(jù)文件;
觀察spark程序運(yùn)行效果。
sparkStreaming import org.apache.log4j.{LoggerLevel} import org.apache.spark.SparkConf import org.apache.spark.streaming.{SecondsStreamingContext} import org.apache.spark.streaming.StreamingContext._ object HdfsWordCount { def main(args: Array[]){ Logger.getLogger("org.apache.spark").setLevel(Level.WARN) Logger.getLogger("org.apache.eclipse.jetty.server").setLevel(Level.OFF) sparkConf = SparkConf().setAppName().setMaster() ssc = StreamingContext(sparkConf()) lines = ssc.textFileStream() words = lines.flatMap(_.split()) wordCounts = words.map(x=>(x)).reduceByKey(_+_) wordCounts.print() ssc.start() ssc.awaitTermination() } }
另外有需要云服務(wù)器可以了解下創(chuàng)新互聯(lián)scvps.cn,海內(nèi)外云服務(wù)器15元起步,三天無理由+7*72小時售后在線,公司持有idc許可證,提供“云服務(wù)器、裸金屬服務(wù)器、高防服務(wù)器、香港服務(wù)器、美國服務(wù)器、虛擬主機(jī)、免備案服務(wù)器”等云主機(jī)租用服務(wù)以及企業(yè)上云的綜合解決方案,具有“安全穩(wěn)定、簡單易用、服務(wù)可用性高、性價比高”等特點(diǎn)與優(yōu)勢,專為企業(yè)上云打造定制,能夠滿足用戶豐富、多元化的應(yīng)用場景需求。