数据基础离线 2_第1页
数据基础离线 2_第2页
数据基础离线 2_第3页
数据基础离线 2_第4页
数据基础离线 2_第5页
已阅读5页,还剩46页未读 继续免费阅读

下载本文档

版权说明:本文档由用户提供并上传,收益归属内容提供方,若内容存在侵权,请进行举报或认领

文档简介

任务2SparkCore离线计算——电影评分数据集实战内容导航01任务概述与数据初探了解任务目标,使用pandas初探数据集02RDD核心操作实战离线计算电影总数、评分统计信息03进阶计算与文件读写跨文件平均值计算与PySpark文件读写04知识链接与总结RDD基础、键值对RDD、广播与累加器任务概述与数据初探01任务2整体目标与实施路径概述任务2的整体目标与实施路径,建立全局认知运用SparkCore对电影评分数据集进行离线计算整体目标核心能力条件计算按条件计算平均值简单统计数据汇总与统计文件读写数据持久化操作实施路径初探数据使用pandas初步探查

3个数据文件工单实践基于RDD通过

3个工单

实现离线计算文件读写掌握PySpark文件读写方法数据集构成u.user观众信息u.item电影信息u.data评分数据工单2.1:基本信息与目标电影评分数据集初探工单基本信息工单编号2.1工单名称电影评分数据集初探建议学时2学时所属任务SparkCore离线计算环境要求部署好的onYARN模式的Spark集群学习目标知识目标掌握常用Python数据分析函数;掌握PySpark交互界面基本使用方法技能目标能够读取并展示多种格式的数据;能够完成复杂的数据分组与排序分析素养目标培养脚踏实地的工作作风;培养交互式数据分析的思维习惯在PySpark交互界面中,完成u.user、u.item、u.data三个文件的读取、显示,以及多文件连接、分组、排序等操作1安装pandas当前环境未安装pandas,运行命令自动安装:[root@master~]#pip3installpandas2上传数据文件将以下3个文件通过Xftp上传至master节点的

/opt/example

目录:文件名说明u.user用户数据u.item电影信息u.data评分数据3启动集群在master节点依次执行:[root@master~]#start-all.sh#启动Hadoop集群[root@master~]#start-master.sh#启动Spark主节点[root@master~]#start-workers.sh#启动Spark工作节点启动完成后,运行

pyspark

进入PySpark交互界面。工单2.1:环境准备与数据上传读取u.user文件文件每行数据以

|

分隔,运行如下代码:>>>importpandasaspd#定义users列名>>>unames=['uid','age','gender','occupation','zip']#按'|'进行分隔>>>users=pd.read_table('/opt/example/u.user',sep='|',...header=None,names=unames)#显示前5条数据>>>print(users.head(5))read_table()

函数说明参数说明第一个参数指定要读取的

文件路径

或URLsep指定分隔符,默认为制表符header指定是否将文件第一行作为列名names指定DataFrame对象的

列名工单2.1:读取u.user文件工单2.1:连接数据并分组计算平均值任务目标连接观众与评分数据,先按

年龄

再按

性别

分组,计算电影评分平均值代码实现importpandasaspd#定义users列名unames=['uid','age','gender','occupation','zip']users=pd.read_table('/opt/example/u.user',sep='|',header=None,names=unames)#定义ratings列名rnames=['uid','mid','rating','timestamp']ratings=pd.read_table('/opt/example/u.data',sep='\t',header=None,names=rnames)#基于观众数据和评分数据连接两个DataFrameframe=pd.merge(ratings,users)#先按年龄分组再按性别分组,最后求电影评分值的平均值print(frame['rating'].groupBy([frame['age'].apply(round,args=[-1]),frame['gender']]).mean())核心APIAPI作用关键参数merge()按条件合并两个DataFrame两个待合并的DataFramegroupBy()按列值分组聚合分组依据的列名apply()批量应用指定函数args传递函数参数mean()计算数值型数据的平均值—工单2.1:按性别计算平均评分并排序步骤6·Pandas目标三表连接后,按性别计算每部电影平均评分,并降序排列task2_1.pyimportpandasaspd#定义三表列名unames=['uid','age','gender','occupation','zip']rnames=['uid','mid','rating','timestamp']mnames=['mid','title','d1','d2','url',

'C1','C2','C3','C4','C5','C6','C7','C8','C9',

'C10','C11','C12','C13','C14','C15','C16','C17',

'C18','C19']#读取三表数据users=pd.read_table('/opt/example/u.user',sep='|',header=None,names=unames)ratings=pd.read_table('/opt/example/u.data',sep='\t',header=None,names=rnames)movies=pd.read_table('/opt/example/u.item',sep='|',header=None,names=mnames,encoding='ISO-8859-1')#三表合并→按性别+电影分组求均值→降序排列frame=pd.merge(pd.merge(ratings,users),movies)frame['rating'].groupBy([frame['gender'],frame['mid']]).mean().sort_values(ascending=False)知识链接sort_values()参数说明ascending=True升序排列(默认)ascending=False降序排列工单2.1小结:实践是检验真理的唯一标准工单小结学生通过与

PySpark

交互的方式,针对电影评分数据集完成了以下操作:文件读取连接分组排序求平均值对数据集结构有了一定认识掌握了使用

PySpark进行数据分析的基本函数素养课堂实践是检验真理的唯一标准“01这一观点源于马克思关于真理标准的论述。人的思维是否具有客观的真理性,不是一个纯理论的问题,而是一个实践的问题。021978年5月11日《光明日报》刊登该文,引发全国关于真理标准问题的讨论。03通过实践,人们将理论应用于实际,观察其效果,从而判断是否正确,逐步接近真理。RDD核心操作实战02工单2.2:基本信息与目标12借助RDD的

map()、count()、filter()、countByValue()

等函数完成电影总数统计与年份分组计数。工单基本信息工单编号2.2工单名称离线计算电影总数和不同年份的电影数量建议学时2

课时所属任务SparkCore离线计算环境要求部署好的onYARN模式的Spark集群知识目标掌握分组统计与聚合计算的方法掌握RDD创建与转换操作技能目标能够正确创建与操作RDD完成数据加载能够实现基础统计与计数计算素养目标培养积极探索、勇于创新的科学素养培养与人交流、合作的能力工单2.2:环境准备与计算电影总数1-2启动集群并上传文件bash#启动HDFS与YARNstart-dfs.shstart-yarn.sh#启动Spark主节点与工作节点start-master.shstart-workers.sh#上传数据文件至HDFShdfsdfs-put/opt/example/u.item\/user/zjaf/data3-4编写movie_total.py计算电影总数python核心代码逻辑:读取HDFS电影数据→创建RDD→调用

count()

统计总数frompysparkimportSparkContextif__name__=='__main__':sc=SparkContext(master="spark://master:7077",appName="movie_total")movie_rdd=sc.textFile(

"hdfs://master:9000/user/zjaf/data/u.item")

print("Totalnumberis:%s"%movie_rdd.count())sc.stop()#提交运行命令spark-submit--masterspark://master:7077\--deploy-modeclient--executor-memory1g\--total-executor-cores2movie_total.py🎯

目标提取每条记录的年份,统计不同年份的电影数量📌

数据特征:年份位于第二列右侧,被圆括号括起,靠近分隔符

|📋

数据格式示例记录第二列内容提取年份1ToyStory(1995)19952GoldenEye(1995)19953FourRooms(1995)1995🔧

提取逻辑①

|

分割→②

取第2个元素→③

倒数第5位截取4字符💻

代码实现:movie_total2.pyimportsysfrompysparkimportSparkContext#输入电影的年份,若非法则返回-1

def

convert_year(x):try:returnint(x[-4:])except:return-1

if__name__=="__main__":iflen(sys.argv)!=2:print("Usage:hdfs_wordcount.py<directory>",file=sys.stderr)sys.exit(-1)#创建SparkContext对象sc=SparkContext(master="spark://master:7077",appName="movie_total2")工单2.2:计算不同年份的电影数量(上)步骤5:实现离线计算不同年份的电影数量——先分析数据格式,再搭建程序框架工单2.2:计算不同年份的电影数量(下)核心任务完成

movie_total2.py

后半部分,实现电影年份统计RDD转换操作链1textFile()从HDFS读取文件,构建电影RDD2map()以"|"分割记录,再取年份字段并由convert_year()转换,非法年份置为-13filter()过滤掉非法记录(x!=-1),打印过滤后的记录数4countByValue()自动统计每个Key下的Value个数,以字典格式返回,打印各年份电影数量步骤6提交运行运行spark-submit提交至集群,注意带上参数(待处理文件的路径)spark-submitmovie_total2.py<待处理文件路径>movie_total2.py后半部分#从HDFS读取文件movie_rdd=sc.textFile(sys.argv[1])#以'|'为分隔符分割记录,返回电影RDDmovie_fields=movie_rdd.map(

lambdalines:lines.split("|"))#生成新的电影年份RDD,并将非法的年份置为-1years=movie_fields.map(

lambdafields:fields[2]).map(

lambdax:convert_year(x))#过滤掉电影年份RDD中非法的记录years_filtered=years.filter(lambdax:x!=-1)print("numberafterfilteris:%s"%years_filtered.count())#统计不同年份的电影数量movie_ages=years_filtered.countByValue()#countByValue()是Spark提供的便捷函数,能自动统计#每个Key下面的Value个数,并以字典格式返回fork,vinmovie_ages.items():

print(k,v)sc.stop()工单2.2小结:团结协作,攻坚克难总结工单2.2学习成果,结合红旗渠精神强调团结协作的重要性工单小结RDD函数操作实战回顾通过以下函数操作,实现电影数据的统计分析:map()count()countByValue()filter()实现成果电影总数统计不同年份电影数量分布素养课堂·红旗渠精神一组数据,读懂「工程奇迹」1500

km人工天河全长1250

座削平山头10

年艰苦修筑“自力更生、艰苦创业、

团结协作

、无私奉献”RDD设计理念与红旗渠精神的契合面对大型任务时,科学的协作会大幅提高效率。工单2.3:基本信息与目标工单编号

2.3知知识目标掌握RDD高级转换与行动操作的方法;掌握常用的统计函数技技能目标能够使用

collect()、reduce()、map()、sortByKey()、groupByKey()

等函数素素养目标培养积极探索,勇于创新的科学素养▶本工单将离线计算电影评分数据集中常用的统计信息,包括最高评分最低评分平均评分中位评分,以及每个观众的平均评分值。工单编号2.3工单名称离线计算电影评分统计信息建议学时2

课时所属任务SparkCore离线计算环境要求部署好的onYARN模式的Spark集群基本信息

BasicInformation工单2.3:环境准备与初步探索u.data步骤1-3步骤1-2启动集群并上传文件start-dfs.sh

#启动HDFSstart-yarn.sh

#启动YARNstart-master.sh

#启动Spark主节点start-workers.sh

#启动Spark工作节点hdfsdfs-put/opt/example/u.data/user/zjaf/data步骤3初步探索u.data文件—在master的/opt/example目录下新建movie_rating.pyfrompysparkimportSparkContextif__name__=='__main__':sc=SparkContext(master="spark://master:7077",appName="movie_rating")rating_rdd=sc.textFile("hdfs://master:9000/user/zjaf/data/u.data")print("Totalnumberis:%s"%rating_rdd.count())sc.stop()通过

spark-submit

提交至集群运行,输出

Totalnumber

即数据集总行数计算最高评分、最低评分、平均评分、中位评分Python/PySparkimportsysimportnumpyasnpfrompysparkimportSparkContextif__name__=='__main__':

iflen(sys.argv)!=2:print("Usage:hdfs_wordcount.py<directory>",file=sys.stderr)sys.exit(-1)sc=SparkContext(master="spark://master:7077",appName="movie_rating")rating_rdd=sc.textFile(sys.argv[1])

#获取评分RDD,分隔符为制表符rating_data=rating_rdd.map(

lambdaline:line.split("\t"))ratings=rating_data.map(

lambdafields:int(fields[2]))数据示例(前3条记录)观众ID电影ID评分值时间19624238812509491863023891717742223771878887116📍取每条记录第

3

个数即为评分值📊目标指标:最高/最低/平均/中位

评分工单2.3:计算最高/最低/平均/中位评分步骤4工单2.3:reduce计算统计指标核心计算1使用

reduce

计算最高/最低评分2使用

reduce

计算平均评分3使用

numpy

计算中位评分4输出四项统计结果:minimalmaximalaveragemedianminimalmaximalaveragemedianPython·统计指标计算#计算最高/最低评分max_rating=ratings.reduce(lambdax,y:max(x,y))min_rating=ratings.reduce(lambdax,y:min(x,y))#计算平均/中位评分mean_rating=ratings.reduce(

lambdax,y:x+y)/float(rating_rdd.count())median_rating=np.median(ratings.collect())print("minimal:%d"%min_rating)print("maximal:%d"%max_rating)print("average:%2.2f"%mean_rating)print("median:%d"%median_rating)sc.stop()最低评分最高评分平均评分中位评分步骤5:计算每个观众的平均评分值movie_rating3.pyPYTHON123456789101112131415161718192021importsysfrompysparkimportSparkContextdef

avg_rating(line):rating_list=list(line[1])temp=0foriteminrating_list:temp+=float(item)temp=temp/len(rating_list)

return(int(line[0]),temp)if__name__=='__main__':

if

len(sys.argv)!=2:

print("Usage:hdfs_wordcount.py<directory>",file=sys.stderr)sys.exit(-1)sc=SparkContext(master="spark://master:7077",appName="movie_calc")rating_rdd=sc.textFile(sys.argv[1])

#分词rating_rdd2=rating_rdd.map(

lambdaline:line.split("\t"))新建Python脚本在

/opt/example

目录下新建

movie_rating3.py本步骤程序要点1定义avg_rating函数累计评分并除以条数,返回(观众ID,平均评分)2初始化SparkContext并读取数据连接spark://master:7077,textFile读入后按制表符分词工单2.3:avg_rating函数与程序框架工单2.3:分组计算与排序输出核心操作链map取第1、3列→groupByKey分组→avg_rating计算平均值→sortByKey排序→collect输出user_avg_rating.py#取第1、3列,并按第1列分组user_ratings_grouped=rating_rdd2.map(

lambdafields:(int(fields[0]),int(fields[2]))).groupByKey()#计算每个用户的平均评分值user_ratings_avg=user_ratings_grouped.map(

lambdax:avg_rating(x))#按用户ID升序排序user_ratings_avg_sorted=user_ratings_avg.sortByKey()results=user_ratings_avg_sorted.collect()forresultinresults:

print(result)sc.stop()步骤6提交运行与验证1通过

spark-submit

提交运行2浏览器访问

:8080

查看

SparkWebUI3单击

Application

链接查看详细运行数据和日志SparkRDD实践工单2.3小结:自主可控,创新求变工单技术成果reduce()map()collect()sortByKey()groupByKey()median()通过上述函数,完成电影评分数据集离线计算:最高评分、最低评分平均评分、中位评分每个观众的平均评分值素养课堂—国产芯片的自主创新中星微技术股份有限公司"星光摩尔一号"人工智能芯片十五大核心技术,3000+

项发明专利"星光"系列第一代超大规模集成电路芯片自主知识产权,自主创新突破性进展持之以恒筑牢基础,不断创新迸发灵感自主可控自主可控是科技发展的目标之一。肯下功夫、练得一身硬本事,才能创造更多原创性产品。进阶计算与文件读写03工单2.4:基本信息与目标基本信息工单编号2.4工单名称离线计算跨文件平均值建议学时4

课时所属任务SparkCore离线计算环境要求部署好的onYARN模式的Spark集群工单任务通过RDD实现从多个文件读取数据,利用

join

连接操作完成分析计算:1计算每部电影的平均评分值2最受男性喜爱的电影

Top103最受女性喜爱的电影

Top10学习目标知识目标掌握RDD连接的方法;掌握RDD组合的方法技能目标能够连接多个RDD;能够组合多个RDD素养目标培养积极探索,勇于创新的科学素养工单2.4:计算每部电影平均评分(上)步骤1·启动集群命令作用start-dfs.sh启动HDFSstart-yarn.sh启动YARNstart-master.sh启动Spark主节点start-workers.sh启动Spark工作节点步骤2·计算每部电影平均评分创建

/opt/example/movie_rating4.pyfrompysparkimportSparkContextif__name__=='__main__':sc=SparkContext(master="spark://master:7077",appName="movie_rating_count")

#读取数据rrdd=sc.textFile("/user/zjaf/data/u.data")

#分隔字段rrdd2=rrdd.map(lambdaline:line.split("\t"))

#统计每部电影评分次数rrdd3_1=rrdd2.map(

lambdafields:(int(fields[1]),1)).reduceByKey(lambdax,y:x+y)

#统计每部电影评分总和rrdd3_2=rrdd2.map(

lambdafields:(int(fields[1]),int(fields[2]))).reduceByKey(lambdax,y:x+y)工单2.4:计算每部电影平均评分(下)join连接电影名称→collect转换为List→总和

÷

次数得到平均评分💡核心思路1读取

u.data,分别计算评分次数

rrdd3_1

与评分总和

rrdd3_22读取

u.item,取电影ID和名称

mrdd33通过

join

将次数、总和分别与电影名称连接,按电影ID升序排列,并用

collect

转换为List4评分值总和

÷

评分次数,即得每部电影的平均评分值#读取u.item文件mrdd=sc.textFile("/user/zjaf/data/u.item")#按"|"分隔记录mrdd2=mrdd.map(lambdaline:line.split("|"))#取电影ID和名称mrdd3=mrdd2.map(

lambdafields:(int(fields[0]),fields[1]))#分别将被评分的总次数、评分值的总和#与取得的电影ID和名称连接,并按电影ID升序排列rdd1=mrdd3.join(rrdd3_1).sortByKey()rdd2=mrdd3.join(rrdd3_2).sortByKey()rdd3=rdd2.map(lambdafields:fields[1])rdd4=rdd1.map(lambdafields:fields[1])#将RDD转换为Listzonghe=rdd3.map(

lambdafields:fields[1]).collect()cishu=rdd4.map(

lambdafields:fields[1]).collect()#计算每部电影的平均评分值avg=[x/yforx,yinzip(zonghe,cishu)]print(avg)sc.stop()代码实现步骤3·计算最受男性/女性喜爱的电影TOP10工单2.4:三表关联与性别过滤movie_gender_top10·前半段frompysparkimportSparkContextif__name__=='__main__':sc=SparkContext(master="spark://master:7077",appName="movie_rating_analysis")sc.setLogLevel("WARN")

#读取三个数据文件rating_data=sc.textFile("/user/zjaf/data/u.data")user_data=sc.textFile("/user/zjaf/data/u.user")movie_data=sc.textFile("/user/zjaf/data/u.item")

#解析数据rating_records=rating_data.map(lambdaline:line.split("\t"))user_records=user_data.map(lambdaline:line.split("|"))movie_records=movie_data.map(lambdaline:line.split("|"))

#准备用户性别信息(用户ID→性别)user_gender=user_records.map(

lambdafields:(int(fields[0]),fields[2]))

#准备评分信息(用户ID→(电影ID,评分))movie_ratings=rating_records.map(

lambdafields:(int(fields[0]),(int(fields[1]),float(fields[2]))))

#关联评分和性别信息gender_ratings=movie_ratings.join(user_gender)

#按性别过滤并重组数据male_ratings=gender_ratings.filter(

lambdax:x[1][1]=="M").map(

lambdax:(x[1][0][0],x[1][0][1]))female_ratings=gender_ratings.filter(

lambdax:x[1][1]=="F").map(

lambdax:(x[1][0][0],x[1][0][1]))

#准备电影信息movie_info=movie_records.map(

lambdax:(int(x[0]),x[1]))——待续——3任务说明(第一部分·前半段代码)1读取三个数据文件加载

u.data(评分)、u.user(用户)、u.item(电影)三个数据集2解析数据按分隔符拆分字段:u.data用制表符,u.user/u.item用「|」3关联评分与性别信息以用户ID为键,将评分记录与用户性别join关联4按性别过滤分别筛出男性(M)与女性(F)评分,重组为(电影ID,评分)以备后续分组统计Top10u.data评分+u.user性别+u.item电影→按性别统计Top10步骤3后半部分:定义calculate_top_movies函数·聚合评分求平均·关联电影名·排序取Top10并输出计算Top10函数1聚合评分→2求平均→3关联电影名→4排序取Top10task_2_4.py—calculate_top_movies()def

calculate_top_movies(ratings_rdd,movie_info_rdd):

returnratings_rdd.map(

lambdax:(x[0],(x[1],1))).reduceByKey(

lambdax,y:(x[0]+y[0],x[1]+y[1])).map(

lambdax:(x[0],x[1][0]/x[1][1])).join(movie_info_rdd).map(

lambdax:(x[1][0],x[1][1])).sortByKey(ascending=False).take(10)男性/女性Top10输出#输出男性Top10male_top10=calculate_top_movies(male_ratings,movie_info)print("\n最受男性欢迎的Top10电影:")fori,(avg_rating,title)in

enumerate(male_top10,1):

print(f"{i}.{title}:{avg_rating:.2f}")#输出女性Top10female_top10=calculate_top_movies(female_ratings,movie_info)print("\n最受女性欢迎的Top10电影:")fori,(avg_rating,title)in

enumerate(female_top10,1):

print(f"{i}.{title}:{avg_rating:.2f}")sc.stop()工单2.4:计算Top10并输出结果29工单2.4小结:自力更生,勇于探索工单小结利用

join

实现RDD连接,完成三项核心任务:1计算每部电影的平均评分值2最受男性喜爱的电影

Top103最受女性喜爱的电影

Top10素养课堂—中国超算的自主之路20世纪80年代,面对外部技术壁垒,我国科研工作者自力更生,逐步建立高性能计算体系起步期"银河"系列超级计算机发展期"曙光"超级计算机跨越期"天河"超级计算机领先期"神威·太湖之光"我们在学习过程中要重视自主研发能力与创新能力的培养,共同推动国家科技事业迈向新高度。工单2.5:基本信息与目标基本信息工单编号2.5工单名称PySpark文件读写建议学时2课时所属任务SparkCore离线计算环境要求部署好的onYARN模式的Spark集群工单简介文件读写是实现数据离线计算的重要步骤。本工单将通过

PySpark交互界面实现从本地文件系统和HDFS读写文本文件,并编写

Python程序读取和解析存储在HDFS中的JSON文件。学习目标知识目标掌握文本文件和JSON文件在Spark中的解析方法技能目标能够读写本地文件系统和HDFS中的文件;能够读写文本文件和JSON文件素养目标培养积极探索、勇于创新的科学素养;培养一丝不苟、规范操作的能力在master节点创建词频统计测试文件,随后启动集群并进入PySpark1准备测试文件master·/opt/exampletest2.txt满江红【作者】岳飞

【朝代】宋怒发冲冠,凭栏处、潇潇雨歇。抬望眼、仰天长啸,壮怀激烈。三十功名尘与土,八千里路云和月。莫等闲、白了少年头,空悲切。靖康耻,犹未雪。臣子恨,何时灭。驾长车踏破,贺兰山缺。壮志饥餐胡虏肉,笑谈渴饮匈奴血。待从头、收拾旧山河,朝天阙。注意:如遇中文乱码,运行

localectlset-localeLANG=zh_CN

解决2启动集群并进入PySpark1start-dfs.sh启动HDFS2start-yarn.sh启动YARN3start-master.sh启动Spark主节点4start-workers.sh启动Spark工作节点5pyspark--masterspark://master:7077工单2.5:准备测试文件与启动集群工单2.5:本地文件系统读写文本文件步骤3-4:读取与写入本地文本文件3读取本地文本文件在PySpark交互界面运行:>>>rdd=sc.textFile("file:///opt/example/test2.txt")>>>rdd.take(3)#显示前3段知识链接▸读取本地文件必须采用以

"file:///"

开头的文件路径▸Spark采用

惰性机制,需要行动操作触发文件读取4写入本地文本文件在PySpark交互界面运行:>>>rdd.saveAsTextFile("file:///opt/example/backup")知识链接▸saveAsTextFile()

是一个

行动操作▸新建终端窗口,运行

cd/opt/example

查看保存结果。Spark将结果写入backup目录,生成part文件。▸

终端验证:cd/opt/example→backup目录→part文件工单2.5:HDFS文件系统读写文本文件读取与写入HDFS文本文件(步骤5—6)5读取HDFS文本文件上传文件至HDFShdfsdfs-put/opt/example/test2.txt

/user/zjaf/datahdfsdfs-ls/user/zjaf/dataPySpark读取并显示第一行>>>rdd2=sc.textFile(...

"hdfs://master:9000/user/zjaf/data/test2.txt")>>>rdd2.first()6写入HDFS文本文件保存RDD至HDFS>>>rdd2.saveAsTextFile("/backup")验证结果Spark自动创建目录

/user/root/backup

并写入结果文件工单2.5:JSON文件读取与解析7JSON文件读取在

/opt/example

下新建

passport.json:{"Name":"Zhangsan","Age":30,"Nation":"HANZU"}{"Name":"Lisi","Age":28,"Nation":"TUJIAZU"}{"Name":"Wangwu","Age":32,"Nation":"CHAOXIANZU"}使用

textFile()

读取并显示原始文本内容:>>>data=sc.textFile("/user/zjaf/data/passport.json")8-9解析JSON文件在

/opt/example

下新建

json_parse.py:frompysparkimportSparkConf,SparkContextimportjsonconf=SparkConf().setMaster("spark://:7077").setAppName("parsejson")sc=SparkContext(conf=conf)sc.setLogLevel("WARN")inputFile="/user/zjaf/data/passport.json"jsonStrs=sc.textFile(inputFile)result=jsonStrs.map(lambdas:json.loads(s))result.foreach(print)通过

spark-submit

提交运行:$spark-submitjson_parse.py✓浏览器访问

:8080

查看

SparkWebUI

验证结果✓单击

parsejson

的Application链接查看详细日志工单2.5小结工单2.5小结:腹有诗书气自华工单小结通过PySpark交互界面完成:1本地文件系统

HDFS

文本文件读写2Python程序

读取解析HDFS中的JSON文件❝

春眠不觉晓,床前明月光,明月几时有

❞素养课堂中华诗词与文化自信中国是诗词的国度。《诗经》《楚辞》《乐府诗集》等无数名篇,造就无比灿烂的中华诗词文化。诗词

“随风潜入夜,润物细无声”——给心灵以美的熏陶,给生命以丰厚的馈赠,给人生以深沉的激励。闲暇时观看节目,增强文化素养和文化自信:《中国诗词大会》《中华好诗词》《诗意中国》腹有诗书气自华志气·骨气·底气知识链接与总结04RDD产生的背景对比传统框架的不足,说明RDD的设计目标RDD改进RDD为应对大规模数据处理挑战而生,满足迭代式算法与交互式数据挖掘场景需求✕传统框架的局限中间结果写入

HDFS,产生大量数据复制、磁盘I/O与序列化开销迭代式算法(如

ALS

交替最小二乘法、凸优化梯度下降)性能差难以满足交互式数据挖掘场景需求✓RDD的设计目标实现数据操作

管道化应用逻辑表达为一系列转换操作RDD转换形成依赖关系,避免中间结果存储显著减少数据复制、磁盘I/O和序列化开销RDD是Spark的核心抽象概念,代表一个

不可变、可分区

元素可并行计算

的集合。RResilient弹性数据可保存在内存或磁盘中,具有

容错能力,可通过依赖关系和DAG重新构建。DDistributed分布式数据分布式存储,便于分布式计算,至少被分为

一个分区。DDataset数据集由记录组成的数据集,各分区包含不同的记录子集,可

独立分析。RDD的概念RDD的类型体系类型名称数据结构特征典型应用场景特殊操作方法PairRDD(K,V)键值对结构聚合操作、数据关联keys()join()groupByKey()mapValues()sortByKey()DoubleRDD纯数值型RDD统计分析、机器学习特征处理mean()variance()histogram()stdev()DataFrameSchemaRDD带Schema的二维表结构SQL查询、结构化数据分析select()filter()groupBy()agg()HadoopRDD自定义结构原生Hadoop格式读取、兼容旧Hadoop生态saveAsHadoopFile()saveAsNewAPIHadoopDataset()RDD的依赖关系:窄依赖与宽依赖RDD的依赖关系决定数据在集群中的流动与计算方式,分为

窄依赖

宽依赖

两类。窄依赖无Shuffle依赖模式一对一/多对一核心特点无

shuffle,分区可并行处理计算效率高,各节点独立执行典型函数map()filter()union()mapPartitions()mapValues()宽依赖引发Shuffle依赖模式多对多核心特点引发

shuffle,跨节点传输同步计算效率较低,涉及全局数据重分配典型函数groupByKey()partitionBy()reduceByKey()sortByKey()创建RDD的三种方式1读取外部数据集从本地文件加载,或从HDFS、HBase、Cassandra等外部数据源加载;支持文本文件、JSON文件和符合HadoopInputFormat格式的文件。rdd=sc.textFile("file:///opt/example/test2.txt")2parallelize()并行化集合调用SparkContext的parallelize()函数,将驱动程序中的现有集合加载到并行化RDD中。data=[1,2,3,4,5,6,7,8,9,10,11,12]rdd=spark.sparkContext.parallelize(data)3创建空RDDrdd=spark.sparkContext.emptyRDD两类操作操作类型特点执行方式转换操作操作RDD并返回新RDD惰性执行行动操作在数据集上运算,返回计算值触发实际计算转换操作是惰性的,只有在调用行动操作之后才会真正执行。核心区别核心转换操作速查函数功能说明典型场景filter(func)筛选满足条件的元素,返回新数据集数据清洗,剔除无效记录map(func)对每个元素应用func,返回新数据集数据格式转换,如解析JSONflatMap(func)类似map,但每个元素可映射到

0个或多个

输出文本分词、行转列groupByKey()按Key分组,返回(K,Iterable<V>)需保留全部分组值的场景reduceByKey(func)按Key聚合计算,返回(K,V)求和、计数等聚合操作PySpark代码实践准备数据words=sc.parallelize(["scala","java","hadoop","spark","akka","sparkvshadoop","pyspark","pysparkandspark"])filter—过滤含"spark"的元素words.filter(lambdax:'spark'inx).collect()#['spark','sparkvshadoop','pyspark','pysparkandspark']map—转为(元素,1)键值对words.map(lambdax:(x,1)).collect()#[('scala',1),('java',1),('hadoop',1),...]RDD转换操作速查RDD转换操作代码实践(续)1flatMap拆分元素>>>words.flatMap(lambdax:(x,1)).collect()输出['scala',1,'java',1,'hadoop',1,...]2groupByKey按首字母分组>>>words.map(lambdax:(x[0],x)).groupByKey().collect()输出's':['scala','spark','sparkvshadoop']3reduceByKey统计首字母出现次数>>>words.map(lambdax:(x[0],1)).reduceByKey(lambdaa,b:a+b).collect()输出[('s',3),('j',1),('h',1),('a',1),('p',2)]RDD行动操作详解(上)行动操作是计算真正被触发的位置。Spark程序执行到行动操作时,才会执行真正的计算。常见行动操作函数count()返回数据集中的元素个数collect()以数组的形式返回数据集中的所有元素first()返回数据

温馨提示

  • 1. 本站所有资源如无特殊说明,都需要本地电脑安装OFFICE2007和PDF阅读器。图纸软件为CAD,CAXA,PROE,UG,SolidWorks等.压缩文件请下载最新的WinRAR软件解压。
  • 2. 本站的文档不包含任何第三方提供的附件图纸等,如果需要附件,请联系上传者。文件的所有权益归上传用户所有。
  • 3. 本站RAR压缩包中若带图纸,网页内容里面会有图纸预览,若没有图纸预览就没有图纸。
  • 4. 未经权益所有人同意不得将文件中的内容挪作商业或盈利用途。
  • 5. 人人文库网仅提供信息存储空间,仅对用户上传内容的表现方式做保护处理,对用户上传分享的文档内容本身不做任何修改或编辑,并不能对任何下载内容负责。
  • 6. 下载文件中如有侵权或不适当内容,请与我们联系,我们立即纠正。
  • 7. 本站不保证下载资源的准确性、安全性和完整性, 同时也不承担用户因使用这些下载资源对自己和他人造成任何形式的伤害或损失。

评论

0/150

提交评论