版权说明:本文档由用户提供并上传,收益归属内容提供方,若内容存在侵权,请进行举报或认领
文档简介
任务3SparkSQL离线计算内容导航01任务概述与准备明确任务目标、思维导图与整体框架02工单实操演练5个工单逐步掌握DataFrame操作与SQL编程03知识链接梳理系统讲解SparkSQL基础、编程与调优04任务总结收束回顾核心要点,巩固学习成果任务概述与准备01编写Python程序,运用SparkSQL完成电影评分数据集的离线计算本任务将编写
Python
程序,运用
SparkSQL
对
电影评分数据集
进行离线计算,包括单数据离线计算、多数据连接离线计算、SparkSQL调优等操作。1读取数据集使用SparkSQL读取数据集,形成DataFrame,并练习常用函数SparkSQLDataFrame2两种编码风格针对DataFrame分别使用SQL和DSL两种风格编写代码SQLDSL3单数据离线计算完成5个操作,实现单数据文件DSL风格离线计算5个操作DSL风格4多数据连接计算采用SQL风格对多数据文件连接结果进行离线计算SQL风格多文件连接5SparkSQL调优最后进行SparkSQL调优调优任务描述与整体框架工单实操演练02工单3.1:基本信息与目标DataFrame基本操作工单工单3.1建议学时2学时环境要求onYARN模式的Spark集群(已部署)工单说明1对从
HDFS
读取的文件进行split()和map()操作,然后创建
DataFrame2在
PySpark交互界面
中读取本地JSON文件,写入本地CSV文件3根据没有列名的数据结构创建DataFrame,然后进行简单的数据清洗工单目标知识目标掌握从原始数据到DataFrame的转换方法掌握DataFrame基础操作技能目标能创建DataFrame对象能操作DataFrame进行数据查询、输出素养目标培养数据质量优先的处理意识培养解决实际问题的能力启动集群与RDD转DataFrame集群启动命令与RDD转DataFrame实操代码1步骤1:启动Hadoop和Spark集群在master节点执行以下命令:start-dfs.sh#启动HDFSstart-yarn.sh#启动YARNstart-master.sh#启动Spark主节点start-workers.sh#启动Spark工作节点2步骤2:RDD转DataFrame读取
u.user
文件,抽取观众ID和性别创建DataFrame。在
/opt/example
目录新建
movie_sql1.py:frompyspark.sqlimportSparkSessionif__name__=='__main__':spark=SparkSession.builder.appName("createdf")\.master("spark://master:7077")\.getOrCreate()sc=spark.sparkContextrdd=sc.textFile("/user/zjaf/data/u.user")\.map(lambdaline:line.split("|"))\.map(lambdafields:(int(fields[0]),fields[2]))df=spark.createDataFrame(rdd,schema=['UserID','Gender'])df.printSchema()df.show()提交运行:spark-submitmovie_sql1.py步骤3文件读写:JSON与CSV使用SparkSQL读取本地JSON文件,然后写入本地CSV文件①在master节点启动PySpark$pyspark--masterspark://master:7077②运行代码,读取JSON文件>>>
frompyspark.sqlimportSparkSession>>>spark=SparkSession.builder.appName("df1").master(
"spark://master:7077").getOrCreate()>>>infile='file:///opt/spark/examples/src/main/resources/people.json'>>>df=spark.read.json(infile)>>>df.show()③输出结果
df.show()agenamenullMichael30Andy19Justin④写入CSV文件>>>outfile=r"file:///opt/test.csv">>>df.write.csv(path=outfile)步骤4根据无列名数据表创建DataFrame并进行数据清洗创建DataFrame(指定Schema与分隔符)frompyspark.sqlimportSparkSessionfrompyspark.sql.typesimportStructType,StructField,StringType,IntegerTypespark=SparkSession.builder.appName("dfop").master(
"spark://master:7077").getOrCreate()schema=StructType().add("user_id",StringType(),True).\add("movie_id",IntegerType(),True).\add("rank",IntegerType(),True).\add("ts",StringType(),True)df2=spark.read.format("csv").\option("sep","\t").\option("header",False).\schema(schema).\load("/user/zjaf/data/u.data")df2.show()StructType()
定义表头,add()参数依次为:字段名、字段类型、是否允许null数据清洗:填充空值&去重df3=df2.fillna('unknown')df4=df3.dropDuplicates(subset=['user_id'])清洗方法说明方法作用fillna()替换满足条件的null值dropDuplicates()数据去重,重复时保留第一条数据清洗操作工单3.1小结与素养课堂✓工单小结1在
PySpark交互界面
中使用
SparkSQL
完成读写JSON文件2创建
DataFrame
并进行简单的
数据清洗★素养课堂迎难而上詹天佑的工程智慧京张铁路修建面临三大困难,詹天佑带领团队实地勘测、周密计算,于
1909年
建成中国人自主设计的第一条铁路。地势险要气候恶劣缺乏设备创新工法"竖井开凿法"
"人"字形线路数据处理启示缺失值处理需检查、删除或填充——正需
"迎难而上"
之精神。工单3.2建议学时
2|环境要求:部署好的
onYARN
模式的Spark集群SparkSQL单词计数工单说明1使用
SQL
和
DSL
两种风格编写SparkSQL代码,实现单词计数功能。2编写
Python程序并提交至Spark集群执行,得到单词计数的结果。3完成效果如图3.3
所示。学习目标知识目标掌握SQL风格掌握DSL风格技能目标能够使用SQL风格编程实现单词计数能够使用DSL风格编程实现单词计数素养目标培养多风格编程的开发能力培养与人交流、合作的能力工单3.2:基本信息与目标SQL风格单词计数SQL风格实现单词计数:注册临时视图→编写SQL查询split()
将文本每行按指定分隔符分割,得到Array类型的DataFrame;explode()
将分割结果展开为单个单词,再分组统计数量并按降序排列。1①
切分·split()按分隔符分割文本,得到Array类型2②
展开·explode()将一行扩展为多行,提取单个单词3③
统计·groupby+count()+orderby分组计数并按数量降序排列PYTHON·PYSPARKSQLfrompyspark.sqlimportSparkSessionspark=SparkSession.builder.appName("dfop").master(
"spark://master:7077").getOrCreate()#读取文本文件df=spark.read.text("/user/zjaf/data/test/test1.txt")df.printSchema()df.show()#注册临时视图,视图名为words_tdf.createOrReplaceTempView('words_t')#在临时视图上编写SQL语句进行查询spark.sql('''selectword,count(1)ascntfrom(selectexplode(split(value,''))aswordfromwords_t)tgroupbywordorderbycntdesc''').show()DSL风格单词计数直接调用DataFrame与SQL对应的API完成split()→explode()→groupBy()→count()→orderBy()完整代码Python/PySparkfrompyspark.sqlimportSparkSessionfrompyspark.sqlimportfunctionsasFspark=SparkSession.builder.appName("dfop").master(
"spark://master:7077").getOrCreate()df=spark.read.text("/user/zjaf/data/test1.txt")df.printSchema()df.show()df2=df.select(F.split('value','').alias('arr'))df3=df2.select(F.explode('arr').alias('word'))df4=df3.groupBy('word').count()df5=df4.orderBy('count',ascending=False)df5.show()运行结果df5.show()wordcountI3love3my2college1China1family1工单小结本工单通过
PySpark交互界面,使用两种风格编程实现单词计数功能:SQL查询风格编程DSL领域语言风格编程目标实现单词计数功能素养课堂2024年6月25日14时7分,嫦娥六号返回器携带来自月背的月球样品安全着陆专注投入像航天人一样执着目标精益求精从"神舟"到"嫦娥",追求极致一丝不苟严谨对待每一个技术细节追求卓越在钻研中勇于创新网络时代娱乐应用纷繁,年轻人尤需保持
专注学习的态度工单3.2小结与素养课堂工单3.3:基本信息与目标单数据离线计算工单编号3.3工单主题单数据离线计算建议学时3学时实验环境onYARN模式Spark集群工单说明对电影评分数据集中的
3个文件(
u.user
、
u.data
、
u.item
)
分别使用
DSL风格编写SparkSQL代码进行离线计算。通过计算每个观众的平均评分、每部电影的平均评分、评分大于平均评分的电影数量、每个观众的简单统计信息、被评分超过100次的电影平均分,帮助学生熟练使用DSL风格编写SparkSQL代码。教学目标知识目标掌握DSL风格SparkSQL代码的编写格式;掌握关联查询的性能优化思路技能目标能够实现基于条件的聚合计算与筛选素养目标培养面向业务的DSL风格代码设计思维;培养解决实际问题的能力读取数据与观众平均分启动集群与PySpark,读取HDFS中的u.data,按观众ID分组计算平均分并降序排列步骤1-3启动集群·启动PySpark·读取数据读取HDFS中的u.data并添加表头frompyspark.sqlimportSparkSessionfrompyspark.sql.typesimportStructType,StructField,StringType,IntegerTypefrompyspark.sqlimportfunctionsasfuncspark=SparkSession.builder.appName("dfsql2").master(
"spark://master:7077").getOrCreate()schema=StructType().add("userID",StringType(),True).\add("movieID",IntegerType(),True).\add("rating",IntegerType(),True).\add("time",StringType(),True)df=spark.read.format("csv").\option("sep","\t").\option("header",False).\schema(schema).\load("/user/zjaf/data/u.data")步骤4计算观众平均分按观众ID分组,计算平均分并降序排列df.groupBy("userID").\avg("rating").\withColumnRenamed("avg(rating)","avg_rating").\withColumn("avg_rating",func.round("avg_rating",2)).\orderBy("avg_rating",ascending=False).\show()groupBy()分组→avg()平均→重命名→orderBy()降序电影平均分与评分筛选5步骤5按电影ID分组,计算平均分并降序排列df.groupBy("movieID").\.avg("rating").\.withColumnRenamed("avg(rating)","avg_rating").\.withColumn("avg_rating",func.round("avg_rating",2)).\.orderBy("avg_rating",ascending=False).\.show()6步骤6计算评分大于平均分的电影数量num=df.where(df['rating']>df.select(func.avg(df['rating'])).first()['avg(rating)']).count()print("movienumber(>avg_rating):%d"%num)执行顺序func.avg()计算平均分→where()过滤→count()统计→print输出7计算每个观众的平均评分、最低评分、最高评分df.groupBy("userID").\.agg(
func.round(func.avg('rating'),2).alias("avg_rating"),
func.min('rating').alias("min_rank"),
func.max('rating').alias("max_rank")).show()agg()与groupBy()配合使用,对分组数据进行聚合计算alias()为计算结果列指定别名8计算被评分超100次电影的平均分,显示Top10df.groupBy("movieID").\.agg(
func.count("movieID").alias("cnt"),
func.round(func.avg("rating"),2).alias("avg_rating")).where("cnt>100").\.orderBy("avg_rating",ascending=False).\.limit(10).show()观众统计与热门电影评分✓工单小结采用
PySpark交互界面+DSL风格
完成5个离线计算操作:1按观众ID分组,计算平均分并
降序排列2按电影ID分组,计算平均分并
降序排列3计算
评分>平均分
的电影数量4计算每个观众的
平均/最低/最高
评分5筛选被评分
>100次
的电影,取平均分
Top10★素养课堂交通信号灯程序设计引导遵守规则,提高安全意识质量第一,及时
Debug维护社会秩序,促进社会文明充分
测试用例(Case)
验证社会道德与秩序的重要体现严谨认真、精益求精所有逻辑推导和方案设计,都必须通过以下两个环节:程序调试测试用例验证工单3.3小结与素养课堂工单3.4:基本信息与目标工单3.4多数据连接离线计算建议学时4课时环境要求部署好的onYARN模式的Spark集群工单说明对电影评分数据集中的
3
个文件
u.user、
u.data、
u.item
采用
SQL风格
编写SparkSQL代码进行离线计算。通过计算以下三项内容,帮助学生掌握SQL风格编写SparkSQL代码的方法:三项计算内容1看过指定电影的观众的年龄和性别2男性看过最多的10部电影3年龄45~49岁男性观众中评分最高的10部电影学习目标知识目标掌握SQL风格SparkSQL代码的编写格式技能目标能够使用SparkSQL对多个数据文件进行联合查询能够编写清晰、可维护的
SparkSQL代码素养目标培养面向业务的
SQL风格代码设计思维培养解决实际问题的能力201启动集群→2启动PySpark→3读取HDFS三文件read_u_data.pyfrompyspark.sqlimportSparkSessionfrompyspark.sql.typesimportStructType,StructField,StringType,IntegerTypefrompyspark.sqlimportfunctionsasfuncspark=SparkSession.builder.appName("dfsql3").master(
"spark://master:7077").getOrCreate()#读取u.data(评分数据,制表符分隔)schema1=StructType().add("userID",IntegerType(),True).\add("movieID",IntegerType(),True).\add("rating",IntegerType(),True).\add("time",StringType(),True)df1=spark.read.format("csv").\option("sep","\t").\option("header",False).\schema(schema1).\load("/user/zjaf/data/u.data")HDFS三个数据文件文件数据内容分隔符u.user观众ID、年龄、性别、职业、邮编竖线|u.item电影ID、名称、拍摄时间等竖线|u.data观众ID、电影ID、评分值、时间制表符\t使用
StructType()
构建表头,通过
read()
读取文件并以
schema()
指定表头:读取u.data并创建DataFrame读取u.item与u.user文件定义Schema,以竖线「|」分隔读取CSV,创建电影数据与观众数据DataFrameuu.item——电影数据竖线分隔·含movieID+名称+19个分类字段#读取u.item(电影数据,竖线分隔,含movieID+名称+19个分类字段)schema2=StructType().add("movieID",IntegerType(),True)\.add("movieName",StringType(),True)\.add("D1",StringType(),True)\.add("D2",StringType(),True)\.add("URL",StringType(),True)#...C1至C19字段df2=spark.read.format("csv")\.option("sep","|")\.option("header",False)\.schema(schema2)\.load("/user/zjaf/data/u.item")uu.user——观众数据竖线分隔·用户属性字段#读取u.user(观众数据,竖线分隔)schema3=StructType().add("userID",IntegerType(),True)\.add("age",IntegerType(),True)\.add("gender",StringType(),True)\.add("career",StringType(),True)\.add("zip",StringType(),True)df3=spark.read.format("csv")\.option("sep","|")\.option("header",False)\.schema(schema3)\.load("/user/zjaf/data/u.user")列名重命名&结果预览df3=df3.toDF("userID","age","gender","career","zip")df1.show(2)df2.show(2)df3.show(2)创建视图与查询ToyStory观众4创建临时视图对3个DataFrame分别通过
createTempView()
创建临时视图:df1.createTempView("t_rating")df2.createTempView("t_movie")df3.createTempView("t_user")5计算看过"ToyStory(1995)"的观众年龄和性别查询思路:根据名称求电影ID→借助u.data求观众ID→根据观众ID求性别和年龄spark.sql('''selectc.userID,c.age,casec.genderwhen'M'then'male'else'female'endasgenderfromt_ratingaleftjoint_usercona.userID=c.userIDleftjoint_moviebona.movieID=b.movieIDwhereb.movieName='ToyStory(1995)'''').show()运行结果如图3.11所示。步骤6男性看过最多的10部电影⭐TopN问题·降序取前N条查询流程逻辑链条1子查询:筛选男性观众ID(gender='M')2连接查询:评分表↔电影表↔用户表3按电影ID和名称
groupby
分组4count()
统计观影人数5orderby
降序排列6limit10
取前10条结论运行结果如
图3.12
所示SparkSQL·Pythonspark.sql('''selecta.movieID,count(b.movieName)ascountfromt_ratingaleftjoint_moviebona.movieID=b.movieIDleftjoint_usercona.userID=c.userIDwherea.userIDin(selectu.userIDfromt_useruwheregender='M')groupbya.movieID,b.movieNameorderbycountdesclimit10''').show()步骤745-49岁男性评分最高的10部电影计算年龄为
45~49岁
男性观众评分最高的
10部
电影——查询思路(五步递进),运行结果如图3.13所示查询思路·五步递进1子查询求得45-49岁
男性观众ID2根据观众ID求电影ID3按电影ID和名称分组4按评分降序排列5limit
取10条嵌套子查询逐层过滤,最终聚合排序取Top10SparkSQL查询代码Pythonspark.sql('''selecta.movieID,sum(a.rating)ratefrom
t_rating
aleftjoin
t_movie
bona.movieID=b.movieIDwherea.movieIDin(selecta.movieIDfrom
t_rating
aleftjoin
t_movie
bona.movieID=b.movieIDleftjoin
t_user
cona.userID=c.userIDwherea.userIDin(selectu.userIDfrom
t_user
uwhereage>=
45
andage<=
49
andgender=
'M'))groupbya.movieID,b.movieNameorderbyratedesclimit
10''').show()工单3.4小结与素养课堂SKILLSSUMMARY工单3.4小结通过PySpark交互界面,采用
SQL风格编写SparkSQL代码,完成
3
个离线计算操作:1计算看过电影
"ToyStory(1995)"
的观众的年龄和性别2计算男性看过最多的
10
部电影3计算年龄45~49岁的男性观众中评分最高的
10
部电影LEGALLITERACY素养课堂《中华人民共和国民法典》第三编
对技术合同作出系统规定,旨在规范技术合同活动,促进科技成果转化与应用,推动科技进步和创新。技术合同类型权利归属规则技术开发职务技术成果:使用权、转让权属于法人/非法人组织技术转让非职务技术成果:使用权、转让权归完成技术成果的个人技术咨询技术服务工单3.5:基本信息与目标工单3.5SparkSQL调优建议学时2环境要求:部署好的onYARN模式的Spark集群工单说明u.useru.datau.item←电影评分数据集3个文件对电影评分数据集中的3个文件
u.user、u.data、u.item
采用SQL风格编写SparkSQL代码进行优化处理。通过优化前和优化后的代码实现计算
1995年最受欢迎的前3部电影的观众的年龄和性别数量,帮助学生初步掌握SparkSQL调优的基本方法。学习目标知识目标掌握常用的SparkSQL调优方法技能目标能够实现SparkSQL调优素养目标培养积极探索、勇于创新的科学素养;培养不畏困难、勇攀高峰的精神准备数据与读取u.data步骤1-3:启动集群→启动PySpark→准备数据上传数据通过hdfs命令上传数据文件BASHhdfsdfs-put/opt/example/u.user/user/zjaf/data读取数据读取
u.data
文件,定义schemaPYTHON/PYSPARKfrompyspark.sqlimportSparkSessionfrompyspark.sql.typesimportStructType,StructField,StringType,IntegerTypespark=SparkSession.builder.appName("dfsql3").master("spark://master:7077").getOrCreate()#u.data—用户评分数据schema1=StructType().add("userID",IntegerType(),True).\add("movieID",IntegerType(),True).\add("rating",IntegerType(),True).\add("time",StringType(),True)df1=spark.read.format("csv").\option("sep","\t").schema(schema1).\load("/user/zjaf/data/u.data")u.item电影信息数据(含movieID、名称、分类等)schema2=StructType().add("movieID",IntegerType(),True).\add("movieName",StringType(),True).\add("D1",StringType(),True).add("D2",StringType(),True).\add("URL",StringType(),True)#C1-C19分类字段省略df2=spark.read.format("csv").\option("sep","|").schema(schema2).\load("/user/zjaf/data/u.item")u.user用户属性数据schema3=StructType().add("userID",IntegerType(),True).\add("age",IntegerType(),True).\add("gender",StringType(),True).\add("career",StringType(),True).\add("zip",StringType(),True)df3=spark.read.format("csv").\option("sep","|").schema(schema3).\load("/user/zjaf/data/u.user")验证并注册临时视图df1.show(2);df2.show(2);df3.show(2)df1.createTempView("t_rating");df2.createTempView("t_movie");df3.createTempView("t_user")读取u.item与u.user并注册视图优化前后代码对比步骤4同学A·优化前sql='''selectcasefirst(cc.gender)when'M'then'male'else'female'endasgender,first(cc.age)agefromt_ratingaajoint_moviebbonaa.movieID=bb.movieIDjoint_usercconaa.userID=cc.userIDwheresubstr(aa.time,-5,4)=1995groupbyaa.movieIDorderbyavg(aa.rating)desclimit3'''spark.sql(sql).show(2)步骤5同学B·优化后sql='''selectcasegenderwhen'M'then'male'else'female'endasgender,c.ageagefromt_ratingaleftjoint_moviebona.movieID=b.movieIDleftjoint_usercona.userID=c.userIDwhereb.movieIDin(selecta.movieIDfromt_ratingaleftjoint_moviebona.movieID=b.movieIDwheresubstr(a.time,-5,4)=1995groupbya.movieIDorderbyavg(a.rating)desclimit3)'''spark.sql(sql).show(2)在写SQL语句时,不要轻易使用select*,应该指明具体列名,如spark.sql("selectmovieIDfromt_movie").show()查看DAG与工单小结技术操作1访问
:4040/jobs/
查看计算效率对比2Spark应用UI端口:4040
|JobTracker端口:80883单击Description超链接→选择「DAGVisualization」查看DAG工单小结通过PySpark交互界面完成SparkSQL调优,实现SQL查询效率提升。素养课堂·诚实守信商鞅立木建信——赏金从十金提至五十金,壮士搬木,如约兑现,终取民信。诚实守信,对民族、国家有利,对自己也有益。工作中:诚实劳动、求真务实、遵纪守法生活中:与人为善、坦诚相待、团结互爱、助人为乐诚信,既是道德规范的重要内容,更是立身之本、成事之基。知识链接梳理03SparkSQL的五大特点1高度融合性SQL查询与Spark程序集成,结构化数据视为RDD查询,支持Python、Scala、Java等语言API2统一的数据访问机制DataFrameAPI统一接口,支持Hive表、本地文件、JSON文件等多种数据源3Hive兼容性重用Hive前端与Metastore,直接计算生成Hive表,复用现有Hive资源4标准化的连接支持支持
JDBC/ODBC
标准协议,提供便捷数据库交互能力5出色的可扩展性交互式分析与批处理统一引擎,基于RDD容错机制高效处理大规模数据SparkSQL简介SparkSQL是Spark框架中专门用于处理结构化数据的模块,前身是Shark。关键节点说明起源加州大学伯克利分校实验室研发转型2014年起停止Shark维护,转向SparkSQL开发动因Shark对Hive过度依赖,难以实现新优化执行机制数据处理逻辑转换为
RDD,提交Spark集群执行应用场景离线开发、数据仓库、科学计算、数据分析SparkSQL简介与特点SparkSQL架构与运行流程架构SparkSQL三层核心架构编程语言API层支持Python、Scala、Java及HQL,提供SparkSQL使用入口SchemaRDD层基于RDD设计,作为临时表角色,使结构化数据处理更高效灵活数据源层支持Parquet、JSON、Hive表、Cassandra等多种数据源流程SparkSQL七步运行流程1通过API读取用户提交的SQL语句2利用
ANTLR4
解析SQL,生成未经验证的UnresolvedLogicalPlan3访问
Catalog
验证元数据,生成ResolvedLogicalPlan4Optimizer
通过RBO/CBO优化,生成OptimizedLogicalPlan5Catalyst
生成多候选PhysicalPlan,经costmodel评估选优6将最优PhysicalPlan编译为可执行代码7通过
WholeStageCodegen
编译为高效执行代码,生成RDD由集群执行DataFrame结构与层级关系DataFrame分布式数据存储结构,类似二维表格,记录数据与结构信息支持嵌套数据结构:StructArrayMap相比RDDAPI更直观易用SparkSQL可识别列名和数据类型结构层面层级关系层级对象作用1表结构StructType描述整体表结构2列定义StructField定义列名、数据类型、可空性3单行数据Row表示一行记录4单列数据Column表示一列数据Dataset与SparkSessionDataset强类型集合与Spark统一入口对象DDatasetSpark1.6引入·强类型集合Spark1.6引入的强类型集合,支持函数式编程和并行转换每个Dataset都有非类型化视图即DataFrameDataFrame是Dataset的特殊类型,编译时不进行模式检测Dataset将逐渐取代RDD成为主流选择SSparkSessionSpark2.0起统一入口对象·SQLContext与HiveContext的组合create_df.pyfrompyspark.sqlimportSparkSessionif__name__=='__main__':spark=SparkSession.builder.\
appName("createdf").\
master("local[*]").\
config("spark.sql.shuffle.partitions","2").\
getOrCreate()sc=spark.sparkContextDataFrame创建方式(一)共7种方式·本页介绍前3种1根据变量创建通过
createDataFrame()
将RDD转换为DataFrame,类型自动推断data=[(123,"Katie",19,"brown"),(234,"Michael",22,"green"),(345,"Simone",23,"blue")]df=spark.createDataFrame(data,schema=['id','name','age','eyecolor'])df.show()df.count()2读取JSON文件创建通过
read.json()
读取外部JSON文件file=r"file:///opt/spark/examples/src/main/resources/people.json"df=spark.read.json(file)df.show()3读取CSV文件创建通过
read.csv()
读取外部CSV文件df=spark.read.csv(file,header=True,inferSchema=True)df.show()DataFrame创建方式(二)4种进阶创建方式,覆盖主流数据源4读取MySQL数据指定JDBC的URL、数据表、用户名及密码,通过load()读取df=spark.read.format('jdbc').options(url='jdbc:mysql://',dbtable='mysql.db',user='root',password='123456').load()df.show()5pandasDataFrame转换通过pandas创建后转为SparkDataFramedf=pd.DataFrame(np.random.random((4,4)))spark_df=spark.createDataFrame(df,schema=['a','b','c','d'])DataFrame创建方式(三)读取Parquet文件与Hive数据两种外部数据源的DataFrame创建方式6读取Parquet文件读取列式存储的Parquet格式文件代码示例file=r"file:///opt/spark/examples/src/main/resources/users.parquet"df=spark.read.parquet(file)df.show()7读取Hive数据配置Spark连接Hive后执行SQL查询代码示例spark=SparkSession\.builder\.enableHiveSupport()\.master("70:7077")\.appName("my_first_app_name")\.getOrCreate()df=spark.sql("select*fromhive_tb_name")df.show()1数据排序df.sort('zip',ascending=False)#按指定列降序排序2数据连接df1=spark.createDataFrame([('a','1'),('b','2'),('c','3')],['xx','yy'])df2=spark.createDataFrame([('a','T'),('b','F'),('d','T')],['xx','zz'])df1.join(df2,on='xx').show()#按xx列内连接{}初始化示例frompyspark.sqlimportSparkSessionfrompyspark.sql.typesimportStructType,StructField,StringTypespark=SparkSession.builder.appName("dfop").master(
"spark://master:7077").getOrCreate()schema=StructType([StructField("area",StringType(),True),StructField("zip",StringType(),True)])data=[("Guangzhou","852"),("Chongqing","853"),("Fuzhou","886"),("Beijing","10"),("Shanghai","21"),("Shenzhen","57")]df=spark.createDataFrame(data=data,schema=schema)DataFrame常用操作(二)DataFrame是SparkSQL的核心数据结构,掌握其常用操作是进行数据离线计算的基础⌗显示表结构和数据df.printSchema()#查看表头df.describe().show()#基本统计df.columns#查看列名df.count()#查看行数df.show(truncate=False)#显示全部数据⧩数据筛选df.select('area').show()#单列选择df.select('area','zip').show()#多列筛选df.select(df.columns[0:1]).show()#列序号筛选df.filter((df['area']=='Fuzhou')&(df['zip']=='886')).show()#行筛选DataFrame常用操作(一)DataFrame数据写入(一)DataFrame数据写入使用
write
算子,常见方式如下:写入CSV文件、写入Parquet文件、写入Hive📄写入CSV文件write.csv()借助
write.csv()
函数,需指定文件名、是否写入表头、分隔符、写入模式file=r"file:///opt/spark/examples/src/main/resources/test.csv"spark_df.write.csv(path=file,header=True,sep=",",mode='overwrite')📦写入Parquet文件write.parquet()借助
write.parquet()
函数写入文件file=r"file:///opt/spark/examples/src/main/resources/test.parquet"spark_df.write.parquet(path=file,mode='overwrite')🐝写入Hivespark.sql()先设置动态分区,再通过
SQL语句写入分区表spark.sql("sethive.exec.dynamic.partition.mode=nonstrict")spark.sql("sethive.exec.dynamic.partition=true")spark.sql("""insertoverwritetableai.da_aipurchase_dailysale_hivepartition(saledate)selectproductid,propertyid,processcenterid,saleplatform,sku,poa,salecount,saledatefromszy_aipurchase_tmp_szy_dailysaledistributebysaledate""")DataFrame数据写入(二)1写入HDFSjdbcDF.write.mode("overwrite").options(header="true").csv("/user/sample.txt")2写入MySQLmode()
设置覆盖或追加,options()
指定连接参数:overwrite覆盖spark_df.write.mode("overwrite").format("jdbc").options(url='jdbc:mysql://',user='root',password='123456',dbtable="test.test",batchsize="1000").save()append追加spark_df.write.mode("append").format("jdbc").options(url='jdbc:mysql://',user='root',password='123456',dbtable="test.test",batchsize="1000").save()常见的数据清洗函数包括
去重、删除缺失值、填充缺失值
三类:表三类清洗函数函数功能说明dropDuplicates()去重对DataFrame数据去重,重复数据保留第一条dropna()删除缺失值how="any"
有一个缺失值即删除;how="all"
全部缺失才删除;thresh
指定非缺失值最小数量fillna()填充缺失值value
指定替换值,subset
指定目标列;支持指定值、均值、中位数、众数填充</>代码示例>>>frompyspark.sqlimportSparkSession>>>spark=SparkSession.builder.getOrCreate()>>>data=[(1,"Alice",25),(2,"Bob",None),(3,None,30),(4,"David",35),(5,"David",35)]>>>df=spark.createDataFrame(data,["id","name","age"])#按name列去重,保留第一条记录
>>>df2=df.dropDuplicates(subset=['name'])>>>df2.show()#删除包含缺失值的行
>>>df_cleaned=df2.dropna(how="any")>>>df_cleaned.show()#使用指定值填充name列的缺失值
>>>df_filled=df2.fillna(value="Jim",subset=['name'])>>>df_filled.show()数据清洗函数详解SQL执行顺序与优化原则SparkSQL优化的本质是尽可能减少无效数据计算、缓存数据以及均匀分布数据,从而最大化资源利用效率。SQL关键字的执行顺序优化原则1需求明确从零开始或优化SQL时,首要任务是明确具体需求2数据探查检查数据完整性、评估数据量级、研究表间关联关系、分析数据分布状况3数据缩减列裁剪(只选必要列,避免
*)和行过滤(WHERE限定范围);多表连接时先过滤再连接,避免过滤条件放在ON子句4连接条件优化连接键应唯一,优先等值连接,避免不等值连接(<>)和模糊匹配(LIKE),减少链式连接5善用分析函数充分利用
SUM()、AVG()、MIN()/MAX()、RANK()
等分析函数,提升查询效率和准确性01FROM<left_table>02ON<join_condition>03JOIN<right_table>04WHERE<where_condition>05GROUPBY<group_by_list>06WITH<CUBE|ROLLUP>07HAVING<having_condition>08SELECT<select_list>09DISTINCT10ORDERBY<order_by_list>→→→→→→→→→SQL列裁剪与谓词下推SELECTa.movieID,b.movieNameFROM(SELECTmovieID,ratingFROMt_ratingWHERErating>3)aJOINt_moviebONa.movieID=b.movieID先过滤再连接:避免过滤条件放在ON子句,减少参与连接的数据量PythonBroadcastJoin广播小表frompyspark.sql.functionsimportbroadcastresult=large_df.join(broadcast(small_df),"key")result.explain()#确认执行计划出现BroadcastHashJoini默认阈值
spar
温馨提示
- 1. 本站所有资源如无特殊说明,都需要本地电脑安装OFFICE2007和PDF阅读器。图纸软件为CAD,CAXA,PROE,UG,SolidWorks等.压缩文件请下载最新的WinRAR软件解压。
- 2. 本站的文档不包含任何第三方提供的附件图纸等,如果需要附件,请联系上传者。文件的所有权益归上传用户所有。
- 3. 本站RAR压缩包中若带图纸,网页内容里面会有图纸预览,若没有图纸预览就没有图纸。
- 4. 未经权益所有人同意不得将文件中的内容挪作商业或盈利用途。
- 5. 人人文库网仅提供信息存储空间,仅对用户上传内容的表现方式做保护处理,对用户上传分享的文档内容本身不做任何修改或编辑,并不能对任何下载内容负责。
- 6. 下载文件中如有侵权或不适当内容,请与我们联系,我们立即纠正。
- 7. 本站不保证下载资源的准确性、安全性和完整性, 同时也不承担用户因使用这些下载资源对自己和他人造成任何形式的伤害或损失。
最新文档
- 组合创新管理练习题及答案
- 重庆高性价比二手安卓手机去哪选购比较好
- 2026年幼儿美术色彩与构图能力测试
- 2026年学前儿童空间感知能力测试卷
- 2026年健康生活方式与疾病预防测试
- 2026年环境监测与环境保护实践试卷
- 2026年福建省部编版初中英语下册第11单元写作专项训练
- 2026年教师招聘教育法规知识点巩固习题
- 2026年安全技术管理专项训练题库
- 2026食品安全知识挑战赛题库
- 2026-2027学年第一学期学校1530安全教育记录
- 2026年北师大版小学六年级数学上册课时《数学建模》教案
- 2026译林版九年级英语上册暑假预习:Unit1 Know yourself 导学案(知识点+语法+重点短语)
- 2026秋小学英语外研版(三起)(孙有中)(新教材) 四年级上册教学计划附教学进度表
- 道路施工组织技术方案
- 2026年高考生物(湖北卷)真题详细解读及评析
- 2026新版神经内科考试题库(含完整答案+解析)
- 2026年江苏公务员申论高分范文
- 新疆留疆战士试题及答案
- 2025年闽侯县公安局招聘警务辅助人员真题
- JJF 2216-2025 电磁流量计在线校准规范
评论
0/150
提交评论