资讯动态

spark

发布时间:2026/9/17 11:47:52 来源:尧图企业网站定制
******sparkrdd怎么体现他的弹性和特征rdd就是弹性分布式数据集特征属性1.一组分区 2.一个计算每个分区的函数 3.RDD之间的依赖关系 4.一个分区器 5.一个列表弹性体现1.自动进行内存和磁盘的切换2.基于lineage的高效容错3.task如果失败会进行特点次数的重试4.Stage如果失败会自动进行特点次数的重试而且只会计算失败的分片5.spark的缓存机制将数据缓存在内存或磁盘里对数据进行重用6.数据的调度弹性DAG Task与资源管理无关7.数据分片的高度弹性RDD 的“弹性”本质上不是简单的“能放内存也能放磁盘”而是它把数据拆成分区并记录完整血缘所以在节点故障、task 失败、部分数据丢失时仍然能按分区级别恢复和继续计算。******说下spark中的transform和action为什么spark要把操作分为transform和action转换算子单valuemapflatmapgroupbyrepartitionfilterdistinct双valueintersectionunionsubtractzipk-vpartitionbyreducebykeygroupbykeyjoin行动算子reducecollectfirsttakesaveforeach转换算子是对一个rdd操作得到一个新的rdd不会立即执行计算行动算子是触发对rdd的计算得到一个结果转换算子都是采用的懒策略只有在行动算子被提交的时候才被触发延迟计算的方式可以让 Spark 更好地优化计算过程。它可以分析整个 DAG找出最有效的执行路径合并多个操作减少中间数据的存储和传输从而提高计算效率。******repartition和coalesce的区别二者都是用来改变rdd的partition数量的repartition底层调用的就是coalesce方法repartition一定会发生shufflecoalesce根据传入参数来判断是否有shuffle一般情况增大分区数用repartition减少用coalesce******map和mappartitions区别map每次出来一条mp每次处理一个分区数据******reducebykey和groupbykey区别r有预聚合g没用预聚合在不影响业务的前提下用r******spark的stage任务怎么划分的spark任务会根据rdd之间的依赖关系形成一个DAG有向无环图然后DAGScheduler会把DAG划分成一系列的stage划分stage的依据就是rdd之间的宽窄依赖遇到宽依赖就划分stagestage的个数宽依赖个数1******宽窄依赖有shuffle的就是宽依赖其实也就是看父rdd的一个partition流向如果流向子RDD的多个分区就是窄依赖流向多个就是宽依赖窄依赖没有shuffle因为数据直接中父RDD流向子RDD宽依赖有shuffle因为数据需要在节点间重新分布窄依赖指的是父 RDD 的一个 partition 只会被一个子 RDD partition 依赖通常不需要 shuffle比如 map、filter。宽依赖指的是父 RDD 的一个 partition 可能被多个子 RDD partition 依赖通常需要 shuffle比如 reduceByKey、groupByKey、join。Spark 也正是基于宽依赖作为边界来切分 stage。******spark的shuffer有哪些讲讲各自的特点hashshuffle2.0之后就不用了优化前下游有多少个task就会生成多少个文件优化后通过复用buffer来优化shuffle过程中产生小文件的数量sortshuffle优势减少了小文件数据首先写入内存当内存装不下会溢写到磁盘在溢写磁盘前会先根据key排序然后分批写入磁盘文件中写入磁盘是是通过缓冲区溢写的方式每次溢写都会产生一个磁盘文件也就是说一个task过程会产生多个临时文件最后在每个task中将所有的临时文件合并这就是merge过程map阶段原始rdd的每个元素都会经过多个转换操作并生成键值对shuffle阶段首先会将这些键值对根据key排序然后按照key的hash值被分到不同的文件中在这个过程中数据会被写入到本地磁盘。然后启动新的任务这些任务会读取上有中他们需要的那部分数据reduce阶段每个任务从shuffle阶段获得所有数据并进行聚合操作当shuffle read task的数量小于等于默认的200个时并且不是聚合类的shuffle算子就会启动bypass机制这种机制为每个下游任务创建一个临时磁盘文件并将数据根据key的哈希值写入对应的文件中。完成后所有临时文件会被合并成一个文件并创建一个索引文件这种方式减少了磁盘I/O操作也减少了小文件不会排序效率高Spark 现在默认是 SortShuffle 。它的核心流程是map 端先把 shuffle 数据写入内存缓冲区内存不足时 spill 到磁盘spill 前会按分区组织数据可能还会按 key 排序。一个 task 可能产生多个 spill 文件task 结束时再 merge 成最终的数据文件和索引文件。相比早期 HashShuffle 它最大的优势是减少 shuffle 文件数量降低磁盘和文件句柄压力。对于 reducer 数量较少且不是聚合类算子的场景还可能走 bypass 机制直接跳过排序提高效率。******spark提交client模式和cluster模式的区别client模式Driver运行在client上通过AM向RM获取资源本地Driver负责与所有的executor container交互如果client挂了任务也就挂了cluster模式Driver运行在AM上就算client挂了任务也能正常运行******spark大体的一个执行流程1.driver执行main方法懒执行action算子触发job2.根据宽窄依赖划分stage3.每个stage会被整理成多个task4.每个task分发到具体的executor去执行******spark的行动算子和其他算子有什么不同转换算子都是采用的懒策略转换算子是对一个rdd操作得到一个新的rdd不会立即执行计算行动算子是触发对rdd的计算得到一个结果******spark统一内存模型spark的统一内存模型由统一内存和其他内存组成默认的是统一内存占0.6其它内存占0.4统一内存storage内存用于缓存数据默认占统一内存的0.5execution内存用于缓存在shuffle过程中产生的中间数据二者之间可通过动态占用机制在必要时去占用对方空余的内存但execution被占用可强制收回其他用户定义的数据结构或spark内部数据元******spark和mr的对比1.内存硬盘spark更强调中间数据在优先内存中进行计算减少反复落魄迭代计算效率更高mr具有内存缓存机制但它不是基于内存计算的框架因为其核心计算过程仍然依赖于磁盘作为数据的主要存储和传输媒介2.资源申请粒度MR 的 task 更偏独立进程式执行Spark 是先拿到 Executor 进程再在线程级别执行多个 task。开启和调度进程的代价一般大于线程的代价3.容错性在容错上Spark 可以基于 RDD 血缘做 分区 级重算恢复更灵活MR 也有 task 重试机制但整体更依赖磁盘落地恢复成本通常更高。4.mr优点运行很稳定适合长期后台运行5.spark缺点不稳定因为是基于内存运算的如果数据量超出内存会出现挂掉的现象******sprak和mr的shuffle区别1.hadoop不用等所有的maptask都结束后开启reducetask可以在 map task 尚未全部结束时提前启动并拉取已经完成的 map 输出而spark必须要等到父Stage都完成才能去fetch数据2.hadoop的shuffle是必须排序的不管是map输出还是reduce输出都是分区内有序而spark不一定要求排序******Spark会产生的shuffle算子重分区算子repartition、coalescebykey算子groupbykey、reducebykeyjoin算子join去重算子distinct******spark持久化缓存机制有以下三种方法1.cache是persist的一种简化形式设定将数据保存在内存里面2.persist主要目的是在内存或者磁盘中缓存 RDD 的数据使得在后续的操作中能够快速地访问该 RDD持久化这个 RDD可以避免每次使用时都重新计算它3.checkpoint一般是将数据保存在hdfs上的重点在于容错和故障恢复当 Spark 作业因为节点故障等原因需要重新计算时我们从最近的检查点位置开始计算而不是从最初的数据源和依赖关系开始重新计算整个 RDD 的转换链******spark数据倾斜1.主要表现单个Executor执行时间特别长整体任务卡在了某个阶段不能结束spark任务oom异常退出2.解决1.过滤少数导致倾斜的key如果发现导致倾斜的key就少数几个而且对计算本身影响并不大那就可以将这些key给过滤掉2.提高shuffle操作的并行度最简单效果差增加shuffle read task的数量可以让原本分配给一个task的多个key分配到多个task上从而让每个task处理比原来更少的数据这张方式是最简单的但一般只能缓解数据倾斜如果某个key值倾斜的非常严重那还是会导致拿到这个key的task会比其他的task处理更多的数据3.两阶段聚合局部聚和全局聚合第一次给每个key都加上一个随机数前缀进行一次局部聚合第二次将加上的随机数前缀都取消掉再进行一次全局聚合4.将reduce join转换成map join小RDD join 大RDD的时候就是将小表广播出去然后在大表操作时使用map算子从广播变量种获取小表数据进行合并这样就在map端完成了join避免了shuffle5.使用随机前缀和扩容RDD进行join适用rdd中有大量的key都会导致数据倾斜1.将造成数据倾斜的rdd的每条数据都打上n以内的随机前缀2.同时对另外一个正常的rdd进行扩容将每条数据都扩容成n条数据并依次打上0-n的前缀3.最后将两个处理后的rdd进行join即可******spark什么时候发生shuffle在stage任务阶段划分的时候会shuffle从算子角度来说像joinreduceByKey可能发生shuffle******spark架构及提交流程架构driver执行main方法跟踪executor的执行情况等executor负责运行具体任务master负责资源调度和分配类似yarn中的RMworker具体执行处理和计算的过程APPmaster用于向资源调度器申请执行任务的资源容器提交流程以yarn-cluster为例1.客户端向RM提交请求上传jar包到hdfs2.rm在集群中选择一个nodemanager在其上启动appmaster在appmaster中启动driver线程并实例化sparkcontext3.appmaster向rm注册应用程序并申请资源rm监控appmaster的状态直到appmaster结束4.appmaster申请到资源后与nodemanager通信在container中启动executor进程5.executor向driver反向注册申请任务6.dirver对应用进行解析最后将task发送到executor上7.在executor中执行task并将执行结果或状态汇报driver8.执行完毕后appmaster通知rm注销应用回收资源******spark的使用场景适用于迭代计算和实时性要求高的数据分析******spark的广播变量广播变量的优势driver每次分发任务的时候会把task和相关变量发送给executor如果不使用广播变量的话在每个executor中有多少个task就有多少个driver端变量副本。这样会导致消耗大量的内存好处使用广播变量后不需要每个task都带上一份变量副本而是变成每个节点的executor一份副本各个task去所在的executor上获取副本******spark怎么获得一个rdd1.读取内存数据创建RDD比如makeRDD方法2.读取文件创建RDD比如textfile方法******sparkstreaming精准一次消费1.开启事务把处理消息和提交偏移量放到一个事务里2.手动提交偏移量开启幂等性先确保真正处理完数据后再提交偏移量幂等性保证数据无论被保存多少次效果都是一样的******sparkstreaming工作机制在sparkstreaming中会有一个组件receiver作为一个长期运行的任务运行在一个executor上每个receiver都会负责一个DStream输入流。receiver组件接收到数据源发来的数据后会提交给sparkstreaming程序处理。处理后的结果可以进行可视化展示或者写入到hdfs中******spark的checkpoint原理底层是由一个generateJobs方法负责streaming job的产生产生并提交后就会发送docheckpoint事件然后调用checkpointwriter将checkpoint信息写到checkpoint目录下然后有个updatecheckpointdata方法对每个dstream信息转化成checkpoint******hive on spark 和spark on hive的区别hive on spark元数据存储在mysql执行引擎rdd语法hive sqlspark on hive元数据存储在mysql执行引擎df ds语法spark sql******spark调优1.资源调优比如提交一个复杂的任务的时候可以调executor的内存、核数还有driver的内存大小等在资源比较充足的情况下尽可能的使用更多的计算资源合理设置并行度2.开发调优1.rdd重用和持久化可以把多次使用到的rdd进行次持久化避免后续需要再次计算持久化的时候可以采用序列化这样可以大大减少内存空间的占用2.使用广播变量背一遍广播变量的好处3.在进行两个rddjoin的时候可以采用广播变量加mapjoin的方式来完成join操作避免了shuffle。启用map-side预聚合功能减少shuffle数据量提高性能******sparkSQL如何合并小文件1.如果任务不复杂数据量少可以直接降低并行度2.单独开一个任务专门合并小文件******spark背压机制有三个组件监听器主要监听集群所有作业的提交、运行、完成情况可以把一些信交给速率估算器速率估算器把收集到的数据和一个设定值进行比较然后用它们之间的差来计算新的输入值估算出一个用于下一批次的流量阈值限流器根据阈值进行限流******spark挂了怎么办看日志吗了解访问yarn页面看日志找到对应的任务可以进入到spark UI页面去查看一些情况******spark中有了RDD为什么还要有Dataframe和DataSet我的理解时dfds有比rdd更丰富更方便的功能rdd不支持sparksqlds、df都支持sparksql操作dsdf还支持一些特别方便的保存方式比如可以保存成csv文件可以带上表头这样每一列的字段名都一目了然******Spark任务和Spark Streaming任务的差别不是很了解我的理解就是spark的计算是基于rdd的也就是一次操作针对一个rddspark streaming是将一个时间区间内的rdd封装成一个ds对这个ds进行操作******sparkstreaming如何实现容错对于executor发送故障可以通过wal将数据预先写到dhfs上对于driver故障1.配置driver程序自动重启2.重启时从宕机的地方重启******flink和sparkStreaming区别1.数据处理flink是标准的实时处理引擎基于事件驱动支持流批一体sparkStreaming是微批次处理的2.架构sparkStreaming运行时的主要角色包括masterworkerdriverexecutorflink包括jobmanagertaskmanagerslot3.时间机制sparkStreaming只支持处理时间flink支持处理时间、时间时间、进入时间4.容错机制对于ss可以设置checkpoint如果发送故障并重启可以从上次checkpoint处恢复但这种行为只能保证数据不丢可能会重复处理对于flink则使用两阶段提交协议来解决这个问题

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价