ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

pyspark开发

pyspark开发 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
返回列表