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

pyspark對Mysql數(shù)據(jù)庫進(jìn)行讀寫的實現(xiàn)

 更新時間:2020年12月30日 11:40:29   作者:FTDdata  
這篇文章主要介紹了pyspark對Mysql數(shù)據(jù)庫進(jìn)行讀寫的實現(xiàn),文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧

pyspark是Spark對Python的api接口,可以在Python環(huán)境中通過調(diào)用pyspark模塊來操作spark,完成大數(shù)據(jù)框架下的數(shù)據(jù)分析與挖掘。其中,數(shù)據(jù)的讀寫是基礎(chǔ)操作,pyspark的子模塊pyspark.sql 可以完成大部分類型的數(shù)據(jù)讀寫。文本介紹在pyspark中讀寫Mysql數(shù)據(jù)庫。

1 軟件版本

在Python中使用Spark,需要安裝配置Spark,這里跳過配置的過程,給出運(yùn)行環(huán)境和相關(guān)程序版本信息。

  • win10 64bit
  • java 13.0.1
  • spark 3.0
  • python 3.8
  • pyspark 3.0
  • pycharm 2019.3.4

2 環(huán)境配置

pyspark連接Mysql是通過java實現(xiàn)的,所以需要下載連接Mysql的jar包。

下載地址

在這里插入圖片描述

選擇下載Connector/J,然后選擇操作系統(tǒng)為Platform Independent,下載壓縮包到本地。

在這里插入圖片描述

然后解壓文件,將其中的jar包mysql-connector-java-8.0.19.jar放入spark的安裝目錄下,例如D:\spark\spark-3.0.0-preview2-bin-hadoop2.7\jars。

在這里插入圖片描述

環(huán)境配置完成!

3 讀取Mysql

腳本如下:

from pyspark.sql import SQLContext, SparkSession

if __name__ == '__main__':
  # spark 初始化
  spark = SparkSession. \
    Builder(). \
    appName('sql'). \
    master('local'). \
    getOrCreate()
  # mysql 配置(需要修改)
  prop = {'user': 'xxx', 
      'password': 'xxx', 
      'driver': 'com.mysql.cj.jdbc.Driver'}
  # database 地址(需要修改)
  url = 'jdbc:mysql://host:port/database'
  # 讀取表
  data = spark.read.jdbc(url=url, table='tb_newCity', properties=prop)
  # 打印data數(shù)據(jù)類型
  print(type(data))
  # 展示數(shù)據(jù)
  data.show()
  # 關(guān)閉spark會話
  spark.stop()
  • 注意點(diǎn):
  • prop參數(shù)需要根據(jù)實際情況修改,文中用戶名和密碼用xxx代替了,driver參數(shù)也可以不需要;
  • url參數(shù)需要根據(jù)實際情況修改,格式為jdbc:mysql://主機(jī):端口/數(shù)據(jù)庫;
  • 通過調(diào)用方法read.jdbc進(jìn)行讀取,返回的數(shù)據(jù)類型為spark DataFrame;

運(yùn)行腳本,輸出如下:

在這里插入圖片描述

4 寫入Mysql

腳本如下:

import pandas as pd
from pyspark import SparkContext
from pyspark.sql import SQLContext, Row

if __name__ == '__main__':
  # spark 初始化
  sc = SparkContext(master='local', appName='sql')
  spark = SQLContext(sc)
  # mysql 配置(需要修改)
  prop = {'user': 'xxx',
      'password': 'xxx',
      'driver': 'com.mysql.cj.jdbc.Driver'}
  # database 地址(需要修改)
  url = 'jdbc:mysql://host:port/database'

  # 創(chuàng)建spark DataFrame
  # 方式1:list轉(zhuǎn)spark DataFrame
  l = [(1, 12), (2, 22)]
  # 創(chuàng)建并指定列名
  list_df = spark.createDataFrame(l, schema=['id', 'value']) 
  
  # 方式2:rdd轉(zhuǎn)spark DataFrame
  rdd = sc.parallelize(l) # rdd
  col_names = Row('id', 'value') # 列名
  tmp = rdd.map(lambda x: col_names(*x)) # 設(shè)置列名
  rdd_df = spark.createDataFrame(tmp) 
  
  # 方式3:pandas dataFrame 轉(zhuǎn)spark DataFrame
  df = pd.DataFrame({'id': [1, 2], 'value': [12, 22]})
  pd_df = spark.createDataFrame(df)

  # 寫入數(shù)據(jù)庫
  pd_df.write.jdbc(url=url, table='new', mode='append', properties=prop)
  # 關(guān)閉spark會話
  sc.stop()

注意點(diǎn):

propurl參數(shù)同樣需要根據(jù)實際情況修改;

寫入數(shù)據(jù)庫要求的對象類型是spark DataFrame,提供了三種常見數(shù)據(jù)類型轉(zhuǎn)spark DataFrame的方法;

通過調(diào)用write.jdbc方法進(jìn)行寫入,其中的model參數(shù)控制寫入數(shù)據(jù)的行為。

model 參數(shù)解釋
error 默認(rèn)值,原表存在則報錯
ignore 原表存在,不報錯且不寫入數(shù)據(jù)
append 新數(shù)據(jù)在原表行末追加
overwrite 覆蓋原表

5 常見報錯

Access denied for user …

在這里插入圖片描述

原因:mysql配置參數(shù)出錯
解決辦法:檢查user,password拼寫,檢查賬號密碼是否正確,用其他工具測試mysql是否能正常連接,做對比檢查。

No suitable driver


原因:沒有配置運(yùn)行環(huán)境
解決辦法:下載jar包進(jìn)行配置,具體過程參考本文的2 環(huán)境配置。

到此這篇關(guān)于pyspark對Mysql數(shù)據(jù)庫進(jìn)行讀寫的實現(xiàn)的文章就介紹到這了,更多相關(guān)pyspark Mysql讀寫內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

最新評論

郧西县| 天水市| 广宗县| 马关县| 共和县| 卢湾区| 连云港市| 山东省| 乡城县| 建宁县| 迁西县| 平度市| 扶沟县| 中宁县| 云和县| 师宗县| 孟津县| 托里县| 新密市| 女性| 龙里县| 余干县| 曲周县| 牙克石市| 精河县| 江山市| 阿尔山市| 文昌市| 盐山县| 杨浦区| 恩施市| 都匀市| 桃园市| 措美县| 炉霍县| 尉犁县| 孝感市| 诏安县| 普格县| 唐海县| 遂川县|