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

SparkStreaming-Kafka通過(guò)指定偏移量獲取數(shù)據(jù)實(shí)現(xiàn)

 更新時(shí)間:2023年06月20日 11:55:04   作者:spark打醬油  
這篇文章主要為大家介紹了SparkStreaming-Kafka通過(guò)指定偏移量獲取數(shù)據(jù),有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪

SparkStreaming-Kafka通過(guò)指定偏移量獲取數(shù)據(jù)

1.數(shù)據(jù)源

'310999003001', '3109990030010220140820141230292','00000000','','2017-08-20 14:09:35','0',255,'SN', 0.00,'4','','310999','310999003001','02','','','2','','','2017-08-20 14:12:30','2017-08-20 14:16:13',0,0,'2017-08-21 18:50:05','','',' '
'310999003102', '3109990031020220140820141230266','粵BT96V3','','2017-08-20 14:09:35','0',21,'NS', 0.00,'2','','310999','310999003102','02','','','2','','','2017-08-20 14:12:30','2017-08-20 14:16:13',0,0,'2017-08-21 18:50:05','','',' '

2.生產(chǎn)者

import java.util.Properties
import com.google.gson.{Gson, GsonBuilder}
import org.apache.kafka.clients.producer.{KafkaProducer, ProducerRecord}
import org.apache.spark.rdd.RDD
import org.apache.spark.{SparkConf, SparkContext}
/**
  * Date 2022/11/8 9:49
  */
object KafkaEventProducer {
  def main(args: Array[String]): Unit = {
    val conf: SparkConf = new SparkConf().setAppName("KafkaEventProducer").setMaster("local[*]")
    val sc = new SparkContext(conf)
    sc.setLogLevel("ERROR")
    val topic = "ly_test"
    val props = new Properties()
    props.put("bootstrap.servers","node01:9092,node02:9092,node03:9092")
    props.put("key.serializer","org.apache.kafka.common.serialization.StringSerializer")
    props.put("value.serializer","org.apache.kafka.common.serialization.StringSerializer")
    props.put("acks","all")
//    props.put("security.protocol","SASL_PLAINTEXT")
//    props.put("sasl.mechanism","PLAIN")
//    props.put("sasl.jaas.config","org.apache.kafka.common.security.plain.PlainLoginModule required username='sjf' password='123sjf';")
    val kafkaProducer = new KafkaProducer[String,String](props)
    val srcRDD: RDD[String] = sc.textFile("file:///F:\\work\\sun\\lywork\\hbaseoper\\datas\\kafkaproducerdata.txt")
    val records: Array[Array[String]] = srcRDD.filter(!_.startsWith(";")).map(_.split(",")).collect()
    //對(duì)數(shù)據(jù)進(jìn)行預(yù)處理形成json形式
    for(record<-records){
      val trafficInfo = new TrafficInfo(record(0),record(2),record(4),record(6),record(13))
      // 不能用new Gson()   會(huì)出現(xiàn) \u0027
      // val trafficInfoJson: String = new Gson().toJson(trafficInfo)
      //使用Gson gson = new Gson(),進(jìn)行對(duì)象轉(zhuǎn)化json格式時(shí),單引號(hào)會(huì)被轉(zhuǎn)換成u0027代碼。使用以下方法進(jìn)行替換
      val gson: Gson = new GsonBuilder().disableHtmlEscaping().create()
      val trafficInfoJson: String = gson.toJson(trafficInfo)
      kafkaProducer.send(new ProducerRecord[String,String](topic,trafficInfoJson))
      println("Message Sent:"+trafficInfoJson)
      Thread.sleep(2000)
    }
    sc.stop()
    kafkaProducer.flush()
    kafkaProducer.close()
  }
  //相機(jī)編號(hào)
  //車牌號(hào)
  //時(shí)間
  //速度
  //車道編號(hào)
  case class TrafficInfo(camer_id:String,car_id:String,event_time:String,car_speed:String,car_code:String)
}

3.消費(fèi)者獲取指定偏移量

import java.text.SimpleDateFormat
import java.util
import java.util.Date
import org.apache.kafka.clients.consumer.ConsumerRecord
import org.apache.kafka.common.record.TimestampType
import org.apache.kafka.common.serialization.StringDeserializer
import org.apache.spark.rdd.RDD
import org.apache.spark.streaming.dstream.{DStream, InputDStream}
import org.apache.spark.streaming.kafka010.{ConsumerStrategies, KafkaUtils, LocationStrategies, OffsetRange}
import org.apache.spark.streaming.{Seconds, StreamingContext}
import org.apache.spark.{SparkConf, SparkContext}
/**
  * Date 2022/11/5 16:38
  */
/**
  * 通過(guò)偏移量獲取數(shù)據(jù)
  */
object AttainDataFromOffset {
  def main(args: Array[String]): Unit = {
    val conf: SparkConf = new SparkConf().setAppName("AttainDataFromOffset").setMaster("local[*]")
    val sc = new SparkContext(conf)
    sc.setLogLevel("ERROR")
    val ssc = new StreamingContext(sc,Seconds(5))
    val kafkaParams: Map[String, Object] = Map[String, Object](
      "bootstrap.servers" -> "node01:9092,node02:9092,node03:9092",
      "key.deserializer" -> classOf[StringDeserializer],
      "value.deserializer" -> classOf[StringDeserializer],
      "group.id" -> "ly",
      "auto.offset.reset" -> "earliest",
      "enable.auto.commit" -> (false: java.lang.Boolean)
      // kafka 帶有賬號(hào)密碼sasl協(xié)議的認(rèn)證
//      "security.protocol" -> "SASL_PLAINTEXT",
//      "sasl.mechanism" -> "PLAIN",
//      "sasl.jaas.config" -> "org.apache.kafka.common.security.plain.PlainLoginModule required username='sjf' password='123sjf';"
    )
    val topics = Array("ly_test")
    val stream: InputDStream[ConsumerRecord[String, String]] = KafkaUtils.createDirectStream[String, String](
      ssc,
      LocationStrategies.PreferConsistent,
      ConsumerStrategies.Subscribe[String,String](topics, kafkaParams)
    )
    val res: DStream[(String, String, Int, Long)] = stream.map(recoed => {
      val key: String = recoed.key()
      val value: String = recoed.value()
      val partionId: Int = recoed.partition()
      val offset: Long = recoed.offset()
      (key, value, partionId, offset)
      //println(key+"\t"+value+"\t"+partionId+"\t"+offset)
    })
    // 指定偏移量
    val offsetRanges = Array(
      // topic, partition, inclusive starting offset, exclusive ending offset
      OffsetRange("lawyee_test", 0, 1L, 10L)
    )
    // 獲取指定偏移量的數(shù)據(jù)
    import scala.collection.JavaConverters._
    val jkafkaParams: util.Map[String, Object] = kafkaParams.asJava
    val offsetRDD: RDD[ConsumerRecord[String, String]] = KafkaUtils.createRDD[String,String](
      sc,
      jkafkaParams,
      offsetRanges,
      LocationStrategies.PreferConsistent
    )
    val resRDD: RDD[(String, String, Int, Long,String,TimestampType)] = offsetRDD.map(recoed => {
      val key: String = recoed.key()
      val value: String = recoed.value()
      val partionId: Int = recoed.partition()
      val offset: Long = recoed.offset()
      var time: Long = recoed.timestamp()
      val timeStr = timeStampToDate(time)
      val timestampType: TimestampType = recoed.timestampType()
      (key, value, partionId, offset,timeStr,timestampType)
      //println(key+"\t"+value+"\t"+partionId+"\t"+offset)
    })
    resRDD.foreach(println(_))
    res.print()
    ssc.start()
    ssc.awaitTermination()
  }
  // 時(shí)間格式時(shí)間 轉(zhuǎn)換為字符串時(shí)間
  def dateToString(date:Date): String ={
    val simpleDateFormat = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss")
    val strDate: String = simpleDateFormat.format(date)
    strDate
  }
  // 字符串時(shí)間轉(zhuǎn)換為時(shí)間格式時(shí)間
  def strToDate(str:String):Date = {
    val simpleDateFormat = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss")
    val date: Date = simpleDateFormat.parse(str)
    date
  }
  // 時(shí)間戳轉(zhuǎn)化為字符串時(shí)間
  def timeStampToDate(timeStamp:Long): String ={
    val date = new Date(timeStamp)
    val strDate: String = dateToString(date)
    strDate
  }
  //字符串時(shí)間轉(zhuǎn)化為時(shí)間戳
  def dateToTimeStamp(strDate:String): Long ={
    val date: Date = strToDate(strDate)
    val timeStamp: Long = date.getTime
    timeStamp
  }
}

以上就是SparkStreaming-Kafka通過(guò)指定偏移量獲取數(shù)據(jù)的詳細(xì)內(nèi)容,更多關(guān)于SparkStreaming Kafka獲取數(shù)據(jù)的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • java線程優(yōu)先級(jí)原理詳解

    java線程優(yōu)先級(jí)原理詳解

    這篇文章主要介紹了java線程優(yōu)先級(jí)原理詳解,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2019-10-10
  • SpringBoot中Elasticsearch的連接配置原理與使用詳解

    SpringBoot中Elasticsearch的連接配置原理與使用詳解

    Elasticsearch是一種開(kāi)源的分布式搜索和數(shù)據(jù)分析引擎,它可用于全文搜索、結(jié)構(gòu)化搜索、分析等應(yīng)用場(chǎng)景,本文主要介紹了SpringBoot中Elasticsearch的連接配置原理與使用詳解,感興趣的可以了解一下
    2023-09-09
  • Spring框架+jdbcTemplate實(shí)現(xiàn)增刪改查功能

    Spring框架+jdbcTemplate實(shí)現(xiàn)增刪改查功能

    這篇文章主要介紹了Spring框架+jdbcTemplate實(shí)現(xiàn)增刪改查功能,本文通過(guò)實(shí)例代碼給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2021-09-09
  • Java實(shí)現(xiàn)排隊(duì)論的原理

    Java實(shí)現(xiàn)排隊(duì)論的原理

    這篇文章主要為大家詳細(xì)介紹了Java實(shí)現(xiàn)排隊(duì)論的原理,對(duì)排隊(duì)論感興趣的小伙伴們可以參考一下
    2016-02-02
  • SpringIOC框架的簡(jiǎn)單實(shí)現(xiàn)步驟

    SpringIOC框架的簡(jiǎn)單實(shí)現(xiàn)步驟

    這篇文章主要介紹了SpringIOC框架簡(jiǎn)單實(shí)現(xiàn)步驟,幫助大家更好的理解和學(xué)習(xí)使用Spring,感興趣的朋友可以了解下
    2021-05-05
  • Java使用freemarker實(shí)現(xiàn)word下載方式

    Java使用freemarker實(shí)現(xiàn)word下載方式

    文章介紹了如何使用FreeMarker實(shí)現(xiàn)Word文件下載,包括引用依賴、創(chuàng)建Word模板、將Word文件存為XML格式、更改后綴為FTL模板、處理圖片和代碼實(shí)現(xiàn)
    2025-02-02
  • 詳解Spring中Bean的加載的方法

    詳解Spring中Bean的加載的方法

    本篇文章主要介紹了Spring中Bean的加載的方法,小編覺(jué)得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧
    2017-04-04
  • HashMap vs TreeMap vs Hashtable vs LinkedHashMap

    HashMap vs TreeMap vs Hashtable vs LinkedHashMap

    這篇文章主要介紹了HashMap vs TreeMap vs Hashtable vs LinkedHashMap的相關(guān)知識(shí),非常不錯(cuò),具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2019-07-07
  • iBatis習(xí)慣用的16條SQL語(yǔ)句

    iBatis習(xí)慣用的16條SQL語(yǔ)句

    iBatis 是apache 的一個(gè)開(kāi)源項(xiàng)目,一個(gè)O/R Mapping 解決方案,iBatis 最大的特點(diǎn)就是小巧,上手很快.這篇文章主要介紹了iBatis習(xí)慣用的16條SQL語(yǔ)句的相關(guān)資料,需要的朋友可以參考下
    2016-10-10
  • IDEA中的pom.xml文件無(wú)法識(shí)別問(wèn)題及解決

    IDEA中的pom.xml文件無(wú)法識(shí)別問(wèn)題及解決

    這篇文章主要介紹了IDEA中的pom.xml文件無(wú)法識(shí)別問(wèn)題及解決方案,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2023-11-11

最新評(píng)論

泰宁县| 甘南县| 湖北省| 庆云县| 遂昌县| 察雅县| 勃利县| 祁东县| 丹巴县| 金秀| 丰都县| 博客| 尚义县| 奉节县| 嵩明县| 象山县| 弋阳县| 海南省| 洞口县| 通渭县| 固安县| 海宁市| 闽清县| 阳原县| 离岛区| 盈江县| 交城县| 巫山县| 建阳市| 姚安县| 扎鲁特旗| 康保县| 平果县| 莱州市| 南平市| 苏尼特左旗| 象山县| 清水县| 清水县| 鄂尔多斯市| 忻州市|