林子雨、郑海山、赖永炫编著《Spark编程基础(Python版)》(教材官网)教材中的代码,在纸质教材中的印刷效果,可能会影响读者对代码的理解,为了方便读者正确理解代码或者直接拷贝代码用于上机实验,这里提供全书配套的所有代码。

查看所有章节代码

第5章 Spark SQL

from pyspark import SparkContext,SparkConf

from pyspark.sql import SparkSession

spark = SparkSession.builder.config(conf = SparkConf()).getOrCreate()

>>> df=spark.read.json("file:///usr/local/spark/examples/src/main/resources/people.json")

>>> df.show()

>>> peopleDF = spark.read.format("json").\

... load("file:///usr/local/spark/examples/src/main/resources/people.json")

>>> peopleDF.select("name", "age").write.format("json").\

... save("file:///usr/local/spark/mycode/sparksql/newpeople.json")

>>> peopleDF.select("name").write.format("text").\

... save("file:///usr/local/spark/mycode/sparksql/newpeople.txt")

>>> peopleDF = spark.read.format("json").\

... load("file:///usr/local/spark/mycode/sparksql/newpeople.json")

>>> peopleDF.show()

>>> df=spark.read.json("file:///usr/local/spark/examples/src/main/resources/people.json")

>>>df.printSchema()

>>>df.select(df["name"],df["age"]+1).show()

>>> df.filter(df["age"]>20).show()

>>> df.groupBy("age").count().show()

>>> df.sort(df["age"].desc()).show()

>>> df.sort(df["age"].desc(),df["name"].asc()).show()

>>> from pyspark.sql import Row

>>> people = spark.sparkContext.\

... textFile("file:///usr/local/spark/examples/src/main/resources/people.txt").\

... map(lambda line: line.split(",")).\

... map(lambda p: Row(name=p[0], age=int(p[1])))

>>> schemaPeople = spark.createDataFrame(people)

#必须注册为临时表才能供下面的查询使用

>>> schemaPeople.createOrReplaceTempView("people")

>>> personsDF = spark.sql("select name,age from people where age > 20")

#DataFrame中的每个元素都是一行记录,包含name和age两个字段,分别用p.name和p.age来获取值

>>> personsRDD=personsDF.rdd.map(lambda p:"Name: "+p.name+ ","+"Age: "+str(p.age))

>>> personsRDD.foreach(print)

Name: Michael,Age: 29

Name: Andy,Age: 30

>>> from pyspark.sql.types import *

>>> from pyspark.sql import Row

#下面生成“表头”

>>> schemaString = "name age"

>>> fields = [StructField(field_name, StringType(), True) for field_name in schemaString.split(" ")]

>>> schema = StructType(fields)

#下面生成“表中的记录”

>>> lines = spark.sparkContext.\

... textFile("file:///usr/local/spark/examples/src/main/resources/people.txt")

>>> parts = lines.map(lambda x: x.split(","))

>>> people = parts.map(lambda p: Row(p[0], p[1].strip()))

#下面把“表头”和“表中的记录”拼装在一起

>>> schemaPeople = spark.createDataFrame(people, schema)

#注册一个临时表供下面查询使用

>>> schemaPeople.createOrReplaceTempView("people")

>>> results = spark.sql("SELECT name,age FROM people")

>>> results.show()

+-------+---+

| name|age|

+-------+---+

|Michael| 29|

| Andy| 30|

| Justin| 19|

+-------+---+

sudo service mysql start

mysql -u root -p #屏幕会提示输入密码

mysql> create database spark;

mysql> use spark;

mysql> create table student (id int(4), name char(20), gender char(4), age int(4));

mysql> insert into student values(1,'Xueqian','F',23);

mysql> insert into student values(2,'Weiliang','M',24);

mysql> select * from student;

>>> jdbcDF = spark.read \

.format("jdbc") \

.option("driver","com.mysql.jdbc.Driver") \

.option("url", "jdbc:mysql://localhost:3306/spark") \

.option("dbtable", "student") \

.option("user", "root") \

.option("password", "123456") \

.load()

>>> jdbcDF.show()

+---+--------+------+---+

| id| name|gender|age|

+---+--------+------+---+

| 1| Xueqian| F| 23|

| 2|Weiliang| M| 24|

+---+--------+------+---+

mysql> use spark;

mysql> select * from student;

InsertStudent.py

#!/usr/bin/env python3

from pyspark.sql import Row

from pyspark.sql.types import *

from pyspark import SparkContext,SparkConf

from pyspark.sql import SparkSession

spark = SparkSession.builder.config(conf = SparkConf()).getOrCreate()

#下面设置模式信息

schema = StructType([StructField("id", IntegerType(), True), \

StructField("name", StringType(), True), \

StructField("gender", StringType(), True), \

StructField("age", IntegerType(), True)])

#下面设置两条数据,表示两个学生的信息

studentRDD = spark \

.sparkContext \

.parallelize(["3 Rongcheng M 26","4 Guanhua M 27"]) \

.map(lambda x:x.split(" "))

#下面创建Row对象,每个Row对象都是rowRDD中的一行

rowRDD = studentRDD.map(lambda p:Row(int(p[0].strip()), p[1].strip(), p[2].strip(), int(p[3].strip())))

#建立起Row对象和模式之间的对应关系,也就是把数据和模式对应起来

studentDF = spark.createDataFrame(rowRDD, schema)

#写入数据库

prop = {}

prop['user'] = 'root'

prop['password'] = '123456'

prop['driver'] = "com.mysql.jdbc.Driver"

studentDF.write.jdbc("jdbc:mysql://localhost:3306/spark",'student','append', prop)

mysql> select * from student;

+------+-----------+--------+------+

| id | name | gender | age |

+------+-----------+--------+------+

| 1 | Xueqian | F | 23 |

| 2 | Weiliang | M | 24 |

| 3 | Rongcheng | M | 26 |

| 4 | Guanhua | M | 27 |

+------+-----------+--------+------+

4 rows in set (0.00 sec)

Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐