俯瞰旭日升 發表於 2023-5-22 00:00:00

Spark 数据倾斜及其解决方案

<p data-id="pd157317-40ABmmfx">
        本文从数据倾斜的危害、现象、原因等方面,由浅入深阐述Spark数据倾斜及其解决方案。</p>
<h3 data-id="h26976cb-GeoUArjo" id="h26976cb-GeoUArjo">
        一、什么是数据倾斜</h3>
<p data-id="pd157317-olzEEwtP">
        对 Spark/Hadoop 这样的分布式大数据系统来讲,数据量大并不可怕,可怕的是数据倾斜。</p>
<p data-id="pd157317-IJr3WsOx">
        对于分布式系统而言,理想情况下,随着系统规模(节点数量)的增加,应用整体耗时线性下降。如果一台机器处理一批大量数据需要120分钟,当机器数量增加到3台时,理想的耗时为120 / 3 = 40分钟。但是,想做到分布式情况下每台机器执行时间是单机时的1 / N,就必须保证每台机器的任务量相等。不幸的是,很多时候,任务的分配是不均匀的,甚至不均匀到大部分任务被分配到个别机器上,其它大部分机器所分配的任务量只占总得的小部分。比如一台机器负责处理 80% 的任务,另外两台机器各处理 10% 的任务。</p>
<p data-id="pd157317-W71RBjpD">
        『不患多而患不均』,这是分布式环境下最大的问题。意味着计算能力不是线性扩展的,而是存在短板效应: 一个 Stage 所耗费的时间,是由最慢的那个 Task 决定。</p>
<p data-id="pd157317-OmT3acg4">
        由于同一个 Stage 内的所有 task 执行相同的计算,在排除不同计算节点计算能力差异的前提下,不同 task 之间耗时的差异主要由该 task 所处理的数据量决定。所以,要想发挥分布式系统并行计算的优势,就必须解决数据倾斜问题。</p>
<h3 data-id="h26976cb-Xn5m52FJ" id="h26976cb-Xn5m52FJ">
        二、数据倾斜的危害</h3>
<p data-id="pd157317-GIZkmdIy">
        当出现数据倾斜时,小量任务耗时远高于其它任务,从而使得整体耗时过大,未能充分发挥分布式系统的并行计算优势。</p>
<p data-id="pd157317-b8qtF2NP">
        另外,当发生数据倾斜时,部分任务处理的数据量过大,可能造成内存不足使得任务失败,并进而引进整个应用失败。</p>
<h3 data-id="h26976cb-V5Un9JK4" id="h26976cb-V5Un9JK4">
        三、数据倾斜的现象</h3>
<p data-id="pd157317-TaHjQ3F2">
        当发现如下现象时,十有八九是发生数据倾斜了:</p>
<ul data-id="ucd67dc5-4Eroh3Fl">
<li data-id="l20de63f-uQyBcSWm">
                绝大多数 task 执行得都非常快,但个别 task 执行极慢,整体任务卡在某个阶段不能结束。</li>
        <li data-id="l20de63f-jwxj9Aj7">
                原本能够正常执行的 Spark 作业,某天突然报出 OOM(内存溢出)异常,观察异常栈,是我们写的业务代码造成的。这种情况比较少见。</li>
</ul>
<p data-id="pd157317-B3U5GAT3">
        <strong>TIPS</strong></p>
<p data-id="pd157317-YFtMXI0U">
        在 Spark streaming 程序中,数据倾斜更容易出现,特别是在程序中包含一些类似 sql 的 join、group 这种操作的时候。因为 Spark Streaming 程序在运行的时候,我们一般不会分配特别多的内存,因此一旦在这个过程中出现一些数据倾斜,就十分容易造成 OOM。</p>
<h3 data-id="h26976cb-8DP2ECCP" id="h26976cb-8DP2ECCP">
        四、数据倾斜的原因</h3>
<p data-id="pd157317-Z3NegPe6">
        在进行 shuffle 的时候,必须将各个节点上相同的 key 拉取到某个节点上的一个 task 来进行处理,比如按照 key 进行聚合或 join 等操作。此时如果某个 key 对应的数据量特别大的话,就会发生数据倾斜。比如大部分 key 对应10条数据,但是个别 key 却对应了100万条数据,那么大部分 task 可能就只会分配到10条数据,然后1秒钟就运行完了;但是个别 task 可能分配到了100万数据,要运行一两个小时。</p>
<p data-id="pd157317-T6vyMyHc">
        因此出现数据倾斜的时候,Spark 作业看起来会运行得非常缓慢,甚至可能因为某个 task 处理的数据量过大导致内存溢出。</p>
<h3 data-id="h26976cb-9d2A51WU" id="h26976cb-9d2A51WU">
        五、问题发现与定位</h3>
<h4 data-id="h6e90be6-Vqd4vhxl" id="h6e90be6-Vqd4vhxl">
        1、通过 Spark Web UI</h4>
<p data-id="pd157317-k2XzrpD5">
        通过 Spark Web UI 来查看当前运行的 stage 各个 task 分配的数据量(Shuffle Read Size/Records),从而进一步确定是不是 task 分配的数据不均匀导致了数据倾斜。</p>
<p data-id="pd157317-YoG1JRok">
        知道数据倾斜发生在哪一个 stage 之后,接着我们就需要根据 stage 划分原理,推算出来发生倾斜的那个 stage 对应代码中的哪一部分,这部分代码中肯定会有一个 shuffle 类算子。可以通过 countByKey 查看各个 key 的分布。</p>
<p data-id="pd157317-R2KCT6Oq">
        <strong>TIPS</strong></p>
<p data-id="pd157317-LylAMczH">
        数据倾斜只会发生在 shuffle 过程中。这里给大家罗列一些常用的并且可能会触发 shuffle 操作的算子: distinct、groupByKey、reduceByKey、aggregateByKey、join、cogroup、repartition 等。出现数据倾斜时,可能就是你的代码中使用了这些算子中的某一个所导致的。</p>
<h4 data-id="h6e90be6-g0gYGhlm" id="h6e90be6-g0gYGhlm">
        2、通过 key 统计</h4>
<p data-id="pd157317-aaLAWuyB">
        也可以通过抽样统计 key 的出现次数验证。</p>
<p data-id="pd157317-kDvYgI1J">
        由于数据量巨大,可以采用抽样的方式,对数据进行抽样,统计出现的次数,根据出现次数大小排序取出前几个:</p>
<pre>
<span class="cm-variable">df</span>.<span class="cm-property">select</span>(<span class="cm-string">"key"</span>).<span class="cm-property">sample</span>(<span class="cm-atom">false</span>, <span class="cm-number">0.1</span>) <span class="cm-comment">// 数据采样</span> .(<span class="cm-def">k</span> <span class="cm-operator">=&gt;</span> (<span class="cm-variable-2">k</span>, <span class="cm-number">1</span>)).<span class="cm-property">reduceBykey</span>(<span class="cm-variable">_</span> <span class="cm-operator">+</span> <span class="cm-variable">_</span>) <span class="cm-comment">// 统计 key 出现的次数</span> .<span class="cm-property">map</span>(<span class="cm-def">k</span> <span class="cm-operator">=&gt;</span> (<span class="cm-variable-2">k</span>.<span class="cm-property">_2</span>, <span class="cm-variable-2">k</span>.<span class="cm-property">_1</span>)).<span class="cm-property">sortByKey</span>(<span class="cm-atom">false</span>) <span class="cm-comment">// 根据 key 出现次数进行排序</span> .<span class="cm-property">take</span>(<span class="cm-number">10</span>) <span class="cm-comment">// 取前 10 个。</span></pre>
<p data-id="pd157317-RndN4Miu">
        如果发现多数数据分布都较为平均,而个别数据比其他数据大上若干个数量级,则说明发生了数据倾斜。</p>
<h3 data-id="h26976cb-3aMnOnKU" id="h26976cb-3aMnOnKU">
        六、如何缓解数据倾斜</h3>
<p data-id="pd157317-3JLK68oo">
        <strong>基本思路</strong></p>
<p data-id="pd157317-Iwc3yDo2">
        <strong>业务逻辑:</strong> 我们从业务逻辑的层面上来优化数据倾斜,比如要统计不同城市的订单情况,那么我们单独对这一线城市来做 count,最后和其它城市做整合。</p>
<p data-id="pd157317-KCfqwXdJ">
        <strong>程序实现: </strong>比如说在 Hive 中,经常遇到 count(distinct)操作,这样会导致最终只有一个 reduce,我们可以先 group 再在外面包一层 count,就可以了;在 Spark 中使用 reduceByKey 替代 groupByKey 等。</p>
<p data-id="pd157317-BQUa5mxk">
        <strong>参数调优: </strong>Hadoop 和 Spark 都自带了很多的参数和机制来调节数据倾斜,合理利用它们就能解决大部分问题。</p>
<h4 data-id="h6e90be6-fVuFSdLe" id="h6e90be6-fVuFSdLe">
        思路1. 过滤异常数据</h4>
<p data-id="pd157317-4euUXAWo">
        如果导致数据倾斜的 key 是异常数据,那么简单的过滤掉就可以了。</p>
<p data-id="pd157317-zHUdFpYB">
        首先要对 key 进行分析,判断是哪些 key 造成数据倾斜。具体方法上面已经介绍过了,这里不赘述。</p>
<p data-id="pd157317-gVyJAyEq">
        然后对这些 key 对应的记录进行分析:</p>
<ul data-id="ucd67dc5-3jMNFNh2">
<li data-id="l20de63f-uDcLXBMy">
                空值或者异常值之类的,大多是这个原因引起</li>
        <li data-id="l20de63f-aVbl5FJU">
                无效数据,大量重复的测试数据或是对结果影响不大的有效数据</li>
        <li data-id="l20de63f-DuegXWi2">
                有效数据,业务导致的正常数据分布</li>
</ul>
<p data-id="pd157317-YoXARG9j">
        <strong>解决方案</strong></p>
<p data-id="pd157317-lSIpIVtM">
        对于第 1,2 种情况,直接对数据进行过滤即可。</p>
<p data-id="pd157317-6w1tAZjP">
        第3种情况则需要特殊的处理,具体我们下面详细介绍。</p>
<h4 data-id="h6e90be6-FapYfJPL" id="h6e90be6-FapYfJPL">
        思路2. 提高 shuffle 并行度</h4>
<p data-id="pd157317-Ql15Ijn7">
        Spark 在做 Shuffle 时,默认使用 HashPartitioner(非 Hash Shuffle)对数据进行分区。如果并行度设置的不合适,可能造成大量不相同的 Key 对应的数据被分配到了同一个 Task 上,造成该 Task 所处理的数据远大于其它 Task,从而造成数据倾斜。</p>
<p data-id="pd157317-yjxs6vHl">
        如果调整 Shuffle 时的并行度,使得原本被分配到同一 Task 的不同 Key 发配到不同 Task 上处理,则可降低原 Task 所需处理的数据量,从而缓解数据倾斜问题造成的短板效应。</p>
<p data-id="pd157317-2n3WMPlS">
        <strong>(1)操作流程</strong></p>
<ol data-id="o01bedff-f7ULZgQj">
<li data-id="l20de63f-YLZv7ofL">
                RDD 操作 可在需要 Shuffle 的操作算子上直接设置并行度或者使用 spark.default.parallelism 设置。如果是 Spark SQL,还可通过 SET</li>
        <li data-id="l20de63f-wkhModjV">
                spark.sql.shuffle.partitions= 设置并行度。默认参数由不同的 Cluster Manager 控制。</li>
        <li data-id="l20de63f-KYiHvdh5">
                dataFrame 和 sparkSql 可以设置</li>
        <li data-id="l20de63f-0rXi3yXQ">
                spark.sql.shuffle.partitions= 参数控制 shuffle 的并发度,默认为200。</li>
</ol>
<p data-id="pd157317-NCIIMeu9">
        <strong>(2)适用场景</strong></p>
<p data-id="pd157317-aX5WlOlc">
        大量不同的 Key 被分配到了相同的 Task 造成该 Task 数据量过大。</p>
<p data-id="pd157317-d5aUNdOK">
        <strong>(3)解决方案</strong></p>
<p data-id="pd157317-UJeBvDnW">
        调整并行度。一般是增大并行度,但有时如减小并行度也可达到效果。</p>
<p data-id="pd157317-J6pOr7K0">
        <strong>(4)优势</strong></p>
<p data-id="pd157317-tRwYXBQZ">
        实现简单,只需要参数调优。可用最小的代价解决问题。一般如果出现数据倾斜,都可以通过这种方法先试验几次,如果问题未解决,再尝试其它方法。</p>
<p data-id="pd157317-xCu4OYzX">
        (5)劣势</p>
<p data-id="pd157317-24nkiuND">
        适用场景少,只是让每个 task 执行更少的不同的key。无法解决个别key特别大的情况造成的倾斜,如果某些 key 的大小非常大,即使一个 task 单独执行它,也会受到数据倾斜的困扰。并且该方法一般只能缓解数据倾斜,没有彻底消除问题。从实践经验来看,其效果一般。</p>
<p data-id="pd157317-mJNPpFlv">
        TIPS 可以把数据倾斜类比为 hash 冲突。提高并行度就类似于 提高 hash 表的大小。</p>
<h4 data-id="h6e90be6-mI5gK4Kd" id="h6e90be6-mI5gK4Kd">
        思路3. 自定义 Partitioner</h4>
<p data-id="pd157317-abtKZiqq">
        <strong>(1)原理</strong></p>
<p data-id="pd157317-C303LSMj">
        使用自定义的 Partitioner(默认为 HashPartitioner),将原本被分配到同一个 Task 的不同 Key 分配到不同 Task。</p>
<p data-id="pd157317-slvmvrAa">
        例如,我们在 groupByKey 算子上,使用自定义的 Partitioner:</p>
<pre>
.<span class="cm-variable">groupByKey</span>(<span class="cm-keyword">new</span> <span class="cm-variable">Partitioner</span>() { <span class="cm-operator">@</span><span class="cm-variable">Override</span> <span class="cm-variable">public</span> <span class="cm-variable">int</span> <span class="cm-variable">numPartitions</span>() { <span class="cm-keyword">return</span> <span class="cm-number">12</span>;
} <span class="cm-operator">@</span><span class="cm-variable">Override</span> <span class="cm-variable">public</span> <span class="cm-variable">int</span> <span class="cm-variable">getPartition</span>(<span class="cm-variable">Object</span> <span class="cm-variable">key</span>) { <span class="cm-variable">int</span> <span class="cm-variable">id</span> <span class="cm-operator">=</span> <span class="cm-variable">Integer</span>.<span class="cm-property">parseInt</span>(<span class="cm-variable">key</span>.<span class="cm-property">toString</span>()); <span class="cm-keyword">if</span>(<span class="cm-variable">id</span> <span class="cm-operator">&gt;=</span> <span class="cm-number">9500000</span> <span class="cm-operator">&amp;&amp;</span> <span class="cm-variable">id</span> <span class="cm-operator">&lt;=</span> <span class="cm-number">9500084</span> <span class="cm-operator">&amp;&amp;</span> ((<span class="cm-variable">id</span> <span class="cm-operator">-</span> <span class="cm-number">9500000</span>) <span class="cm-operator">%</span> <span class="cm-number">12</span>) <span class="cm-operator">==</span> <span class="cm-number">0</span>) { <span class="cm-keyword">return</span> (<span class="cm-variable">id</span> <span class="cm-operator">-</span> <span class="cm-number">9500000</span>) <span class="cm-operator">/</span> <span class="cm-number">12</span>;
    } <span class="cm-keyword">else</span> { <span class="cm-keyword">return</span> <span class="cm-variable">id</span> <span class="cm-operator">%</span> <span class="cm-number">12</span>;
    }
}
})</pre>
<p data-id="pd157317-aWvVuPic">
        <strong>TIPS</strong> 这个做法相当于自定义 hash 表的 哈希函数。</p>
<p data-id="pd157317-o50sqent">
        <strong>(2)适用场景</strong></p>
<p data-id="pd157317-B1kvwcUc">
        大量不同的 Key 被分配到了相同的 Task 造成该 Task 数据量过大。</p>
<p data-id="pd157317-Gb5ugQZv">
        <strong>(3)解决方案</strong></p>
<p data-id="pd157317-cDnDpUkF">
        使用自定义的 Partitioner 实现类代替默认的 HashPartitioner,尽量将所有不同的 Key 均匀分配到不同的 Task 中。</p>
<p data-id="pd157317-2eQWtfNp">
        (4)优势</p>
<p data-id="pd157317-qOop7KBM">
        不影响原有的并行度设计。如果改变并行度,后续 Stage 的并行度也会默认改变,可能会影响后续 Stage。</p>
<p data-id="pd157317-2OiOFqss">
        <strong>(5)劣势</strong></p>
<p data-id="pd157317-ItoMrv4H">
        适用场景有限,只能将不同 Key 分散开,对于同一 Key 对应数据集非常大的场景不适用。效果与调整并行度类似,只能缓解数据倾斜而不能完全消除数据倾斜。而且需要根据数据特点自定义专用的 Partitioner,不够灵活。</p>
<h4 data-id="h6e90be6-xKGNVPZG" id="h6e90be6-xKGNVPZG">
        思路4. Reduce 端 Join 转化为 Map 端 Join</h4>
<p data-id="pd157317-gFqCWZfD">
        通过 Spark 的 Broadcast 机制,将 Reduce 端 Join 转化为 Map 端 Join,这意味着 Spark 现在不需要跨节点做 shuffle 而是直接通过本地文件进行 join,从而完全消除 Shuffle 带来的数据倾斜。</p>
<p data-id="pd157317-B6FVgknR">
        <img title="Spark 数据倾斜及其解决方案" alt="Spark 数据倾斜及其解决方案" src="https://zhuji.jb51.net/uploads/img/202305/796bef8819f789c59eda233061bcb6b8.jpg"></p>
<pre>
<span class="cm-variable">from</span> <span class="cm-variable">pyspark</span>.<span class="cm-property">sql</span>.<span class="cm-property">functions</span> <span class="cm-keyword">import</span> <span class="cm-def">broadcast</span> <span class="cm-variable">result</span> <span class="cm-operator">=</span> <span class="cm-variable">broadcast</span>(<span class="cm-variable">A</span>).<span class="cm-property">join</span>(<span class="cm-variable">B</span>, [<span class="cm-string">"join_col"</span>], <span class="cm-string">"left"</span>)</pre>
<p data-id="pd157317-5dZrrZQi">
        其中 A 是比较小的 dataframe 并且能够整个存放在 executor 内存中。</p>
<p data-id="pd157317-icdjH0Yf">
        <strong>(1)适用场景</strong></p>
<p data-id="pd157317-yeELpoyu">
        参与Join的一边数据集足够小,可被加载进 Driver 并通过 Broadcast 方法广播到各个 Executor 中。</p>
<p data-id="pd157317-GWZjMcT4">
        <strong>(2)解决方案</strong></p>
<p data-id="pd157317-1tjUkhyF">
        在 Java/Scala 代码中将小数据集数据拉取到 Driver,然后通过 Broadcast 方案将小数据集的数据广播到各 Executor。或者在使用 SQL 前,将 Broadcast 的阈值调整得足够大,从而使 Broadcast 生效。进而将 Reduce Join 替换为 Map Join。</p>
<p data-id="pd157317-j3R02gd6">
        <strong>(3)优势</strong></p>
<p data-id="pd157317-HBWGeP55">
        避免了 Shuffle,彻底消除了数据倾斜产生的条件,可极大提升性能。</p>
<p data-id="pd157317-5kUiYkJJ">
        <strong>(4)劣势</strong></p>
<p data-id="pd157317-ZVdKR1b6">
        因为是先将小数据通过 Broadcase 发送到每个 executor 上,所以需要参与 Join 的一方数据集足够小,并且主要适用于 Join 的场景,不适合聚合的场景,适用条件有限。</p>
<p data-id="pd157317-da7gVeQQ">
        <strong>NOTES</strong></p>
<ul data-id="ucd67dc5-IDfT1rXh">
<li data-id="l20de63f-Wnxv835y">
                使用Spark SQL时需要通过 SET</li>
        <li data-id="l20de63f-Bc39qviw">
                spark.sql.autoBroadcastJoinThreshold=104857600 将 Broadcast 的阈值设置得足够大,才会生效。</li>
</ul>
<h4 data-id="h6e90be6-7ffw4l8Y" id="h6e90be6-7ffw4l8Y">
        思路5. 拆分 join 再 union</h4>
<p data-id="pd157317-lBoSbJzM">
        思路很简单,就是将一个 join 拆分成 倾斜数据集 Join 和 非倾斜数据集 Join,最后进行 union:</p>
<ol data-id="o01bedff-l6RhRSol">
<li data-id="l20de63f-rifcB6m7">
                对包含少数几个数据量过大的 key 的那个 RDD (假设是 leftRDD),通过 sample 算子采样出一份样本来,然后统计一下每个 key 的数量,计算出来数据量最大的是哪几个 key。具体方法上面已经介绍过了,这里不赘述。</li>
        <li data-id="l20de63f-I0RJJleZ">
                然后将这 k 个 key 对应的数据从 leftRDD 中单独过滤出来,并给每个 key 都打上 1~n 以内的随机数作为前缀,形成一个单独的 leftSkewRDD;而不会导致倾斜的大部分 key 形成另外一个 leftUnSkewRDD。</li>
        <li data-id="l20de63f-K83NE2sB">
                接着将需要 join 的另一个 rightRDD,也过滤出来那几个倾斜 key 并通过 flatMap 操作将该数据集中每条数据均转换为 n 条数据(这 n 条数据都按顺序附加一个 0~n 的前缀),形成单独的 rightSkewRDD;不会导致倾斜的大部分 key 也形成另外一个 rightUnSkewRDD。</li>
        <li data-id="l20de63f-rLkLOhyC">
                现在将 leftSkewRDD 与 膨胀 n 倍的 rightSkewRDD 进行 join,且在 Join 过程中将随机前缀去掉,得到倾斜数据集的 Join 结果 skewedJoinRDD。注意到此时我们已经成功将原先相同的 key 打散成 n 份,分散到多个 task 中去进行 join 了。</li>
        <li data-id="l20de63f-hYnV5EgY">
                对 leftUnSkewRDD 与 rightUnRDD 进行Join,得到 Join 结果 unskewedJoinRDD。</li>
        <li data-id="l20de63f-TdCRP142">
                通过 union 算子将 skewedJoinRDD 与 unskewedJoinRDD 进行合并,从而得到完整的 Join 结果集。</li>
</ol>
<p data-id="pd157317-OEswwW1P">
        <strong>TIPS</strong></p>
<ul data-id="ucd67dc5-doYkCY8j">
<li data-id="l20de63f-QeoOOtWN">
                rightRDD 与倾斜 Key 对应的部分数据,需要与随机前缀集 (1~n) 作笛卡尔乘积 (即将数据量扩大 n 倍),从而保证无论数据倾斜侧倾斜 Key 如何加前缀,都能与之正常 Join。skewRDD 的 join 并行度可以设置为 n * k (k 为 topSkewkey 的个数)。由于倾斜Key与非倾斜Key的操作完全独立,可并行进行。</li>
</ul>
<p data-id="pd157317-AkhkOs8g">
        <strong>(1)适用场景</strong></p>
<p data-id="pd157317-vZF9dT1x">
        两张表都比较大,无法使用 Map 端 Join。其中一个 RDD 有少数几个 Key 的数据量过大,另外一个 RDD 的 Key 分布较为均匀。</p>
<p data-id="pd157317-abxPlc7X">
        <strong>(2)解决方案</strong></p>
<p data-id="pd157317-ggEprxMi">
        将有数据倾斜的 RDD 中倾斜 Key 对应的数据集单独抽取出来加上随机前缀,另外一个 RDD 每条数据分别与随机前缀结合形成新的RDD(相当于将其数据增到到原来的N倍,N即为随机前缀的总个数),然后将二者Join并去掉前缀。然后将不包含倾斜Key的剩余数据进行Join。最后将两次Join的结果集通过union合并,即可得到全部Join结果。</p>
<p data-id="pd157317-7BC4wcMB">
        (3)优势</p>
<p data-id="pd157317-fcyAez1k">
        相对于 Map 则 Join,更能适应大数据集的 Join。如果资源充足,倾斜部分数据集与非倾斜部分数据集可并行进行,效率提升明显。且只针对倾斜部分的数据做数据扩展,增加的资源消耗有限。</p>
<p data-id="pd157317-oVLD3cTs">
        <strong>(4)劣势</strong></p>
<p data-id="pd157317-HaKYG01D">
        如果倾斜 Key 非常多,则另一侧数据膨胀非常大,此方案不适用。而且此时对倾斜 Key 与非倾斜 Key 分开处理,需要扫描数据集两遍,增加了开销。</p>
<h4 data-id="h6e90be6-2Wdjednm" id="h6e90be6-2Wdjednm">
        思路6. 大表 key 加盐,小表扩大 N 倍 jion</h4>
<p data-id="pd157317-AQL868zh">
        如果出现数据倾斜的 Key 比较多,上一种方法将这些大量的倾斜 Key 分拆出来,意义不大。此时更适合直接对存在数据倾斜的数据集全部加上随机前缀,然后对另外一个不存在严重数据倾斜的数据集整体与随机前缀集作笛卡尔乘积(即将数据量扩大N倍)。</p>
<p data-id="pd157317-RBvZYYx4">
        其实就是上一个方法的特例或者简化。少了拆分,也就没有 union。</p>
<p data-id="pd157317-fVQ8EDEj">
        <strong>(1)适用场景</strong></p>
<p data-id="pd157317-OTdxcZpX">
        一个数据集存在的倾斜 Key 比较多,另外一个数据集数据分布比较均匀。</p>
<p data-id="pd157317-mUMNDa17">
        <strong>(2)优势</strong></p>
<p data-id="pd157317-I3B0IIeL">
        对大部分场景都适用,效果不错。</p>
<p data-id="pd157317-omuKGONf">
        <strong>(3)劣势</strong></p>
<p data-id="pd157317-oH2HMurf">
        需要将一个数据集整体扩大 N 倍,会增加资源消耗。</p>
<h4 data-id="h6e90be6-uAFvB8q9" id="h6e90be6-uAFvB8q9">
        思路7. map 端先局部聚合</h4>
<p data-id="pd157317-GPKz5Wv7">
        在 map 端加个 combiner 函数进行局部聚合。加上 combiner 相当于提前进行 reduce ,就会把一个 mapper 中的相同 key 进行聚合,减少 shuffle 过程中数据量 以及 reduce 端的计算量。这种方法可以有效的缓解数据倾斜问题,但是如果导致数据倾斜的 key 大量分布在不同的 mapper 的时候,这种方法就不是很有效了。</p>
<p data-id="pd157317-QbgYqale">
        TIPS 使用 reduceByKey 而不是 groupByKey。</p>
<h4 data-id="h6e90be6-uxE6bmha" id="h6e90be6-uxE6bmha">
        思路8. 加盐局部聚合 + 去盐全局聚合</h4>
<p data-id="pd157317-mXsYtiuS">
        这个方案的核心实现思路就是进行两阶段聚合。第一次是局部聚合,先给每个 key 都打上一个 1~n 的随机数,比如 3 以内的随机数,此时原先一样的 key 就变成不一样的了,比如 (hello, 1) (hello, 1) (hello, 1) (hello, 1) (hello, 1),就会变成 (1_hello, 1) (3_hello, 1) (2_hello, 1) (1_hello, 1) (2_hello, 1)。接着对打上随机数后的数据,执行 reduceByKey 等聚合操作,进行局部聚合,那么局部聚合结果,就会变成了 (1_hello, 2) (2_hello, 2) (3_hello, 1)。然后将各个 key 的前缀给去掉,就会变成 (hello, 2) (hello, 2) (hello, 1),再次进行全局聚合操作,就可以得到最终结果了,比如 (hello, 5)。</p>
<pre>
<span class="cm-variable">def</span> <span class="cm-variable">antiSkew</span>(): <span class="cm-variable">RDD</span>[(<span class="cm-variable">String</span>, <span class="cm-variable">Int</span>)] <span class="cm-operator">=</span> { <span class="cm-property">val</span> <span class="cm-variable">SPLIT</span> <span class="cm-operator">=</span> <span class="cm-string">"-"</span> <span class="cm-variable">val</span> <span class="cm-variable">prefix</span> <span class="cm-operator">=</span> <span class="cm-keyword">new</span> <span class="cm-variable">Random</span>().<span class="cm-variable">nextInt</span>(<span class="cm-number">10</span>) <span class="cm-variable">pairs</span>.<span class="cm-property">map</span>(<span class="cm-def">t</span> <span class="cm-operator">=&gt;</span> ( <span class="cm-variable">prefix</span> <span class="cm-operator">+</span> <span class="cm-variable">SPLIT</span> <span class="cm-operator">+</span> <span class="cm-variable-2">t</span>.<span class="cm-property">_1</span>, <span class="cm-number">1</span>))
      .<span class="cm-property">reduceByKey</span>((<span class="cm-def">v1</span>, <span class="cm-def">v2</span>) <span class="cm-operator">=&gt;</span> <span class="cm-variable-2">v1</span> <span class="cm-operator">+</span> <span class="cm-variable-2">v2</span>)
      .<span class="cm-property">map</span>(<span class="cm-def">t</span> <span class="cm-operator">=&gt;</span> (<span class="cm-variable-2">t</span>.<span class="cm-property">_1</span>.<span class="cm-property">split</span>(<span class="cm-variable">SPLIT</span>)(<span class="cm-number">1</span>), <span class="cm-variable">t2</span>.<span class="cm-property">_2</span>))
      .<span class="cm-property">reduceByKey</span>((<span class="cm-def">v1</span>, <span class="cm-def">v2</span>) <span class="cm-operator">=&gt;</span> <span class="cm-variable-2">v1</span> <span class="cm-operator">+</span> <span class="cm-variable-2">v2</span>)
}</pre>
<p data-id="pd157317-V4wXdn2k">
        不过进行两次 mapreduce,性能稍微比一次的差些。</p>
<h3 data-id="h26976cb-8I7KQRZb" id="h26976cb-8I7KQRZb">
        七、Hadoop 中的数据倾斜</h3>
<p data-id="pd157317-kuX44hp0">
        Hadoop 中直接贴近用户使用的是 Mapreduce 程序和 Hive 程序,虽说 Hive 最后也是用 MR 来执行(至少目前 Hive 内存计算并不普及),但是毕竟写的内容逻辑区别很大,一个是程序,一个是Sql,因此这里稍作区分。</p>
<p data-id="pd157317-DgkiZv3d">
        Hadoop 中的数据倾斜主要表现在 ruduce 阶段卡在99.99%,一直99.99%不能结束。</p>
<p data-id="pd157317-UN08VNXz">
        这里如果详细的看日志或者和监控界面的话会发现:</p>
<ol data-id="o01bedff-pWGNYmmk">
<li data-id="l20de63f-nyoR4eCB">
                有一个多几个 reduce 卡住</li>
        <li data-id="l20de63f-iGgIZMaM">
                各种 container报错 OOM</li>
        <li data-id="l20de63f-iGEbodJJ">
                读写的数据量极大,至少远远超过其它正常的 reduce</li>
        <li data-id="l20de63f-zW60aRrE">
                伴随着数据倾斜,会出现任务被 kill 等各种诡异的表现。</li>
</ol>
<p data-id="pd157317-GEnSuZfK">
        <strong>经验:</strong> Hive的数据倾斜,一般都发生在 Sql 中 Group 和 On 上,而且和数据逻辑绑定比较深。</p>
<p data-id="pd157317-rgW47dSD">
        <strong>优化方法</strong></p>
<ol data-id="o5cd1cdd-begSnTA7">
<li data-id="l20de63f-adXZXWP0">
                这里列出来一些方法和思路,具体的参数和用法在官网看就行了。</li>
        <li data-id="l20de63f-BnzGoNy1">
                map join 方式</li>
        <li data-id="l20de63f-KGXXwPWo">
                count distinct 的操作,先转成 group,再 count</li>
        <li data-id="l20de63f-vFh6GfzX">
                参数调优</li>
        <li data-id="l20de63f-NTROeadw">
                set hive.map.aggr=true</li>
        <li data-id="l20de63f-Y9wK9Rzb">
                set hive.groupby.skewindata=true</li>
        <li data-id="l20de63f-fNpM1lS3">
                left semi jion 的使用</li>
        <li data-id="l20de63f-fwQfTSEB">
                设置 map 端输出、中间结果压缩。(不完全是解决数据倾斜的问题,但是减少了 IO 读写和网络传输,能提高很多效率)</li>
</ol>
<p data-id="pd157317-K1Rzcnup">
        <strong>说明</strong></p>
<p data-id="pd157317-RRnZlSpv">
        hive.map.aggr=true: 在map中会做部分聚集操作,效率更高但需要更多的内存。</p>
<p data-id="pd157317-FE9nuPOz">
        hive.groupby.skewindata=true: 数据倾斜时负载均衡,当选项设定为true,生成的查询计划会有两个MRJob。第一个MRJob 中,Map的输出结果集合会随机分布到Reduce中,每个Reduce做部分聚合操作,并输出结果,这样处理的结果是相同的GroupBy Key有可能被分发到不同的Reduce中,从而达到负载均衡的目的;第二个MRJob再根据预处理的数据结果按照GroupBy Key分布到Reduce中(这个过程可以保证相同的GroupBy Key被分布到同一个Reduce中),最后完成最终的聚合操作。</p>
<p>
        原文地址:https://www.toutiao.com/a7067694657785905679/</p>
頁: [1]
查看完整版本: Spark 数据倾斜及其解决方案