大数据Shuffle原理与实践

简介: 大数据Shuffle原理与实践

💨Shuffle概述

🎈在开源实现的MapReduce中,存在Map、 Shuffle、 Reduce三个阶段。

Shuffle过程是MapReduce的核心。

Map阶段:是在单机上进行的针对-一小块数据的计算过程。Shuffle阶段:在map阶段的基础,上,进行数据移动,为后续的reduce阶段做准备。reduce阶段:对移动后的数据进行处理,依然是在单机上处理一小份数据。

🎈为什么需要Shuffle?

在分布式计算框架中,数据本地化是一个很重要的考虑,即计算需要被分发到数据所在的位置,从而减少数据的移动,提高运行效率。 hadoop中,map负责数据的初级拆分获取解析, reduce负责最终数据的集总,除了业务逻辑的功能外,其他的核心数据处理都是由shuffle来支持。

🎈Shuffle有什么?

简单来说,shuffle中有三次的数据排序

🚩第一次是 快速排序,这是因为第一次的数据全部在内存中开辟了一个缓冲区,数据从map出来后,分批进入缓冲区,对它们的索引进行排序,并且按照map的逻辑进行分区,在出缓冲区落盘的时候,完成排序。🚩第二次是归并排序,将第一次分批出来的文件进行区内归并排序。🚩第三次也是归并排序,将所有的map Task第二次产生的文件进行区内归并排序

这三次可以看做是一个整体的过程,从这里应该可以看出,shuffle是一个比较耗费资源并且时间开销比较大的环节。

🎈为什么Shuffle对性能非常重要?

🚩M * R次网络连接🚩大量的数据移动🚩数据丢失风险🚩可能存在大量的排序操作🚩大量的数据序列化、反序列化操作🚩数据压缩


🍳在大数据场景下,数据shuffle表示了不同分区数据交换的过程,不同的shufle策略性能差异较大。 目前在各个引擎中shuffle都是优化的重点,在spark框架中,shuffle 是支撑spark进行大规模复杂 数据处理的基石。


shuffle的数据来源于map,所以可以对map端出来的数据进行处理,我们可以采用压缩的方式尽量减少数据的规模。

💨Shuffle算子

💨分类

spark中会导致shuffle操作的有以下几种算子:

🚩repartition类的操作:比如repartition、repartitionAndSortWithinPartitions、coalesce等

🚩byKey类的操作:比如reduceByKey、groupByKey、sortByKey等

🚩join类的操作:比如join、cogroup等

🚩Distinct类的操作:distinct

💨Spark中对shuffle的抽象

窄依赖:父RDD的每个分片至多被子RDD中的一个分片所依赖

宽依赖:父RDD中的分片可能被子RDD中的多个分片所依赖

💨算子内部的依赖关系

ShuffleDependency:

🍳CoGroupedRDD含Cogroup

🚩fullOuterJoin、rightOuterJoin、 leftOuterJoin🚩join

🍳ShuffledRDD

🚩combineByKeyWithClassTag       combineByKey       reduceByKey🚩Coalesce🚩sortByKey            sortBy

💨Shuffle过程

💨Write

Spark中需要Shuffle 输出的Map任务会为每个Reduce创建对应的bucket,Map产生的结果会根据设置的partitioner得到对应的bucketId,然后填 充到相应的bucket中去。每个Map的输出结果可能包含所有的Reduce所需要的数据,所以每个Map会创建R个bucket(R是reduce的 个数),M个Map总共会创建M*R个bucket。

💨Fether

Reduce去拖Map的输出数据,Spark提供了两套不同的拉取数据框架:通过socket连接去取数据;使用netty框架去取数据。Spark Map输出的数据没有经过排序,Spark Shuffle过来的数据也不会进行排序,Spark认为Shuffle过程中的排序不是必须的,并不是所有类型的Reduce需要的数据都需要排序,强 制地进行排序只会增加Shuffle的负担。educe拖过来的数据会放在一个HashMap中,HashMap中存储的也是对,key是Map输出的key,Map输出对应这个key的所有value组成HashMap的value。Spark将 Shuffle取过来的每一个对插入或者更新到HashMap中,来一个处理一个。HashMap全部放在内存中。

💨Shuffle Handle创建

🎈Register Shuffle时做的最重要的事情是根据不同条件创建不同的shuffle Handle

🎈Shuffle Handle与Shuffle Writer的对应关系

BypassMergeSortShuffleHandle——>BypassMergeSortShuffleWriter

SerializedShuffleHandle——>UnsafeShuffleWriter

BaseShuffleHandle——>SortShuffleWriter

💨Reader实现-网络时序图

🎈使用基于netty的网络通信框架

🎈位置信息记录在MapOutputTracker中

🎈主要会发送两种类型的请求

    🚩OpenBlocks请求

    🚩Chunk请求或Stream请求

💨Shuffle优化使用的技术: Netty Zero Copy

🚩可堆外内存,避免JVM堆内存到堆外内存的数据拷贝。🚩CompositeByteBuf、Unpooled.wrappedBuffer. ByteBuf.slice,可以合并、包装、切分数组,避免发生内存拷贝🚩Netty使用FileRegion实现文件传输,FileRegion底层封装了FileChannel#transferTo()方法,可以将文件缓冲区的数据直接传输到目标Channel, 避免内核缓冲区和用户态缓冲区之间的数据拷贝


在第一次排序之后,此时由于原数据中各个字段可能会有数据分布不均,这样会导致reduce端处理数据时的数据倾斜——各个Task的处理量相差悬殊,可以在此处进行初步的数据合并处理。

应用场景:

去重操作;聚合,byKey类操作;排序操作等

💨常见问题

🚩数据存储在本地磁盘,没有备份

🚩I0并发:大量RPC请求(M*R)

🚩I0吞吐:随机读、写放大(3X)

🚩GC频繁,影响NodeManager

💨Shuffle优化

🚩避免shuffle,使用broadcast替代join

🚩使用可以map-side预聚合的算子


相关实践学习
基于MaxCompute的热门话题分析
Apsara Clouder大数据专项技能认证配套课程:基于MaxCompute的热门话题分析
目录
相关文章
|
12月前
|
存储 分布式计算 大数据
MaxCompute聚簇优化推荐功能发布,单日节省2PB Shuffle、7000+CU!
MaxCompute全新推出了聚簇优化推荐功能。该功能基于 31 天历史运行数据,每日自动输出全局最优 Hash Cluster Key,对于10 GB以上的大型Shuffle场景,这一功能将直接带来显著的成本优化。
462 3
|
12月前
|
存储 数据采集 搜索推荐
Java 大视界 -- Java 大数据在智慧文旅旅游景区游客情感分析与服务改进中的应用实践(226)
本篇文章探讨了 Java 大数据在智慧文旅景区中的创新应用,重点分析了如何通过数据采集、情感分析与可视化等技术,挖掘游客情感需求,进而优化景区服务。文章结合实际案例,展示了 Java 在数据处理与智能推荐等方面的强大能力,为文旅行业的智慧化升级提供了可行路径。
Java 大视界 -- Java 大数据在智慧文旅旅游景区游客情感分析与服务改进中的应用实践(226)
|
数据采集 SQL 搜索推荐
大数据之路:阿里巴巴大数据实践——OneData数据中台体系
OneData是阿里巴巴内部实现数据整合与管理的方法体系与工具,旨在解决指标混乱、数据孤岛等问题。通过规范定义、模型设计与工具平台三层架构,实现数据标准化与高效开发,提升数据质量与应用效率。
3587 0
大数据之路:阿里巴巴大数据实践——OneData数据中台体系
|
存储 SQL 分布式计算
大数据之路:阿里巴巴大数据实践——元数据与计算管理
本内容系统讲解了大数据体系中的元数据管理与计算优化。元数据部分涵盖技术、业务与管理元数据的分类及平台工具,并介绍血缘捕获、智能推荐与冷热分级等技术创新。元数据应用于数据标签、门户管理与建模分析。计算管理方面,深入探讨资源调度失衡、数据倾斜、小文件及长尾任务等问题,提出HBO与CBO优化策略及任务治理方案,全面提升资源利用率与任务执行效率。
800 0
|
11月前
|
存储 SQL 分布式计算
MaxCompute 聚簇优化推荐原理
基于历史查询智能推荐Clustered表,显著降低计算成本,提升数仓性能。
574 4
MaxCompute 聚簇优化推荐原理
|
10月前
|
人工智能 Cloud Native 算法
拔俗云原生 AI 临床大数据平台:赋能医学科研的开发者实践
AI临床大数据科研平台依托阿里云、腾讯云,打通医疗数据孤岛,提供从数据治理到模型落地的全链路支持。通过联邦学习、弹性算力与安全合规技术,实现跨机构协作与高效训练,助力开发者提升科研效率,推动医学AI创新落地。(238字)
620 7
|
存储 监控 大数据
大数据之路:阿里巴巴大数据实践——事实表设计
事实表是数据仓库核心,用于记录可度量的业务事件,支持高性能查询与低成本存储。主要包含事务事实表(记录原子事件)、周期快照表(捕获状态)和累积快照表(追踪流程)。设计需遵循粒度统一、事实可加性、一致性等原则,提升扩展性与分析效率。
897 0
|
存储 搜索推荐 算法
Java 大视界 -- Java 大数据在智慧文旅旅游线路规划与游客流量均衡调控中的应用实践(196)
本实践案例深入探讨了Java大数据技术在智慧文旅中的创新应用,聚焦旅游线路规划与游客流量调控难题。通过整合多源数据、构建用户画像、开发个性化推荐算法及流量预测模型,实现了旅游线路的精准推荐与流量的科学调控。在某旅游城市的落地实践中,游客满意度显著提升,景区流量分布更加均衡,充分展现了Java大数据技术在推动文旅产业智能化升级中的核心价值与广阔前景。
|
存储 分布式计算 大数据
大数据之路:阿里巴巴大数据实践——大数据领域建模综述
数据建模解决数据冗余、资源浪费、一致性缺失及开发低效等核心问题,通过分层设计提升性能10~100倍,优化存储与计算成本,保障数据质量并提升开发效率。相比关系数据库,数据仓库采用维度建模与列式存储,支持高效分析。阿里巴巴采用Kimball模型与分层架构,实现OLAP场景下的高性能计算与实时离线一体化。
1226 0
|
SQL 缓存 监控
大数据之路:阿里巴巴大数据实践——实时技术与数据服务
实时技术通过流式架构实现数据的实时采集、处理与存储,支持高并发、低延迟的数据服务。架构涵盖数据分层、多流关联,结合Flink、Kafka等技术实现高效流计算。数据服务提供统一接口,支持SQL查询、数据推送与定时任务,保障数据实时性与可靠性。
1833 0