Spark Streaming

来源:http://www.cnblogs.com/one--way/archive/2016/09/16/5877552.html
-Advertisement-
Play Games

Spark Streaming Spark Streaming 是Spark為了用戶實現流式計算的模型。 數據源包括Kafka,Flume,HDFS等。 DStream 離散化流(discretized stream), Spark Streaming 使用DStream作為抽象表示。是隨時間推移而 ...


Spark Streaming 

Spark Streaming 是Spark為了用戶實現流式計算的模型。

數據源包括Kafka,Flume,HDFS等。

DStream 離散化流(discretized stream), Spark Streaming 使用DStream作為抽象表示。是隨時間推移而收到的數據的序列。DStream內部的數據都是RDD形式存儲, DStream是由這些RDD所組成的離散序列。

 

編寫Streaming步驟:

1.創建StreamingContext

// Create a local StreamingContext with two working thread and batch interval of 5 second.
// The master requires 2 cores to prevent from a starvation scenario.
val conf = new SparkConf().setMaster("local[2]").setAppName("NetworkWordCount")
val ssc = new StreamingContext(conf, Seconds(5))

創建本地化StreamingContext, 需要至少2個工作線程。一個是receiver,一個是計算節點。

2.定義輸入源,創建輸入DStream

// Create a DStream that will connect to hostname:port, like localhost:9999
val lines = ssc.socketTextStream("node1", 9999)

3.定義流的計算過程,使用transformation和output operation DStream

// Split each line into words
val words = lines.flatMap(_.split(" "))

// Count each word in each batch
val pairs = words.map(word => (word, 1))
val wordCounts = pairs.reduceByKey(_ + _)

// Print the first ten elements of each RDD generated in this DStream to the console
wordCounts.print()

4.開始接收數據及處理數據,使用streamingContext.start()

ssc.start()             // Start the computation

5.等待批處理被終止,使用streamingContext.awaitTermination()

ssc.awaitTermination()  // Wait for the computation to terminate

6.可以手工停止批處理,使用streamingContext.stop()

 

數據源

數據源分為兩種

1.基本源

text,HDFS等

2.高級源

Flume,Kafka等

 

DStream支持兩種操作

一、轉化操作(transformation)

無狀態轉化(stateless):每個批次的處理不依賴於之前批次的數據

 

有狀態轉化(stateful):跨時間區間跟蹤數據的操作;一些先前批次的數據被用來在新的批次中參與運算。

  • 滑動視窗:
  • 追蹤狀態變化:updateStateByKey()

 

transform函數

transform操作允許任意RDD-to-RDD函數被應用在一個DStream中.比如在DStream中的RDD可以和DStream外部的RDD進行Join操作。通常用來過濾垃圾郵件等操作。

不屬於DStream的RDD,在每個批間隔都被調用;允許你做隨時間變化的RDD操作。RDD操作,partitions的數量,廣播變數的值,可以變化在不同的批次中。

例子:

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

/**
  * Created by Edward on 2016/9/16.
  */
object TransformOperation {

  def main(args: Array[String]) {
    val conf: SparkConf = new SparkConf().setMaster("local[2]").setAppName("TransformOperation")
    val ssc = new StreamingContext(conf, Seconds(5))
    val textStream: ReceiverInputDStream[String] = ssc.socketTextStream("node1",9999)

    //黑名單過濾功能,可以將數據存到redis或者資料庫,每批次間隔都會重新取數據並參與運算,保證數據可以動態載入進來。
    val blackList=Array(Tuple2("Tom", true))
    val listRDD  = ssc.sparkContext.parallelize(blackList).persist() //創建RDD

    val map  = textStream.map(x=>(x.split(" ")(1),x))
    
    //通過transform將DStream中的RDD進行過濾操作
    val dStream = map.transform(rdd =>{
      //listRDD.collect()
      //println(listRDD.collect.length)
      
      //通過RDD的左鏈接及過濾函數,對數據進行處理,生成新的RDD
      rdd.leftOuterJoin(listRDD).filter(x =>{ //使用transform操作DStream中的rdd  rdd左鏈接listRDD, 併進行過濾操作
        if(!x._2._2.isEmpty && x._2._2.get)// if(x._2._2.getOrElse(false)) //如果沒取到值則結果為false
          false
        else{
          true
        }
      })
    })

    dStream.print()

    ssc.start()
    ssc.awaitTermination()

  }
}

 

視窗函數

 

 

二、輸出操作(output operation)

 

 


您的分享是我們最大的動力!

-Advertisement-
Play Games
更多相關文章
  • 【等待事件】等待事件系列(3+4)--System IO(控制文件)+日誌類等待 1 BLOG文檔結構圖 2 前言部分 2.1 導讀和註意事項 各位技術愛好者,看完本文後,你可以掌握如下的技能,也可以學到一些其它你所不知道的知識,~O(∩_∩)O~: ① 控制文件類等待 ② 日誌類等待 2.2 相關... ...
  • 每一個資料庫連接可以包括多個資料庫文件,一個主資料庫文件和attached的幾個資料庫文件。每一個資料庫文件都有自己的B-tree和pager。資料庫連接(connection)和事務(transaction)的關係1. 每一個資料庫連接都必須在事務下進行操作2. 每一個資料庫連接不同同時打開多個... ...
  • https://blueprints.launchpad.net/myconnpy/+spec/sp-multi-resultsets Calling a stored procedure can produce multiple result sets. They should be retrie ...
  • 12down votefavorite 8 8 http://stackoverflow.com/questions/31748278/how-do-you-install-mysql-connector-python-development-version-through-pip/34027037 ...
  • SQL 是用於訪問和處理資料庫的標準的電腦語言。 什麼是 SQL? SQL 指結構化查詢語言 SQL 使我們有能力訪問資料庫 SQL 是一種 ANSI 的標準電腦語言 編者註:ANSI,美國國家標準化組織 什麼是 SQL? SQL 指結構化查詢語言 SQL 使我們有能力訪問資料庫 SQL 是一種 ...
  • 之前一直沒用到這個函數,因為一般的情況下分類比較少的情況下,比如 : 男 : 0 女 : 1 我們一般會這麼設計資料庫,前臺顯示的時候一般會用select 或者 redio value = "0" 或=“1”顯示這樣確實是比較方便 其實用decode函數直接就可以使用了 select decode( ...
  • 一、登錄 打開終端,輸入/usr/local/mysql/bin/mysql -u root -p 初次進入mysql,密碼為空。當出現mysql>提示符時,表示你已經進入mysql中。鍵入exit退出mysql。 二、更改Mysqlroot用戶密碼 更改mysql root 用戶密碼,在終端輸入/ ...
  • oracle使用java source調用外部程式 需求 Oracle調用第三方外部程式。Oracle使用sqluldr2快速導出大批量數據,然後用winrar壓縮後發送郵件。 原碼 java source create or replace and compile java source name... ...
一周排行
    -Advertisement-
    Play Games
  • 示例項目結構 在 Visual Studio 中創建一個 WinForms 應用程式後,項目結構如下所示: MyWinFormsApp/ │ ├───Properties/ │ └───Settings.settings │ ├───bin/ │ ├───Debug/ │ └───Release/ ...
  • [STAThread] 特性用於需要與 COM 組件交互的應用程式,尤其是依賴單線程模型(如 Windows Forms 應用程式)的組件。在 STA 模式下,線程擁有自己的消息迴圈,這對於處理用戶界面和某些 COM 組件是必要的。 [STAThread] static void Main(stri ...
  • 在WinForm中使用全局異常捕獲處理 在WinForm應用程式中,全局異常捕獲是確保程式穩定性的關鍵。通過在Program類的Main方法中設置全局異常處理,可以有效地捕獲並處理未預見的異常,從而避免程式崩潰。 註冊全局異常事件 [STAThread] static void Main() { / ...
  • 前言 給大家推薦一款開源的 Winform 控制項庫,可以幫助我們開發更加美觀、漂亮的 WinForm 界面。 項目介紹 SunnyUI.NET 是一個基於 .NET Framework 4.0+、.NET 6、.NET 7 和 .NET 8 的 WinForm 開源控制項庫,同時也提供了工具類庫、擴展 ...
  • 說明 該文章是屬於OverallAuth2.0系列文章,每周更新一篇該系列文章(從0到1完成系統開發)。 該系統文章,我會儘量說的非常詳細,做到不管新手、老手都能看懂。 說明:OverallAuth2.0 是一個簡單、易懂、功能強大的許可權+可視化流程管理系統。 有興趣的朋友,請關註我吧(*^▽^*) ...
  • 一、下載安裝 1.下載git 必須先下載並安裝git,再TortoiseGit下載安裝 git安裝參考教程:https://blog.csdn.net/mukes/article/details/115693833 2.TortoiseGit下載與安裝 TortoiseGit,Git客戶端,32/6 ...
  • 前言 在項目開發過程中,理解數據結構和演算法如同掌握蓋房子的秘訣。演算法不僅能幫助我們編寫高效、優質的代碼,還能解決項目中遇到的各種難題。 給大家推薦一個支持C#的開源免費、新手友好的數據結構與演算法入門教程:Hello演算法。 項目介紹 《Hello Algo》是一本開源免費、新手友好的數據結構與演算法入門 ...
  • 1.生成單個Proto.bat內容 @rem Copyright 2016, Google Inc. @rem All rights reserved. @rem @rem Redistribution and use in source and binary forms, with or with ...
  • 一:背景 1. 講故事 前段時間有位朋友找到我,說他的窗體程式在客戶這邊出現了卡死,讓我幫忙看下怎麼回事?dump也生成了,既然有dump了那就上 windbg 分析吧。 二:WinDbg 分析 1. 為什麼會卡死 窗體程式的卡死,入口門檻很低,後續往下分析就不一定了,不管怎麼說先用 !clrsta ...
  • 前言 人工智慧時代,人臉識別技術已成為安全驗證、身份識別和用戶交互的關鍵工具。 給大家推薦一款.NET 開源提供了強大的人臉識別 API,工具不僅易於集成,還具備高效處理能力。 本文將介紹一款如何利用這些API,為我們的項目添加智能識別的亮點。 項目介紹 GitHub 上擁有 1.2k 星標的 C# ...