91超碰碰碰碰久久久久久综合_超碰av人澡人澡人澡人澡人掠_国产黄大片在线观看画质优化_txt小说免费全本

溫馨提示×

溫馨提示×

您好,登錄后才能下訂單哦!

密碼登錄×
登錄注冊×
其他方式登錄
點擊 登錄注冊 即表示同意《億速云用戶服務條款》

Spark-Streaming如何處理數據到mysql中

發布時間:2021-12-04 11:33:23 來源:億速云 閱讀:295 作者:iii 欄目:大數據

本篇內容主要講解“Spark-Streaming如何處理數據到mysql中”,感興趣的朋友不妨來看看。本文介紹的方法操作簡單快捷,實用性強。下面就讓小編來帶大家學習“Spark-Streaming如何處理數據到mysql中”吧!

1.說明

數據表如下:

create database test;
use test;
DROP TABLE IF EXISTS car_gps;
CREATE TABLE IF NOT EXISTS car_gps(
deployNum VARCHAR(30) COMMENT '調度編號',
plateNum VARCHAR(10) COMMENT '車牌號',
timeStr VARCHAR(20) COMMENT '時間戳',
lng VARCHAR(20) COMMENT '經度',
lat VARCHAR(20) COMMENT '緯度',
dbtime TIMESTAMP DEFAULT CURRENT_TIMESTAMP COMMENT '數據入庫時間',
PRIMARY KEY(deployNum, plateNum, timeStr))

2.編寫程序

首先引入mysql的驅動

  <dependency>
      <groupId>mysql</groupId>
      <artifactId>mysql-connector-java</artifactId>
      <version>5.1.44</version>
  </dependency>

2.1 jdbc寫入mysql

package com.hoult.Streaming.work

import java.sql.{Connection, DriverManager, PreparedStatement}
import java.util.Properties

import com.hoult.structed.bean.BusInfo
import org.apache.spark.sql.ForeachWriter

class JdbcHelper extends ForeachWriter[BusInfo] {
  var conn: Connection = _
  var statement: PreparedStatement = _
  override def open(partitionId: Long, epochId: Long): Boolean = {
    if (conn == null) {
      conn = JdbcHelper.openConnection
    }
    true
  }

  override def process(value: BusInfo): Unit = {
    //把數據寫入mysql表中
    val arr: Array[String] = value.lglat.split("_")
    val sql = "insert into car_gps(deployNum,plateNum,timeStr,lng,lat) values(?,?,?,?,?)"
    statement = conn.prepareStatement(sql)
    statement.setString(1, value.deployNum)
    statement.setString(2, value.plateNum)
    statement.setString(3, value.timeStr)
    statement.setString(4, arr(0))
    statement.setString(5, arr(1))
    statement.executeUpdate()
  }

  override def close(errorOrNull: Throwable): Unit = {
    if (null != conn) conn.close()
    if (null != statement) statement.close()
  }
}

object JdbcHelper {
  var conn: Connection = _
  val url = "jdbc:mysql://hadoop1:3306/test?useUnicode=true&characterEncoding=utf8"
  val username = "root"
  val password = "123456"
  def openConnection: Connection = {
    if (null == conn || conn.isClosed) {
      val p = new Properties
      Class.forName("com.mysql.jdbc.Driver")
      conn = DriverManager.getConnection(url, username, password)
    }
    conn
  }
}

2.2 通過foreach來寫入mysql

package com.hoult.Streaming.work
import com.hoult.structed.bean.BusInfo
import org.apache.spark.sql.{Column, DataFrame, Dataset, SparkSession}

object KafkaToJdbc {
  def main(args: Array[String]): Unit = {
    System.setProperty("HADOOP_USER_NAME", "root")
    //1 獲取sparksession
    val spark: SparkSession = SparkSession.builder()
      .master("local[*]")
      .appName(KafkaToJdbc.getClass.getName)
      .getOrCreate()
    spark.sparkContext.setLogLevel("WARN")
    import spark.implicits._
    //2 定義讀取kafka數據源
    val kafkaDf: DataFrame = spark.readStream
      .format("kafka")
      .option("kafka.bootstrap.servers", "linux121:9092")
      .option("subscribe", "test_bus_info")
      .load()
    //3 處理數據
    val kafkaValDf: DataFrame = kafkaDf.selectExpr("CAST(value AS STRING)")
    //轉為ds
    val kafkaDs: Dataset[String] = kafkaValDf.as[String]
    //解析出經緯度數據,寫入redis
    //封裝為一個case class方便后續獲取指定字段的數據
    val busInfoDs: Dataset[BusInfo] = kafkaDs.map(BusInfo(_)).filter(_ != null)

    //將數據寫入MySQL表
    busInfoDs.writeStream
      .foreach(new JdbcHelper)
      .outputMode("append")
      .start()
      .awaitTermination()
  }
}

2.4 創建topic和從消費者端寫入數據

kafka-topics.sh --zookeeper linux121:2181/myKafka --create --topic test_bus_info --partitions 2 --replication-factor 1
kafka-console-producer.sh --broker-list linux121:9092 --topic test_bus_info

到此,相信大家對“Spark-Streaming如何處理數據到mysql中”有了更深的了解,不妨來實際操作一番吧!這里是億速云網站,更多相關內容可以進入相關頻道進行查詢,關注我們,繼續學習!

向AI問一下細節

免責聲明:本站發布的內容(圖片、視頻和文字)以原創、轉載和分享為主,文章觀點不代表本網站立場,如果涉及侵權請聯系站長郵箱:is@yisu.com進行舉報,并提供相關證據,一經查實,將立刻刪除涉嫌侵權內容。

AI

沂源县| 诸暨市| 邛崃市| 延川县| 唐山市| 安龙县| 称多县| 邯郸县| 北票市| 新沂市| 阿瓦提县| 肃宁县| 辽阳市| 上虞市| 博湖县| 凭祥市| 开原市| 安陆市| 大理市| 仙居县| 宾阳县| 大同市| 峡江县| 镇远县| 榆林市| 吉木乃县| 太仓市| 平武县| 德保县| 永宁县| 天全县| 扶风县| 奇台县| 郸城县| 通江县| 安福县| 额济纳旗| 静海县| 棋牌| 囊谦县| 阳新县|