
pyspark開發pyspark認識關于PySpark,它是Python調用Spark的接口,可以通過調用Python API的方式來編寫Spark程序,它支持了大多數的Spark功能,比如SparkDataFrame、Spark SQL、Streaming、MLlib等等。Spark SQL使用這個模塊是Spark中用來處理結構化數據的,提供一個SparkDataFrame的東西并且自動解析為分布式SQL查詢數據。在Python的Pandas庫,也能大致了解了DataFrame,這個其實和它沒有太大的區別,只是調用的API可能有些不同罷了。通過使用Spark SQL來處理數據,比如可以用SQL語句、用SparkDataFrame的API或者Datasets API,可以按照需求隨心轉換,通過SparkDataFrame API 和 SQL 寫的邏輯,會被Spark優化器Catalyst自動優化成RDD,即便寫得不好也可能運行得很快(如果是直接寫RDD可能就掛了)。1、讀取數據RDD創建rdd=sc.parallelize([("Sam",28,88),("Flora",28,90),("Run",1,60)])df=rdd.toDF(["name","age","score"])DataFrame創建df=pd.DataFrame([['Sam',28,88],['Flora',28,90],['Run',1,60]],columns=['name','age','score'])Spark_df=spark.createDataFrame(df)```python3.List創建 ```python list_values=[['Sam',28,88],['Flora',28,90],['Run',1,60]]Spark_df=spark.createDataFrame(list_values,['name','age','score'])```python4.讀取文件創建 ```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")數據庫讀取# (1) 讀取hive數據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數據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簡單處理1. 查看DataFrame的APIs# (1)以列表形式返回行df.collect()# (2)返回統計數量df.count()# (3)返回字段列表df.columns# (4)返回數據類型df.dtypes# (5)返回列的基礎統計信息,describe("非必須")df.describe(['col_name'])# (6)選定指定列并按照一定順序呈現df.select("col_name1","col_name2")# (7)查看第1條數據df.first()df.head(1)# (8)查看指定列的枚舉值df.freqItems(["col_name1","col_name2"])# (9)返回統計摘要df.summary()# (10)按照一定規則從df隨機抽樣數據df.sample(0.5)```python#### 2. 簡單處理DataFrame的APIs```python# (1)對數據集進行去重df.distinct()# (2)對指定列去重df.dropDuplicates(["col_name"])# (3)根據指定的df對df進行去重df1.exceptAll(df2)# exceptAll()進行df1 - df2的差集運算,保留重復項df1.subtract(df2)# subtract()獲取兩個 DataFrame 的行級差集,即找出在第一個 DataFrame 中存在,但在第二個中不存在的行# (4)返回兩個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)根據條件過濾df.filter(df.col_name50)# (11)數據集連接,DataFrame.join(other, on=None, how=None)df1.join(df2,df1.id