發(fā))
pyspark開(kāi)發(fā)pyspark認(rèn)識(shí)關(guān)于PySpark,它是Python調(diào)用Spark的接口,可以通過(guò)調(diào)用Python API的方式來(lái)編寫Spark程序,它支持了大多數(shù)的Spark功能,比如SparkDataFrame、Spark SQL、Streaming、MLlib等等。Spark SQL使用這個(gè)模塊是Spark中用來(lái)處理結(jié)構(gòu)化數(shù)據(jù)的,提供一個(gè)SparkDataFrame的東西并且自動(dòng)解析為分布式SQL查詢數(shù)據(jù)。在Python的Pandas庫(kù),也能大致了解了DataFrame,這個(gè)其實(shí)和它沒(méi)有太大的區(qū)別,只是調(diào)用的API可能有些不同罷了。通過(guò)使用Spark SQL來(lái)處理數(shù)據(jù),比如可以用SQL語(yǔ)句、用SparkDataFrame的API或者Datasets API,可以按照需求隨心轉(zhuǎn)換,通過(guò)SparkDataFrame API 和 SQL 寫的邏輯,會(huì)被Spark優(yōu)化器Catalyst自動(dòng)優(yōu)化成RDD,即便寫得不好也可能運(yùn)行得很快(如果是直接寫RDD可能就掛了)。1、讀取數(shù)據(jù)RDD創(chuàng)建rdd=sc.parallelize([("Sam",28,88),("Flora",28,90),("Run",1,60)])df=rdd.toDF(["name","age","score"])DataFrame創(chuàng)建df=pd.DataFrame([['Sam',28,88],['Flora',28,90],['Run',1,60]],columns=['name','age','score'])Spark_df=spark.createDataFrame(df)```python3.List創(chuàng)建 ```python list_values=[['Sam',28,88],['Flora',28,90],['Run',1,60]]Spark_df=spark.createDataFrame(list_values,['name','age','score'])```python4.讀取文件創(chuàng)建 ```python# (1) CSV文件df=spark.read.option("header","true")\.option("inferSchema","true")\.option("delimiter",",")\.csv("./test/data/titanic/train.csv")# (2) json文件df=spark.read.json("./test/data/hello_samshare.json")數(shù)據(jù)庫(kù)讀取# (1) 讀取hive數(shù)據(jù)spark.sql("CREATE TABLE IF NOT EXISTS src (key INT, value STRING) USING hive")spark.sql("LOAD DATA LOCAL INPATH 'data/kv1.txt' INTO TABLE src")df=spark.sql("SELECT key, value FROM src WHERE key 10 ORDER BY key")# (2) 讀取mysql數(shù)據(jù)url="jdbc:mysql://localhost:3306/test"df=spark.read.format("jdbc")\.option("url",url)\.option("dbtable","runoob_tbl")\.option("user","root")\.option("password","8888")\.load()2、DataFrame簡(jiǎn)單處理1. 查看DataFrame的APIs# (1)以列表形式返回行df.collect()# (2)返回統(tǒng)計(jì)數(shù)量df.count()# (3)返回字段列表df.columns# (4)返回?cái)?shù)據(jù)類型df.dtypes# (5)返回列的基礎(chǔ)統(tǒng)計(jì)信息,describe("非必須")df.describe(['col_name'])# (6)選定指定列并按照一定順序呈現(xiàn)df.select("col_name1","col_name2")# (7)查看第1條數(shù)據(jù)df.first()df.head(1)# (8)查看指定列的枚舉值df.freqItems(["col_name1","col_name2"])# (9)返回統(tǒng)計(jì)摘要df.summary()# (10)按照一定規(guī)則從df隨機(jī)抽樣數(shù)據(jù)df.sample(0.5)```python#### 2. 簡(jiǎn)單處理DataFrame的APIs```python# (1)對(duì)數(shù)據(jù)集進(jìn)行去重df.distinct()# (2)對(duì)指定列去重df.dropDuplicates(["col_name"])# (3)根據(jù)指定的df對(duì)df進(jìn)行去重df1.exceptAll(df2)# exceptAll()進(jìn)行df1 - df2的差集運(yùn)算,保留重復(fù)項(xiàng)df1.subtract(df2)# subtract()獲取兩個(gè) DataFrame 的行級(jí)差集,即找出在第一個(gè) DataFrame 中存在,但在第二個(gè)中不存在的行# (4)返回兩個(gè)DataFrame的交集df1.intersectAll(df2)# (5)丟棄指定列df.drop('col_name')# (6)新增列df.withColumn("col_name",col_value)# (7)重命名列名df.withColumnRenamed("col_ora_name","col_new_name")# (8)丟棄空值,DataFrame.dropna(how='any', thresh=None, subset=None)df.dropna(how='all',subset=['col_name'])# (9)空值填充操作df.fillna({"col_name1":"col_value1","col_name2":col_value2})# (10)根據(jù)條件過(guò)濾df.filter(df.col_name50)# (11)數(shù)據(jù)集連接,DataFrame.join(other, on=None, how=None)df1.join(df2,df1.id