最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

SparkStreaming整合Kafka過程詳解

 更新時間:2023年01月27日 10:48:24   作者:健鑫.  
這篇文章主要介紹了SparkStreaming整合Kafka過程,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)吧

Spark Streaming連接kafka 的兩種方式

Receiver based Approah

  • KafkaUtils.createDstream基于接收器方式,消費(fèi)Kafka數(shù)據(jù),已淘汰
  • Receiver作為Task運(yùn)行在Executor等待數(shù)據(jù),一個Receiver效率低,需要開啟多個,再手動合并數(shù)據(jù),很麻煩
  • Receiver掛了,可能丟失數(shù)據(jù),需要開啟WAL(預(yù)寫日志)保證數(shù)據(jù)安全,效率低
  • 通過Zookeeper來連接kafka,offset存儲再zookeeper中
  • spark消費(fèi)的時候?yàn)榱吮WC數(shù)據(jù)不丟也會保存一份offset,可能出現(xiàn)數(shù)據(jù)不一致

Direct Approach

  • KafkaUtils.createDirectStream直連方式,streaming中每個批次的job直接調(diào)用Simple Consumer API獲取對應(yīng)Topic數(shù)據(jù)
  • Direct方式直接連接kafka分區(qū)獲取數(shù)據(jù),提高了并行能力
  • Direct方式調(diào)用kafka低階API,offset自己存儲和維護(hù),默認(rèn)由spark維護(hù)在checkpoint中
  • offset也可以自己手動維護(hù),保存在mysql/redis中
// 從kafka加載數(shù)據(jù)
val kafkaParams = Map[String, Object](
  "bootstrap.servers" -> "hadoop102:9092",//kafka集群地址
  "key.deserializer" -> classOf[StringDeserializer],//key的反序列化規(guī)則
  "value.deserializer" -> classOf[StringDeserializer],//value的反序列化規(guī)則
  "group.id" -> "sparkdemo",//消費(fèi)者組名稱
  //earliest:表示如果有offset記錄從offset記錄開始消費(fèi),如果沒有從最早的消息開始消費(fèi)
  //latest:表示如果有offset記錄從offset記錄開始消費(fèi),如果沒有從最后/最新的消息開始消費(fèi)
  //none:表示如果有offset記錄從offset記錄開始消費(fèi),如果沒有就報(bào)錯
  "auto.offset.reset" -> "latest",
  "auto.commit.interval.ms"->"1000",//自動提交的時間間隔
  "enable.auto.commit" -> (true: java.lang.Boolean)//是否自動提交
)
val topics = Array("spark_kafka")//要訂閱的主題
//使用工具類從Kafka中消費(fèi)消息
val kafkaDS: InputDStream[ConsumerRecord[String, String]] = KafkaUtils.createDirectStream[String, String](
  ssc,
  LocationStrategies.PreferConsistent, //位置策略,使用源碼中推薦的
  ConsumerStrategies.Subscribe[String, String](topics, kafkaParams) //消費(fèi)策略,使用源碼中推薦的
)

代碼展示

自動提交偏移量

object kafka_Demo01 {
  def main(args: Array[String]): Unit = {
    val conf = new SparkConf().setMaster("local[*]").setAppName("kafka_Demo01")
    val sc = new SparkContext(conf)
    val ssc = new StreamingContext(sc, Seconds(5))
    ssc.checkpoint("data/ckp")
    // 從kafka加載數(shù)據(jù)
    val kafkaParams = Map[String, Object](
      "bootstrap.servers" -> "hadoop102:9092",//kafka集群地址
      "key.deserializer" -> classOf[StringDeserializer],//key的反序列化規(guī)則
      "value.deserializer" -> classOf[StringDeserializer],//value的反序列化規(guī)則
      "group.id" -> "sparkdemo",//消費(fèi)者組名稱
      //earliest:表示如果有offset記錄從offset記錄開始消費(fèi),如果沒有從最早的消息開始消費(fèi)
      //latest:表示如果有offset記錄從offset記錄開始消費(fèi),如果沒有從最后/最新的消息開始消費(fèi)
      //none:表示如果有offset記錄從offset記錄開始消費(fèi),如果沒有就報(bào)錯
      "auto.offset.reset" -> "latest",
      "auto.commit.interval.ms"->"1000",//自動提交的時間間隔
      "enable.auto.commit" -> (true: java.lang.Boolean)//是否自動提交
    )
    val topics = Array("spark_kafka")//要訂閱的主題
    //使用工具類從Kafka中消費(fèi)消息
    val kafkaDS: InputDStream[ConsumerRecord[String, String]] = KafkaUtils.createDirectStream[String, String](
      ssc,
      LocationStrategies.PreferConsistent, //位置策略,使用源碼中推薦的
      ConsumerStrategies.Subscribe[String, String](topics, kafkaParams) //消費(fèi)策略,使用源碼中推薦的
    )
    // 處理消息
    val infoDS = kafkaDS.map(record => {
      val topic = record.topic()
      val partition = record.partition()
      val offset = record.offset()
      val key = record.key()
      val value = record.value()
      val info: String = s"""topic:${topic}, partition:${partition}, offset:${offset}, key:${key}, value:${value}"""
      info
    })
    // 輸出
    infoDS.print()
    ssc.start()
    ssc.awaitTermination()
    ssc.stop(true, true)
  }
}

手動提交

提交代碼

// 處理消息
//注意提交的時機(jī):應(yīng)該是消費(fèi)完一小批就該提交一次offset,而在DStream一小批的體現(xiàn)是RDD
kafkaDS.foreachRDD(rdd => {
  rdd.foreach(record => {
    val topic = record.topic()
    val partition = record.partition()
    val offset = record.offset()
    val key = record.key()
    val value = record.value()
    val info: String = s"""topic:${topic}, partition:${partition}, offset:${offset}, key:${key}, value:${value}"""
    info
    println("消費(fèi)" + info)
  })
  //獲取rdd中offset相關(guān)的信息:offsetRanges里面就包含了該批次各個分區(qū)的offset信息
  val offsetRanges: Array[OffsetRange] = rdd.asInstanceOf[HasOffsetRanges].offsetRanges
  //提交
  kafkaDS.asInstanceOf[CanCommitOffsets].commitAsync(offsetRanges)
  println("當(dāng)前批次的數(shù)據(jù)已消費(fèi)并手動提交")
})

完整代碼

object kafka_Demo02 {
  def main(args: Array[String]): Unit = {
    val conf = new SparkConf().setMaster("local[*]").setAppName("kafka_Demo01")
    val sc = new SparkContext(conf)
    val ssc = new StreamingContext(sc, Seconds(5))
    ssc.checkpoint("data/ckp")
    // 從kafka加載數(shù)據(jù)
    val kafkaParams = Map[String, Object](
      "bootstrap.servers" -> "hadoop102:9092",//kafka集群地址
      "key.deserializer" -> classOf[StringDeserializer],//key的反序列化規(guī)則
      "value.deserializer" -> classOf[StringDeserializer],//value的反序列化規(guī)則
      "group.id" -> "sparkdemo",//消費(fèi)者組名稱
      //earliest:表示如果有offset記錄從offset記錄開始消費(fèi),如果沒有從最早的消息開始消費(fèi)
      //latest:表示如果有offset記錄從offset記錄開始消費(fèi),如果沒有從最后/最新的消息開始消費(fèi)
      //none:表示如果有offset記錄從offset記錄開始消費(fèi),如果沒有就報(bào)錯
      "auto.offset.reset" -> "latest",
//      "auto.commit.interval.ms"->"1000",//自動提交的時間間隔
      "enable.auto.commit" -> (false: java.lang.Boolean)//是否自動提交
    )
    val topics = Array("spark_kafka")//要訂閱的主題
    //使用工具類從Kafka中消費(fèi)消息
    val kafkaDS: InputDStream[ConsumerRecord[String, String]] = KafkaUtils.createDirectStream[String, String](
      ssc,
      LocationStrategies.PreferConsistent, //位置策略,使用源碼中推薦的
      ConsumerStrategies.Subscribe[String, String](topics, kafkaParams) //消費(fèi)策略,使用源碼中推薦的
    )
    // 處理消息
    //注意提交的時機(jī):應(yīng)該是消費(fèi)完一小批就該提交一次offset,而在DStream一小批的體現(xiàn)是RDD
    kafkaDS.foreachRDD(rdd => {
      rdd.foreach(record => {
        val topic = record.topic()
        val partition = record.partition()
        val offset = record.offset()
        val key = record.key()
        val value = record.value()
        val info: String = s"""topic:${topic}, partition:${partition}, offset:${offset}, key:${key}, value:${value}"""
        info
        println("消費(fèi)" + info)
      })
      //獲取rdd中offset相關(guān)的信息:offsetRanges里面就包含了該批次各個分區(qū)的offset信息
      val offsetRanges: Array[OffsetRange] = rdd.asInstanceOf[HasOffsetRanges].offsetRanges
      //提交
      kafkaDS.asInstanceOf[CanCommitOffsets].commitAsync(offsetRanges)
      println("當(dāng)前批次的數(shù)據(jù)已消費(fèi)并手動提交")
    })
    // 輸出
    kafkaDS.print()
    ssc.start()
    ssc.awaitTermination()
    ssc.stop(true, true)
  }
}

到此這篇關(guān)于SparkStreaming整合Kafka過程詳解的文章就介紹到這了,更多相關(guān)SparkStreaming整合Kafka內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Java的Lombok之@Builder使用總結(jié)

    Java的Lombok之@Builder使用總結(jié)

    這篇文章主要介紹了Java的Lombok之@Builder使用總結(jié),當(dāng)不使用@Builder注解到類上,創(chuàng)建T1的有參構(gòu)造函數(shù),入?yún)⒉粌H包括T1中所有的參數(shù),還包括T中所有的參數(shù),T2的屬性由T1在有參構(gòu)造函數(shù)中通過調(diào)用父類構(gòu)造器的方式賦初值,需要的朋友可以參考下
    2023-12-12
  • Java實(shí)戰(zhàn)之在線租房系統(tǒng)的實(shí)現(xiàn)

    Java實(shí)戰(zhàn)之在線租房系統(tǒng)的實(shí)現(xiàn)

    這篇文章主要介紹了利用Java實(shí)現(xiàn)的在線租房系統(tǒng),文中用到了SpringBoot、Redis、MySQL、Vue等技術(shù),文中示例代碼講解詳細(xì),需要的可以參考一下
    2022-02-02
  • Java中Controller、Service、Dao/Mapper層的區(qū)別與用法

    Java中Controller、Service、Dao/Mapper層的區(qū)別與用法

    在Java開發(fā)中,通常會采用三層架構(gòu)(或稱MVC架構(gòu))來劃分程序的職責(zé)和功能,分別是Controller層、Service層、Dao/Mapper層,本文將詳細(xì)給大家介紹了三層的區(qū)別和用法,需要的朋友可以參考下
    2023-05-05
  • Spring詳細(xì)講解FactoryBean接口的使用

    Spring詳細(xì)講解FactoryBean接口的使用

    這篇文章主要為大家介紹了Spring容器FactoryBean工廠實(shí)例詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2022-06-06
  • Java實(shí)現(xiàn)QQ第三方登錄的示例代碼

    Java實(shí)現(xiàn)QQ第三方登錄的示例代碼

    這篇文章主要介紹了Java實(shí)現(xiàn)QQ第三方登錄的示例代碼,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2019-11-11
  • mybatis中使用大于小于等于的正確方法

    mybatis中使用大于小于等于的正確方法

    在mybatis中sql是寫在xml映射文件中的,如果sql中有一些特殊字符的話,在解析xml文件的時候就會被轉(zhuǎn)義,下面我們就一起來看一下大于小于等于是怎么轉(zhuǎn)義的
    2021-04-04
  • Java生成exe可執(zhí)行文件

    Java生成exe可執(zhí)行文件

    本文主要介紹了Java如何生成exe可執(zhí)行文件,想了解更多的小伙伴可以借鑒閱讀這篇文章
    2023-03-03
  • Spring?AI集成DeepSeek實(shí)現(xiàn)流式輸出的操作方法

    Spring?AI集成DeepSeek實(shí)現(xiàn)流式輸出的操作方法

    本文介紹了如何在SpringBoot中使用Sse(Server-SentEvents)技術(shù)實(shí)現(xiàn)流式輸出,后端使用SpringMVC中的SseEmitter對象,前端使用EventSource對象監(jiān)聽SSE接口并展示數(shù)據(jù)流,通過這種方式可以提升用戶體驗(yàn),避免大模型響應(yīng)速度慢的問題,感興趣的朋友一起看看吧
    2025-03-03
  • springboot結(jié)合vue實(shí)現(xiàn)增刪改查及分頁查詢

    springboot結(jié)合vue實(shí)現(xiàn)增刪改查及分頁查詢

    本文主要介紹了springboot結(jié)合vue實(shí)現(xiàn)增刪改查及分頁查詢,文中通過示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2021-09-09
  • Spring MVC 更靈活的控制 json 返回問題(自定義過濾字段)

    Spring MVC 更靈活的控制 json 返回問題(自定義過濾字段)

    本篇文章主要介紹了Spring MVC 更靈活的控制 json 返回問題(自定義過濾字段),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下。
    2017-02-02

最新評論

德钦县| 河东区| 孟州市| 虎林市| 昌黎县| 桦南县| 巴中市| 昌平区| 珠海市| 盱眙县| 福清市| 湘潭市| 延川县| 康马县| 类乌齐县| 兖州市| 曲阜市| 贵溪市| 无为县| 衢州市| 肇东市| 平顺县| 喀喇| 疏勒县| 卢龙县| 岗巴县| 龙川县| 灵璧县| 锡林郭勒盟| 赤水市| 铜梁县| 永安市| 策勒县| 榕江县| 青浦区| 鄂温| 博客| 张掖市| 恭城| 自贡市| 怀仁县|