版权说明:本文档由用户提供并上传,收益归属内容提供方,若内容存在侵权,请进行举报或认领
文档简介
任务5实时计算实现内容导航01任务概述介绍实时计算任务背景与目标02工单实战四个工单逐步掌握SparkStreaming与StructuredStreaming03知识链接系统讲解实时计算基础与框架原理04任务总结回顾本任务核心知识点任务概述01基于StructuredStreaming的电影评分数据集实时计算01StructuredStreaming
是Spark2.0推出的实时流处理技术,相较于SparkStreaming具有
更低的延迟。02本任务将编写Python程序,在任务4搭建好的实时计算环境基础上,实现
电影评分数据集
的实时计算。四工单实施路径工单1SparkStreaming单词计数实时计算目标理解实时计算基本原理1工单2StructuredStreaming单词计数实时计算目标对比两种框架实现方式2工单3SparkStreaming电影评分数据集实时计算目标实战案例入门3工单4StructuredStreaming电影评分数据集实时计算目标深入掌握编程技巧4任务描述工单实战02基本信息工单编号5.1工单名称SparkStreaming实时单词计数建议学时2
学时所属任务实时计算实现环境要求已搭建好的实时计算环境工单说明启动Hadoop集群和Spark集群,编写Python程序使用SparkStreaming实时监听HDFS指定目录变动,对新文件进行单词计数,通过
spark-submit
提交至集群运行。工单目标知识目标1.掌握SparkStreaming的执行机制2.掌握实时监听HDFS目录变动,对新增文件进行单词计数技能目标能使用Python编写SparkStreaming程序监听HDFS指定目录变动素养目标1.培养从本地开发到集群部署的规范意识2.培养解决实际问题的能力工单5.1基本信息与目标工单5.1启动集群与编写说明1启动Hadoop和Spark集群在
master
节点执行:$start-dfs.sh#启动HDFS$start-yarn.sh#启动YARN$start-master.sh#启动Spark主节点$start-workers.sh#启动Spark工作节点运行
jps
验证进程master:NameNodeSecondaryNameNodeMasterResourceManagerslave01/slave02:DataNodeNodeManagerWorker2编写单词计数Python程序D监听目录/user/zjaf/dataT查询方式每隔
5s查询新增文件F创建文件/opt/example/hdfs_streaming.py单词计数程序代码📄/opt/example/hdfs_streaming.py⟡hdfs_streaming.pyPythonimportsysfrompysparkimportSparkContextfrompyspark.streamingimportStreamingContextif__name__=="__main__":iflen(sys.argv)!=2:print("Usage:hdfs_streaming.py<directory>",file=sys.stderr);sys.exit(-1)sc=SparkContext(master="spark://master:7077",appName="hdfs_streaming")sc.setLogLevel("ERROR")#设置日志级别ssc=StreamingContext(sc,5)#创建StreamingContext,间隔5slines=ssc.textFileStream(sys.argv[1])#读取新增文件#单词计数counts=lines.flatMap(lambdaline:line.split(""))\.map(lambdax:(x,1))\.reduceByKey(lambdaa,b:a+b)counts.pprint()#显示结果ssc.start()#启动计算ssc.awaitTermination()#等待执行!💡
在不影响业务逻辑的前提下,优先选用
reduceByKey()
而不是
groupByKey()运行验证与工单小结3master节点运行监听程序在master节点终端运行:spark-submit/opt/example/hdfs_streaming.py"/user/zjaf/data"每隔
5s
产生一个新结果4新建终端上传测试文件新建终端,执行上传命令:hdfsdfs-put/opt/example/test1.txt/user/zjaf/data5返回原终端查看结果查看实时单词计数输出:('love',3)('I',3)('my',2)…等工单小结本工单通过编写Python程序,实现了使用
SparkStreaming
实时监听HDFS指定目录的变化,并对新增文件完成
单词计数。素养课堂:新技术赋能产业新技术正在深刻改变各行各业的生产方式和商业模式智能制造01推动制造业向高端化、智能化、绿色化转型升级02提升产品质量人工智能辅助诊断让诊断更准确、治疗更精细降低医疗费用惠及更多人群技术发展日新月异,在学习和生活中多关注新技术发展趋势,为个人发展指明方向◉基本信息工单编号5.2工单名称SparkStreaming实时计算电影评分数据集建议学时3所属任务实时计算实现环境要求已搭建好的实时计算环境➔任务流程1模拟数据生成编写Python程序,连接u.user和u.data,按
0.5s
间隔产生评分数据并写入文件2流数据接入整合
Flume+HDFS,将模拟数据转换为实时数据流3实时统计编写
SparkStreaming
程序,按性别统计观影次数★学习目标▦知识目标掌握使用SparkStreaming实时监听HDFS文件变动,对新增内容进行词频统计✎技能目标能使用Python编写SparkStreaming程序,实时监听HDFS文件变动❖素养目标培养严谨的实时数据处理验证意识培养解决实际问题的能力工单5.2基本信息与目标实时计算实战▼▼任务准备启动集群与数据说明Step1启动Hadoop和Spark集群在master节点执行以下命令:start-dfs.sh启动HDFSstart-yarn.sh启动YARNstart-master.sh启动Spark主节点start-workers.sh启动Spark工作节点Step2模拟实时评分数据将u.user与u.data连接,提取以下4个字段:u.user+u.data1电影ID2年龄3性别4评分值写入路径/opt/tmp/task5_1.log逐行写入数据,按0.5s间隔模拟实时数据流数据读取与处理Python①12345678910111213141516importpandasaspdimportosimporttime#读取u.userunames=['uid','age','gender','occupation','zip']users=pd.read_table('/opt/example/u.user',sep='|',header=None,names=unames)#读取u.datarnames=['uid','mid','rating','timestamp']ratings=pd.read_table('/opt/example/u.data',sep='\t',header=None,names=rnames)#连接两表frame=pd.merge(ratings,users)#取子集:mid,age,gender,ratingdf=frame[['mid','age','gender','rating']]#若文件存在则先删除fn='/opt/tmp/task5_1.log'ifos.path.exists(fn):逐行写入·模拟实时数据流Python②171819202122232425262728293031os.remove(fn)print('filedeleted')#逐行写入,模拟实时数据流count=0lines=df.shape[0]whilecount<lines:log='{}\t{}\t{}\t{}\n'.format(df.iloc[count,0],df.iloc[count,1],df.iloc[count,2],df.iloc[count,3])with
open(fn,'a+')asfile:file.write(log)time.sleep(0.5)count=count+1模拟评分数据代码1启动ZooKeeper与Kafka集群shell$
zkServer.shstart#启动ZooKeeper集群$
kafka-server-start.sh-daemon/opt/kafka/config/perties#以守护进程启动Kafka2创建Kafka主题
sscreplication-factor1·partitions1·master:2181shell$
kafka-topics.sh--create--zookeepermaster:2181--replication-factor1--partitions1--topic
ssc启动Kafka与创建主题启动ZooKeeper/Kafka集群,并创建Kafka主题编写SparkStreaming计算程序目标读取HDFS文件新增数据,按性别统计观影次数task5_streaming_1.pyimportsysfrompysparkimportSparkContextfrompyspark.streamingimportStreamingContextif__name__=="__main__":
iflen(sys.argv)!=2:
print("Usage:task5_streaming_1.py<directory>",file=sys.stderr)sys.exit(-1)sc=SparkContext(master="spark://master:7077",appName="task5_streaming_1")sc.setLogLevel("ERROR")ssc=StreamingContext(sc,2)
#读取HDFS数据流→取第3列(性别)→分组计数lines=ssc.textFileStream(sys.argv[1])counts=lines.map(lambdaline:line.split("\t")[2])\.map(lambdax:(x,1))\.reduceByKey(lambdaa,b:a+b)counts.pprint()ssc.start()ssc.awaitTermination()计算逻辑1读取HDFS数据流2取第3列(性别)3分组计数测试验证spark-submit/opt/example/task5_streaming_1.py"/record_hdfs_task5"运行结果按性别统计的观影次数'M'20'F'13多终端测试与工单小结8终端2:运行PySpark代码在终端2中打开PySpark交互界面,输入步骤(2)中的代码并保存。9终端3:运行FlumeAgent在终端3中执行以下命令:flume-ngagent-cconf-f/opt/example/task5_flume_1.conf-namea1-Dflume.root.logger=INFO,console报错说明如果报告类似
ERRORFileInputDStream:Filehdfs://master:9000/record_hdfs_task5/FlumeData.1694405047216.tmp
的错误,是写入HDFS的操作还没完成,但Spark监听到文件有变动造成的。执行步骤(8)和步骤(9)后就能看到类似
('M',10)
的结果。工单小结本工单通过编写Python程序实现了使用SparkStreaming实时监听HDFS中指定的文件,并对Flume写入的电影评分数据
按性别分组计算。使用StructuredStreaming实时监听Kafka,并完成实时单词计数5.3工单基本信息工单名称StructuredStreaming实时单词计数建议学时3所属任务实时计算实现环境要求已搭建好的实时计算环境使用
StructuredStreaming
实现实时计算,帮助学生掌握
Spark连接Kafka
的方法。通过编写Python程序,实现StructuredStreaming实时监听Kafka,并对数据进行
单词计数。◎学习目标知知识目标掌握StructuredStreaming的使用方法技技能目标配置Spark环境以支持StructuredStreaming与Kafka集成使用StructuredStreaming实现实时计算素素养目标培养遵循科学发展的规律,制定合理的规划培养解决实际问题的能力工单5.3基本信息与目标启动集群与进程验证1启动Hadoop&Spark集群master节点start-dfs.sh#启动HDFSstart-yarn.sh#启动YARNstart-master.sh#启动Spark主节点start-workers.sh#启动Spark工作节点2启动ZooKeeper&Kafka集群master/slave01/slave02zkServer.shstart#启动ZooKeeper集群#以守护进程方式启动Kafka集群kafka-server-start.sh-daemon/opt/kafka/config/perties进程验证jpsmaster✓NameNode✓SecondaryNameNode✓Master✓QuorumPeerMain✓Kafka✓ResourceManagerslave01/slave02✓DataNode✓NodeManager✓Kafka✓QuorumPeerMain✓WorkerStructuredStreaming程序(上)使用StructuredStreaming实现实时单词计数任务目标12核心流程1Spark监听Kafka指定主题↓2接收Producer数据↓3执行单词计数↓4终端实时输出结果3文件路径master节点/opt/example/kafka_streaming.pykafka_streaming.py12345678910111213141516171819202122232425importsysfrompyspark.sqlimportSparkSessionfrompyspark.sql.functionsimportexplode,splitif__name__=="__main__":
if
len(sys.argv)!=3:
print("""Usage:structured_kafka_wordcount.py<bootstrap-servers><topics>""",file=sys.stderr)sys.exit(-1)bootstrapServers=sys.argv[1]topics=sys.argv[2]
#创建SparkSessionspark=SparkSession\.builder\.appName("StructuredKafkaWordCount")\.getOrCreate()
#从Kafka获取输入并创建Datasetlines=spark\.readStream\.format("kafka")\.option("kafka.bootstrap.servers",bootstrapServers)\.option("subscribe",topics)\.load()\.selectExpr("CAST(valueASSTRING)")StructuredStreaming程序(下)kafka_streaming.py1分隔数据
explode()
将数组元素拆分为单独行2单词计数按
word
分组并统计
count()3执行输出输出模式
complete
,格式
consolekafka_streaming.py
words=lines.select(
explode(
split(lines.value,'')
).alias('word')
)
wordCounts=words.groupBy('word').count()
query=wordCounts\
.writeStream\
.outputMode('complete')\
.format('console')\
.start()
query.awaitTermination()1
#分隔数据:explode()将数组元素拆分为单独行2345678
#单词计数91011
#执行输出12131415161718测试验证步骤5-9:环境准备与程序提交5创建Kafka主题kafka-topics.sh--create--bootstrap-servermaster:9092--replication-factor1--partitions1--topickwc6查看主题列表kafka-topics.sh--zookeepermaster:2181-list7启动生产者(slave01)kafka-console-producer.sh--broker-listmaster:9092--topickwc8启动消费者(slave02)kafka-console-consumer.sh--bootstrap-servermaster:9092--topickwc--from-beginning9提交Spark作业spark-submit/opt/example/kafka_streaming.pymaster:9092kwc步骤10-11:数据输入与结果验证工单小结通过Python程序实现StructuredStreaming实时监听Kafka,接收Flume采集数据并完成单词计数。slave01终端输入步骤10在slave01终端依次输入:>ILoveChina>ILoveMyCollege>ILoveMyFamilymaster终端Batch结果步骤11I3Love3China1My2College1Family1结果验证与工单小结制定合理的规划,并为之付出努力,才能收获胜利的果实高铁发展:从追赶到自主创新动力系统信号系统安全系统轨道系统通过科技手段实现列车自动控制、状态监测和诊断,达成高速、高效、安全、可靠的运行模式。京津城际铁路建设,高铁发展迈入新阶段运营里程超越世界他国,沪杭、郑西等干线开通网络覆盖全国主要城市,实现自主创新学习启示:循序渐进,攻克难点SparkCoreStructuredStreaming学习和实践难度层层递增做好规划,逐步攻克知识难点离线计算实时计算素养课堂:遵循科学发展规律2005年2010年2015年本工单将使用
StructuredStreaming
实现实时计算:掌握窗口函数
window()
的原理、适用场景及使用方法;编写Python程序,实时按性别计算观影次数并写入
Kafka;最后通过
spark-submit
提交程序至Spark集群进行验证。基本信息工单编号5.4工单名称StructuredStreaming实时计算电影评分数据集建议学时4学时所属任务实时计算实现环境要求已经搭建好的实时计算环境工单目标知识目标掌握StructuredStreaming实时计算的原理技能目标能够使用StructuredStreaming将计算结果写入Kafka素养目标培养正面、积极的职业心态和正确的职业价值观培养与人交流、合作的能力工单5.4基本信息与目标1启动Hadoop和Spark集群start-dfs.sh
#启动HDFSstart-yarn.sh
#启动YARNstart-master.sh
#启动Spark主节点start-workers.sh
#启动Spark工作节点2启动ZooKeeper和Kafkamasterslave01slave02zkServer.shstartkafka-server-start.sh-daemon/opt/kafka/config/perties3模拟实时产生评分数据(JSON格式)importpandasaspdimportosimporttime#读取用户数据unames=['uid','age','gender','occupation','zip']users=pd.read_table('/opt/example/u.user',sep='|',header=None,names=unames)#读取评分数据rnames=['uid','mid','rating','timestamp']ratings=pd.read_table('/opt/example/u.data',sep='\t',header=None,names=rnames)#合并与筛选frame=pd.merge(ratings,users)df=frame[['mid','age','gender','rating']]#循环写入JSON格式日志fn='/opt/tmp/task5_2.log'ifos.path.exists(fn):os.remove(fn)count=0lines=df.shape[0]whilecount<lines:log='{"m":%d,"a":%d,"g":"%s","r":%d}\n'%(df.iloc[count,0],df.iloc[count,1],df.iloc[count,2],df.iloc[count,3])
withopen(fn,'a+')asfile:file.write(log)time.sleep(0.5)count=count+1启动集群与模拟数据生成创建Kafka主题4创建Kafka主题ssc2kafka-topics.sh--create--zookeepermaster:2181--replication-factor1--partitions1--topicssc25查看主题列表kafka-topics.sh--zookeepermaster:2181-list准备步骤7安装依赖pip3installkafka-python8创建输出主题kafka-topics.sh--create--zookeepermaster:2181--replication-factor1--partitions1--topicout_topic9编写计算程序/opt/example/task5_streaming_2.pyimportsysimportjsonfromkafkaimportKafkaProducerfrompyspark.sql.typesimportStructType,StructField,StringType,IntegerTypefrompyspark.sqlimportSparkSessionfrompyspark.sql.functionsimportfrom_jsonfrompyspark.sqlimportfunctionsasF#自定义Kafka输出函数,用于显示计算结果def
to_display_kafka(row):
ifrow.count!=0:data=[]tmp={row[0]:row[1]}data.append(tmp)producer=KafkaProducer(bootstrap_servers='master:9092')producer.send("out_topic",json.dumps(data).encode('utf8'))producer.flush()StructuredStreaming计算程序(上)StructuredStreaming计算程序(下)接上文,继续task5_streaming_2.py主程序:Kafka数据读取→窗口聚合→输出到Kafkatask5_streaming_2.pyif__name__=="__main__":
iflen(sys.argv)!=3:
print("""Usage:task5_streaming_2.py<bootstrap-servers><intopic>""",file=sys.stderr)sys.exit(-1)bootstrapServers=sys.argv[1]intopic=sys.argv[2]
#创建SparkSession入口spark=SparkSession\.builder\.appName("task5_streaming_2")\.getOrCreate()
#从Kafka读取数据lines=spark\.readStream\.format("kafka")\.option("kafka.bootstrap.servers",bootstrapServers)\.option("subscribe",intopic)\.load()\.selectExpr("CAST(timestampAStimestamp)","CAST(valueASSTRING)")
#处理数据:水位线+窗口聚合mc=lines\.withWatermark("timestamp","1seconds")\.groupBy(F.window(F.col("timestamp"),"1seconds","1seconds"),F.col("value"))\.count()df=mc.selectExpr("CAST(valueASSTRING)","CAST(countASSTRING)")
print(df)
#输出到Kafkaquery=df\.writeStream\.outputMode('append')\.foreach(to_display_kafka)\.start()query.awaitTermination()主程序执行流程1命令行参数校验传入bootstrap-servers与intopic,否则退出↓2创建SparkSession入口SparkSession.builder.appName("task5_streaming_2")↓3从Kafka读取数据readStream.format("kafka")订阅intopic并做类型转换↓4窗口聚合处理水位线1seconds+按1秒窗口groupBy计数↓5输出到KafkawriteStream,outputMode('append'),foreach写出并等待终止四终端测试流程1终端1提交Spark任务spark-submit/opt/example/task5_streaming_2.pymaster:9092ssc22终端2打开PySpark交互界面,输入步骤(3)代码并保存3终端3运行FlumeAgentflume-ngagent-cconf-f/opt/example/task5_flume_2.conf-namea1-Dflume.root.logger=INFO,console4终端4运行消费者查看结果kafka-console-consumer.sh--bootstrap-servermaster:9092--topicout_topic--from-beginning运行结果示例to_display_kafka()
函数将StructuredStreaming计算结果转化为JSON数据,写入主题
out_topic,启动消费者后即可查看实时计算结果:运行结果[{"M":"1"}]运行结果[{"M":"4"}]运行结果[{"M":"10"}]工单小结编写Python程序StructuredStreaming实时计算计算结果写入Kafka本工单通过编写Python程序,实现了使用
StructuredStreaming
对Kafka实时按性别计算观影次数,并将计算结果写入Kafka,完成了StructuredStreaming实时计算的基本流程。运行结果与工单小结知识链接03实时计算基础什么是实时计算互联网、Web应用、网络监控、传感监测、电信金融、生产制造等领域对数据实时处理的需求日益增加。与传统离线计算相比,实时计算能对海量数据实时处理,从数据采集到数据处理快速完成,确保及时性与准确性。典型应用淘宝、京东等电商平台收集用户点击行为与浏览历史,利用实时计算框架高效处理数据,分析用户购买意图和兴趣爱好,实现精准推荐。常见实时计算框架对比4款框架框架核心特点延迟级别现状ApacheSparkStreaming微批处理架构,按时间窗口切分批次秒级活跃ApacheStorm处理结果可直接持久化到数据库或HDFS毫秒级活跃ApacheFlink真正流式处理,支持事件时间语义和精确一次处理毫秒级活跃雅虎S4分布式流处理,可扩展性和容错能力—已停止维护SparkStreaming
实现了实时数据流的高效、可扩展
与
容错
处理,
支持从
KafkaFlumeKinesisTCP套接字
等多样化数据源获取数据,
通过
map()reduce()join()window()
等高级函数构建复杂处理逻辑,
输出至
文件系统、数据库
或
实时仪表盘。核心抽象DStream(DiscretizedStream)离散化流表示持续的数据流——由一系列时间上连续的RDD
组成,每个RDD包含特定时间间隔内的数据,
使SparkStreaming能以
微批处理方式
处理实时数据流。DStream处理流程1接收实时输入数据流→2按时间间隔划分批次→3提交Spark引擎处理→4生成数据批次内置数据来源类型特点典型来源基础来源StreamingContextAPI直接可用文件系统(HDFS/S3/NFS)、套接字高级来源需额外工具类,依赖外部非Spark库接口Kafka、Flume、Kinesis基础来源开箱即用;高级来源需引入对应外部库配合工具类接入。SparkStreaming简介SparkStreaming特点与优缺点⚡核心特点便捷易用支持
Java、Python、Scala
等多种语言,开发者可像编写离线程序一样编写实时计算程序,零额外学习成本无缝整合Spark体系基于
SparkCore
运行,RDD操作代码可直接复用于批处理,支撑数据分析与交互式应用✓优点✓DAG调度+RDD机制:小批量数据快速处理,实时计算能力强✓"处理且仅处理一次":粗粒度处理确保准确性,简化容错恢复✓DStream继承RDD特性:学习成本低,与RDD交互简单直观!缺点粗粒度数据处理引入
处理延迟——需累积一定量数据后才处理,实时性要求极高的场景可能成为制约因素无状态转换每个批次的处理不依赖于之前批次的数据map()对原始DStream的每个元素进行转换flatMap()每个输入项可映射为多个输出项filter()过滤,返回一个新的DStreamrepartition()改变DStream的并行程度union()合并两个DStreamcount()统计DStream元素数量reduce()聚合,返回单元素RDD的新DStreamcountByValue()统计每个键的出现次数reduceByKey()按给定函数聚合每个键对应的值join()连接两个DStreamDStream无状态转换操作DStream转换操作分为
无状态转换
和
有状态转换
两类DStream有状态转换与窗口操作有状态转换当前批次的处理需要使用之前批次的数据或中间结果滑动窗口转换需设置
窗口长度
和
滑动间隔
两个参数。窗口按滑动间隔在原始DStream上移动,停留时将窗口长度内的数据作为一个数据段处理。批次1批次2批次3批次4批次5批次6批次7批次8当前窗口覆盖的数据段(窗口长度内)窗口外批次(按滑动间隔移动)滑动窗口转换函数window(windowLength,slideInterval)基于源DStream生成固定窗口大小的DStream流countByWindow(windowLength,slideInterval)计算滑动窗口中流元素的数量reduceByWindow(func,windowLength,slideInterval)对滑动窗口中的数据进行聚合计算,返回单元素流StreamingContext与编写步骤StreamingContext核心要点与SparkStreaming程序的基本编写流程编写步骤(单词计数示例)Python#导入SparkContext和StreamingContextfrom
pyspark
import
SparkContextfrom
pyspark.streaming
import
StreamingContext#获取SparkContext对象sc=SparkContext(master="spark://master:7077",appName="task5_streaming_1")sc.setLogLevel("ERROR")#获取StreamingContext对象ssc=StreamingContext(sc,2)#从外部文件读取DStreamlines=ssc.textFileStream("/home/test")#实时数据流处理:提取第3列并统计词频counts=lines.map(lambdaline:line.split("\t")[2])\.map(lambdax:(x,1))\.reduceByKey(lambdaa,b:a+b)counts.pprint()#启动并等待完成ssc.start()ssc.awaitTermination()StreamingContext核心要点1启动后不可添加新的流计算逻辑2停止后无法再次启动,需创建新对象3同一时刻只能有一个活跃实例4stop()
默认同时停止SparkContext;若需保留,设置
stopSparkContext=false性能调优从运行效率与内存使用两个维度,优化SparkStreaming应用性能提升运行效率EFFICIENCY1提升并行处理能力:充分利用集群资源,避免任务集中;涉及shuffle的操作可增大并行度2降低序列化开销:采用
Kryo
等高效序列化方法,或实现自定义序列化接口3设置合理的批处理间隔:避免前序作业执行过长导致后续作业延迟4减轻任务提交和分发负担:Standalone和Coarse-grainedMesos模式通常比Fine-grainedMesos模式延迟更低优化内存使用MEMORY1控制批处理数据量:确保节点内存足够容纳批处理间隔内接收的全部数据2及时释放无用数据:设置合理的
spark.cleaner.ttl
时长自动清理,避免误删正在使用的数据3监控并调整垃圾回收策略:观察GC情况并调整策略,减少垃圾回收对作业执行的干扰StructuredStreaming是基于SparkSQL引擎构建的流处理框架,兼具可扩展性与容错性两大特性。开发者可使用与批处理相似的编程模式处理流式数据,支持ScalaJavaPython的DatasetAPI与DataFrameAPI,实现流聚合、事件时间窗口计算、流批连接等复杂操作。1微批处理Spark2.0起默认采用流数据切分为小批量作业,每批
<100ms实现低延迟与
ExactlyOnce
语义默认模式2持续处理Spark2.3+进一步降低延迟至毫秒级提供
AtLeastOnce
语义保证容错保障检查点(Checkpoint):定期保存中间状态与进度预写日志(Write-AheadLog):记录处理日志故障时从检查点恢复,保证端到端容错性StructuredStreaming简介StructuredStreaming编程模型编程模型StructuredStreaming基于
UnboundedData模型
实现,核心是将数据流抽象为
无限增长的数据表,持续流入的数据增量写入其中。核心机制输入表数据流抽象,每条数据项作为新记录追加Query执行查询操作,产生结果表结果表查询输出,按模式写入外部存储三种输出模式CompleteMode全量模式输出更新后的整个结果表AppendMode追加模式仅输出新增到结果表的数据UpdateMode更新模式输出新增或更新到结果表的数据三者关系与编写步骤SSparkStreaming数据抽象DStream(RDD组成的序列)延迟能力通常秒级延迟VS相同的连续数据流,抽象不同SSStructuredStreaming数据抽象DataFrame/Dataset能力继承继承SparkSQL静态处理能力并扩展到动态数据流延迟能力可实现毫秒级实时响应●两者处理的数据类似,都是连续不断的数据流;StructuredStreaming通过DataFrame/Dataset抽象,充分利用SparkSQL的强大功能处理实时数据流。编写StructuredStreaming程序的基本步骤1导入PySpark模块→2创建SparkSession对象→3创建数据源→4定义流计算过程→5启动流计算并输出结果代码示例:端口单词计数监听本机端口
9999,实现
StructuredStreaming
单词计数StructuredNetworkWordCountPythonfrompyspark.sqlimportSparkSessionfrompyspark.sql.functionsimportsplitfrompyspark.sql.functionsimportexplodeif__name__=="__main__":spark=SparkSession\.builder\.appName("StructuredNetworkWordCount")\.getOrCreate()spark.sparkContext.setLogLevel('WARN')lines=spark\.readStream\.format("socket")\.option("host","localhost")\.option("port",9999)\.load()words=lines.select(
explode(
split(lines.value,"")).alias("word"))wordCounts=words.groupBy("word").count()query=wordCounts\.writeStream\.outputMode("complete")\.format("console")\.trigger(processingTime="8seconds")\.start()query.awaitTermination()输入源:socket·localhost:9999输出模式:complete→console触发间隔:8secondsStructuredStreaming整合Kafka(上)StructuredStreaming提供统一API,实现批处理与流处理无缝衔接。流数据被抽象为持续增长的DataFrame,用户可像处理静态数据一样处理实时数据流,并保证端到端精确一次语义。读取Kafkakafka.bootstrap.servers必填参数subscribe必填参数valdf=spark.readStream.format("kafka").option("kafka.bootstrap.servers","host1:p
温馨提示
- 1. 本站所有资源如无特殊说明,都需要本地电脑安装OFFICE2007和PDF阅读器。图纸软件为CAD,CAXA,PROE,UG,SolidWorks等.压缩文件请下载最新的WinRAR软件解压。
- 2. 本站的文档不包含任何第三方提供的附件图纸等,如果需要附件,请联系上传者。文件的所有权益归上传用户所有。
- 3. 本站RAR压缩包中若带图纸,网页内容里面会有图纸预览,若没有图纸预览就没有图纸。
- 4. 未经权益所有人同意不得将文件中的内容挪作商业或盈利用途。
- 5. 人人文库网仅提供信息存储空间,仅对用户上传内容的表现方式做保护处理,对用户上传分享的文档内容本身不做任何修改或编辑,并不能对任何下载内容负责。
- 6. 下载文件中如有侵权或不适当内容,请与我们联系,我们立即纠正。
- 7. 本站不保证下载资源的准确性、安全性和完整性, 同时也不承担用户因使用这些下载资源对自己和他人造成任何形式的伤害或损失。
最新文档
- 2026年金融监管政策与法规专项训练题库
- 2026年环境保护与生态修复知识测试
- 2026年幼儿问题解决能力评估习题
- T∕CTCA 34-2026 软壳服装分类规范与功能要求
- 保洁员清洁工作心理素质试题及答案
- 探秘保山税务管理考试试题及答案
- 现当代文学上模拟考试试题及答案
- 酒泉小升初各科试题及答案解析
- 2026年农村卫生防疫知识普及考试试题
- 2026年城市公共安全知识及试题
- 2026年秋期人教版部编版小学数学四年级上册全册教学计划教案
- (2026秋新版)北师大版六年级数学上册全册教案
- 2026年云南二级造价师真题及答案解析
- 2026中国医疗器械CDMO行业订单结构变化与利润率走势
- 2026年湖北省黄冈市重点学校小升初入学分班考试语文考试试题及答案
- 化工行业事故案例警示学习课件
- 小区物业整体服务方案投标文件(技术方案)
- MT/T 146-2025树脂锚杆
- 京东生鲜冷链物流合同
- 尾矿库工程施工总体方案设计
- 2026年一级建造师之一建矿业工程实务考试题库300道及完整答案(易错题)
评论
0/150
提交评论