0%

Flink-connector-hippo

broker分拆

获取子任务的index

1
int taskId = this.getRuntimeContext().getIndexOfThisSubtask();

checkpoint

实现CheckpointedFunction

在ListState中保存每个broker的偏移量

1
ListState<Tuple2<String, String>> offsetState;

watermark生成

hippo pullConsumer

1
2
3
4
5
6
7
ConsumerConfig config =
new ConsumerConfig(masterAddress, consumerGroup);
if (!isRestored && bootstrapFromMax) {
config.setConsumeFromMax(true);
}

messagePullConsumer = new PullMessageConsumer(config);

子任务的checkpointLock

往下游放入消息必须加锁

1
SourceContext<byte[]>.getCheckpointLock()

架构

The processes involved in executing a Flink dataflow

当 Flink 集群启动后,首先会启动一个 JobManger 和一个或多个的 TaskManager。
由 Client 提交任务给 JobManager,JobManager 再调度任务到各个 TaskManager 去执行,然后 TaskManager 将心跳和统计信息汇报给 JobManager。
TaskManager 之间以流的形式进行数据的传输。
上述三者均为独立的 JVM 进程。

  • Client 为提交 Job 的客户端,可以是运行在任何机器上(与 JobManager 环境连通即可)。提交 Job 后,Client 可以结束进程(Streaming的任务),也可以不结束并等待结果返回。
  • JobManager 主要负责调度 Job 并协调 Task 做 checkpoint,职责上很像 Storm 的 Nimbus。从 Client 处接收到 Job 和 JAR 包等资源后,会生成优化后的执行计划,并以 Task 的单元调度到各个 TaskManager 去执行。
  • TaskManager 在启动的时候就设置好了槽位数(Slot),每个 slot 能启动一个 Task,Task 为线程。从 JobManager 处接收需要部署的 Task,部署启动后,与自己的上游建立 Netty 连接,接收数据并处理。

可以看到 Flink 的任务调度是多线程模型,并且不同Job/Task混合在一个 TaskManager 进程中。虽然这种方式可以有效提高 CPU 利用率,但是个人不太喜欢这种设计,因为不仅缺乏资源隔离机制,同时也不方便调试。类似 Storm 的进程模型,一个JVM 中只跑该 Job 的 Tasks 实际应用中更为合理。

Graph转换

看起来有点乱,怎么有这么多不一样的图。实际上,还有更多的图。Flink 中的执行图可以分成四层:StreamGraph -> JobGraph -> ExecutionGraph -> 物理执行图。

  • StreamGraph: 是根据用户通过 Stream API 编写的代码生成的最初的图。用来表示程序的拓扑结构。
  • JobGraph: StreamGraph经过优化后生成了 JobGraph,提交给 JobManager 的数据结构。主要的优化为,将多个符合条件的节点 chain 在一起作为一个节点,这样可以减少数据在节点之间流动所需要的序列化/反序列化/传输消耗。
  • ExecutionGraph: JobManager 根据 JobGraph 生成ExecutionGraph。ExecutionGraph是JobGraph的并行化版本,是调度层最核心的数据结构。
  • 物理执行图: JobManager 根据 ExecutionGraph 对 Job 进行调度后,在各个TaskManager 上部署 Task 后形成的“图”,并不是一个具体的数据结构。

例如上文中的2个并发度(Source为1个并发度)的 SocketTextStreamWordCount 四层执行图的演变过程如下图所示(点击查看大图):

img

这里对一些名词进行简单的解释。

  • StreamGraph:

    根据用户通过 Stream API 编写的代码生成的最初的图。

    • StreamNode:用来代表 operator 的类,并具有所有相关的属性,如并发度、入边和出边等。
    • StreamEdge:表示连接两个StreamNode的边。
  • JobGraph:

    StreamGraph经过优化后生成了 JobGraph,提交给 JobManager 的数据结构。

    • JobVertex:经过优化后符合条件的多个StreamNode可能会chain在一起生成一个JobVertex,即一个JobVertex包含一个或多个operator,JobVertex的输入是JobEdge,输出是IntermediateDataSet。
    • IntermediateDataSet:表示JobVertex的输出,即经过operator处理产生的数据集。producer是JobVertex,consumer是JobEdge。
    • JobEdge:代表了job graph中的一条数据传输通道。source 是 IntermediateDataSet,target 是 JobVertex。即数据通过JobEdge由IntermediateDataSet传递给目标JobVertex。
  • ExecutionGraph:

    JobManager 根据 JobGraph 生成ExecutionGraph。ExecutionGraph是JobGraph的并行化版本,是调度层最核心的数据结构。

    • ExecutionJobVertex:和JobGraph中的JobVertex一一对应。每一个ExecutionJobVertex都有和并发度一样多的 ExecutionVertex。
    • ExecutionVertex:表示ExecutionJobVertex的其中一个并发子任务,输入是ExecutionEdge,输出是IntermediateResultPartition。
    • IntermediateResult:和JobGraph中的IntermediateDataSet一一对应。一个IntermediateResult包含多个IntermediateResultPartition,其个数等于该operator的并发度。
    • IntermediateResultPartition:表示ExecutionVertex的一个输出分区,producer是ExecutionVertex,consumer是若干个ExecutionEdge。
    • ExecutionEdge:表示ExecutionVertex的输入,source是IntermediateResultPartition,target是ExecutionVertex。source和target都只能是一个。
    • Execution:是执行一个 ExecutionVertex 的一次尝试。当发生故障或者数据需要重算的情况下 ExecutionVertex 可能会有多个 ExecutionAttemptID。一个 Execution 通过 ExecutionAttemptID 来唯一标识。JM和TM之间关于 task 的部署和 task status 的更新都是通过 ExecutionAttemptID 来确定消息接受者。
  • 物理执行图:

    JobManager 根据 ExecutionGraph 对 Job 进行调度后,在各个TaskManager 上部署 Task 后形成的“图”,并不是一个具体的数据结构。

    • Task:Execution被调度后在分配的 TaskManager 中启动对应的 Task。Task 包裹了具有用户执行逻辑的 operator。
    • ResultPartition:代表由一个Task的生成的数据,和ExecutionGraph中的IntermediateResultPartition一一对应。
    • ResultSubpartition:是ResultPartition的一个子分区。每个ResultPartition包含多个ResultSubpartition,其数目要由下游消费 Task 数和 DistributionPattern 来决定。
    • InputGate:代表Task的输入封装,和JobGraph中JobEdge一一对应。每个InputGate消费了一个或多个的ResultPartition。
    • InputChannel:每个InputGate会包含一个以上的InputChannel,和ExecutionGraph中的ExecutionEdge一一对应,也和ResultSubpartition一对一地相连,即一个InputChannel接收一个ResultSubpartition的输出。

img

首先我们看到,JobGraph 之上除了 StreamGraph 还有 OptimizedPlan。OptimizedPlan 是由 Batch API 转换而来的。StreamGraph 是由 Stream API 转换而来的。为什么 API 不直接转换成 JobGraph?因为,Batch 和 Stream 的图结构和优化方法有很大的区别,比如 Batch 有很多执行前的预分析用来优化图的执行,而这种优化并不普适于 Stream,所以通过 OptimizedPlan 来做 Batch 的优化会更方便和清晰,也不会影响 Stream。JobGraph 的责任就是统一 Batch 和 Stream 的图,用来描述清楚一个拓扑图的结构,并且做了 chaining 的优化,chaining 是普适于 Batch 和 Stream 的,所以在这一层做掉。ExecutionGraph 的责任是方便调度和各个 tasks 状态的监控和跟踪,所以 ExecutionGraph 是并行化的 JobGraph。而“物理执行图”就是最终分布式在各个机器上运行着的tasks了。所以可以看到,这种解耦方式极大地方便了我们在各个层所做的工作,各个层之间是相互隔离的。

后续的文章,将会详细介绍 Flink 是如何生成这些执行图的。由于我目前关注 Flink 的流处理功能,所以主要有以下内容:

  1. 如何生成 StreamGraph
  2. 如何生成 JobGraph
  3. 如何生成 ExecutionGraph
  4. 如何进行调度(如何生成物理执行图)

资源共享链与资源共享组

当我们编写完一个Flink程序,从Client开始执行——>JobManager——>TaskManager——>Slot启动并执行Task的过程中,会对我们提交的执行计划进行优化,其中有两个比较重要的优化过程是:任务链与处理槽共享组,前者是对执行效率的优化,后者是对内存资源的优化。

Graph转换

StreamGraph转换为JobGraph过程中,关键在于将多个 StreamNode 优化为一个 JobVertex,对应的 StreamEdge 则转化为 JobEdge,并且 JobVertexJobEdge 之间通过 IntermediateDataSet (中间数据集)形成一个生产者和消费者的连接关系。每个JobVertex就是JobManger的一个任务调度单位(任务Task)。

为了避免在这个过程中将关联性很强的几个StreamNode(算子)放到不同JobVertexTask)中,从而导致因为Task执行产生的效率问题(数据交换(网络传输)、线程上下文切换),Flink会在StreamGraph转换为JobGraph过程中将可以优化的算子合并为一个算子链(也就是形成一个Task)。这样就可以把这条链上的算子放到一个线程中去执行,这样就提高了任务执行效率。

作业链Task-chain

img

  • Chain:Flink会尽可能地将多个operator链接(chain)在一起形成一个task pipline。每个task pipline在一个线程中执行

  • 优点:它能减少线程之间的切换,减少消息的序列化/反序列化,减少数据在缓冲区的交换,减少了延迟的同时提高整体的吞吐量。

    img

StreamGraph转换为JobGraph过程中,实际上是逐条审查每一个StreamEdge和该SteamEdge两头连接的两个StreamNode的特性,来决定该StreamEdge两头的StreamNode是不是可以合并在一起形成算子链。这个判断过程flink给出了明确的规则,我们看一下StreamingJobGraphGenerator中的isChainable()方法:

img

该方法返回true时两个端点才可以合并到一起,根据源码我们可以得出形成作业链的规则如下:

  1. 上下游的并行度一致(槽一致)
  2. 该节点必须要有上游节点跟下游节点;
  3. 下游StreamNode的输入StreamEdge只能有一个)
  4. 上下游节点都在同一个 slot group 中(下面会解释 slot group)
  5. 下游节点的 chain 策略为 ALWAYS(可以与上下游链接,map、flatmap、filter等默认是ALWAYS)
  6. 上游节点的 chain 策略为 ALWAYS 或 HEAD(只能与下游链接,不能与上游链接,Source默认是HEAD)
  7. 上下游算子之间没有数据shuffle (数据分区方式是 forward)
  8. 用户没有禁用 chain

二、开启/禁用全局作业链

用户能够通过禁用全局作业链的操作来关闭整个Flink的作业链,但是这个操作会影响到这个作业的执行情况,除非我们非常清楚作业的执行过程,否则不建议这么做:StreamExecutionEnvironment.disableOperatorChaining()。全局作业链关闭之后,如果想创建对应Operator的作业链,可以使用startNewChain()方法:someStream.filter(...).map(...).startNewChain().map(...)。注意该方法只对当前操作符及之后的操作符有效,所以上述代码只对两个map进行链条绑定。

三、禁用局部作业链

如果我们只想对某个算子执行禁用作业链,只需调用disableChaining()方法:someSteam.map().disableChaining().filter(),该方法只会禁用当前算子的链条(上述代码中就是map),对其他算子操作不产生影响。

处理槽共享组(出于某中目的将多个Task放到同一个slot中执行)

一、Task Slot

TaskManager 是一个 JVM 进程,并会以独立的线程来执行一个task。为了控制一个 TaskManager 能接受多少个 task,Flink 提出了 Task Slot 的概念,通过 Task Slot 来定义Flink 中的计算资源。solt 对TaskManager内存进行平均分配,每个solt内存都相同,加起来的和等于TaskManager可用内存,但是仅仅对内存做了隔离,并没有对cpu进行隔离。将资源 slot 化意味着来自不同job的task不会为了内存而竞争,而是每个task都拥有一定数量的内存储备。

通过调整 task slot 的数量,用户可以定义task之间是如何相互隔离的。每个 TaskManager 有一个slot,也就意味着每个task运行在独立的 JVM 中。每个 TaskManager 有多个slot的话,也就是说多个task运行在同一个JVM中。而在同一个JVM进程中的task,可以共享TCP连接(基于多路复用)和心跳消息,可以减少数据的网络传输。也能共享一些数据结构,一定程度上减少了每个task的消耗。

二、共享槽

一个TaskManager中至少有一个插槽slot,每个插槽均分内存并且之间是内存隔离的,但是共享CPU。算子根据计算复杂度可以分为资源密集型与非资源密集型算子(可以认为有的算子计算时内存需求大,有些算子内存需求小)。现在有这么个情况:某个Job下的Tasks中既有资源密集型Task(A),又有非资源密集型Task(B),他们被分到不同的slot上,这就会产生一个问题,有的slot内存使用率大,有的slot内存使用率小,这样就很不公平,内存没有得到充分的利用。所以我们可以采用一个方案:将A、B放到同一个slot当中。

默认情况下,Flink 允许subtasks共享slot,条件是它们都来自同一个Job的不同task的subtask。结果可能一个slot持有该job的整个pipeline。允许槽共享,会有以下两个方面的好处:

  • flink计算一个job所需slot数量时,只需要确定所其最大并行度(前提,保持默认SlotSharingGroup),而不用计算每一个任务的并行度的总和;
  • 能更好的利用资源,如果没有solt共享,那些资源需求不大的map子任务将和资源需求更大的window占用相同的资源。

img

Flink相同资源组里的多个Task可以共享一个Slot资源槽。具体共享机制又分两种:

1、CoLocationGroup: 强制将 subtasks 放到同一个 slot 中,是一种硬约束

  • 保证把JobVertices的第n个运行实例和其他相同组内的JobVertices第n个实例运作在相同的slot中(所有的并行度相同的subTasks运行在同一个slot );
  • 主要用于迭代流(训练机器学习模型) ,用来保证迭代头与迭代尾的第i个subtask能被调度到同一个TaskManager上。

2、SlotSharingGroup: 允许不同的JobVertices的部署在相同的Slot中,但这是一种宽约束,只是尽量做到不能完全保证。

  • SlotSharingGroup是Flink中用来实现slot共享的类,它尽可能地让subTasks共享一个slot;
  • 保证同一个group的并行度相同的sub-tasks 共享同一个slots ;
  • 算子的默认group为default(即默认一个job下的subtask都可以共享一个slot)
  • 为了防止不合理的共享,用户可以强制指定operator的共享组,比如: someStream.filter(...).slotSharingGroup("group1");就强制指定了filter的slot共享组为group1;
  • 要想确定一个未做SlotSharingGroup设置的算子的group是什么,可以根据上游算子的 group 和自身是否设置 group共同确定;
  • 适当设置可以减少每个slot运行的线程数,从而整体上减少机?的负载。

【参考文献】

  1. Apache Flink进阶一: Runtime核心机制剖析

窗口会自动管理状态和触发计算,Flink 提供了丰富的窗口函数来进行计算。主要包括以下两种:

  • ProcessWindowFunction,全量计算会把所有数据缓存到状态里,一直到窗口结束时统一计算。相对来说,状态会比较大,计算效率也会低一些;
  • AggregateFunction,增量计算就是来一条数据就算一条,可能我们的状态就会特别的小,计算效率也会比 ProcessWindowFunction 高很多,但是如果状态存储在磁盘频繁访问状态可能会影响性能。

0.1 窗口的触发

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
SingleOutputStreamOperator<ItemEntity> streamOperator = env
.socketTextStream("127.0.0.1", 9091)
.map(new MapFunction<String, ItemEntity>() {
@Override
public ItemEntity map(String s) throws Exception {
String[] split = null;
if (s.isEmpty() || (split = s.split(",")).length != 2) {
return null;
}
ItemEntity itemEntity = ItemEntity.builder().timestamp(split[0])
.eventId(split[1])
.build();
return itemEntity;
}
})
.filter(Objects::nonNull)
.assignTimestampsAndWatermarks(new AscendingTimestampExtractor<ItemEntity>() {
@Override
public long extractAscendingTimestamp(ItemEntity itemEntity) {
long timestamp = itemEntity.getTimestamp();
return timestamp;
}
});
streamOperator
.keyBy(new KeySelector<ItemEntity, String>() {
@Override
public String getKey(ItemEntity itemEntity) throws Exception {
String eventId = itemEntity.getEventId();
return eventId;
}
})
.window(TumblingEventTimeWindows.of(Time.seconds(5)))
.process(
new ProcessWindowFunction<ItemEntity, Tuple2<String, String>, String, TimeWindow>() {
@Override
public void process(String s, Context context,
Iterable<ItemEntity> iterable,
Collector<Tuple2<String, String>> collector) throws Exception {
for (ItemEntity itemEntity : iterable) {
long timestamp = itemEntity.getTimestamp();
Date date = new Date(timestamp);
SimpleDateFormat simpleDateFormat = new SimpleDateFormat(
"yyyy-MM-dd HH:mm:ss");
collector.collect(Tuple2.of(itemEntity.getEventId(),
timestamp + "->" + simpleDateFormat.format(date)));
}
}
})
.print();
env.execute("test-window");

5秒的窗口

1
2
3
4
5
6
7
8
9
10
11
12
13
2020-04-01 00:01:00,1
2020-04-01 00:01:00,1
2020-04-01 00:01:06,1 # 触发计算
(1,1585670460000->2020-04-01 00:01:00)
(1,1585670460000->2020-04-01 00:01:00)
2020-04-01 00:01:00,1 # 窗口已经关闭,旧数据
# WARN org.apache.flink.streaming.api.functions.timestamps.AscendingTimestampExtractor [] - Timestamp monotony violated: 1585670460000 < 1585670466000
2020-04-01 00:01:06,1
2020-04-01 00:01:07,1
2020-04-01 00:01:11,1 # 触发计算
(1,1585670466000->2020-04-01 00:01:06)
(1,1585670466000->2020-04-01 00:01:06)
(1,1585670467000->2020-04-01 00:01:07)

0.2 window的抽象概念

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
Keyed Windows

stream
.keyBy(...) <- keyed versus non-keyed windows
.window(...) <- required: "assigner"
[.trigger(...)] <- optional: "trigger" (else default trigger)
[.evictor(...)] <- optional: "evictor" (else no evictor)
[.allowedLateness(...)] <- optional: "lateness" (else zero)
[.sideOutputLateData(...)] <- optional: "output tag" (else no side output for late data)
.reduce/aggregate/fold/apply() <- required: "function"
[.getSideOutput(...)] <- optional: "output tag"
Non-Keyed Windows

stream
.windowAll(...) <- required: "assigner"
[.trigger(...)] <- optional: "trigger" (else default trigger)
[.evictor(...)] <- optional: "evictor" (else no evictor)
[.allowedLateness(...)] <- optional: "lateness" (else zero)
[.sideOutputLateData(...)] <- optional: "output tag" (else no side output for late data)
.reduce/aggregate/fold/apply() <- required: "function"
[.getSideOutput(...)] <- optional: "output tag"

0.2.1 window assigner

0.2.2 window trigger

0.2.3 window evictor

0.3 windowOperator工作流程

0.3.1 window state

0.4 Session window

Flink 原理与实现:Session Window

SESSION(time_attr, interval)定义一个会话时间窗口。
会话时间窗口没有一个固定的持续时间,但是它们的边界会根据 interval 所定义的不活跃时间所确定;即一个会话时间窗口在定义的间隔时间内没有时间出现,该窗口会被关闭。例如时间窗口的间隔时间是 30 分钟,当其不活跃的时间达到30分钟后,若观测到新的记录,则会启动一个新的会话时间窗口(否则该行数据会被添加到当前的窗口),且若在 30 分钟内没有观测到新纪录,这个窗口将会被关闭。会话时间窗口可以使用事件时间(批处理、流处理)或处理时间(流处理)。

流式数据处理中,很多操作要依赖于时间属性进行,因此时间属性也是流式引擎能够保证准确处理数据的基石。在这篇文章中,我们将对 Flink 中时间属性和窗口的实现逻辑进行分析。

0.5 all window operator’s parallelism is 1

If the parallelism of the environment is set to 3 and you are using a WindowAll operator, only the window operator runs in parallelism 1. The sink will still be running with parallelism 3. Hence, the plan looks as follows:

1
2
3
In_1 -\               /- Out_1
In_2 --- WindowAll_1 --- Out_2
In_3 -/ \- Out_3

The WindowAll operator emits its output to its subsequent tasks using a round-robin strategy. That’s the reason for the different threads emitting the result records of program.

When you set the environment parallelism to 1, all operators run with a single task.

keyed-windown —分布式计算—> all-windown

【参考文献】

  1. flink原理与实现: window机制
  2. flink窗口应用与实现
  3. Flink原理: 窗口原理详解
  4. Flink滑动窗口原理与细粒度滑动窗口的性能问题

1 实时流处理系统反压机制(BackPressure)综述[转]

发表于 2018-11-15 | 更新于 2018-12-03 | 分类于 BigData | 阅读次数 333

本文主要关于实时流处理系统反压机制,最近看到反压问题看到此文章很好,在此分享并mark一下。

本文主要关于实时流处理系统反压机制,最近看到反压问题看到此文章很好,在此分享并mark一下。
(¬_¬)ノ最近菜叶子没自己写见谅。
本文转自 实时流处理系统反压机制(BackPressure)综述
https://blog.csdn.net/qq_21125183/article/details/80708142
开启Back Pressure使生产环境的Spark Streaming应用更稳定、有效

反压机制(BackPressure)被广泛应用到实时流处理系统中,流处理系统需要能优雅地处理反压(backpressure)问题。
反压通常产生于这样的场景:短时负载高峰导致系统接收数据的速率远高于它处理数据的速率。
许多日常问题都会导致反压,例如,垃圾回收停顿可能会导致流入的数据快速堆积,或者遇到大促或秒杀活动导致流量陡增。
反压如果不能得到正确的处理,可能会导致资源耗尽甚至系统崩溃。反压机制就是指系统能够自己检测到被阻塞的Operator,然后系统自适应地降低源头或者上游的发送速率。

目前主流的流处理系统 Apache Storm、JStorm、Spark Streaming、S4、Apache Flink、Twitter Heron都采用反压机制解决这个问题,不过他们的实现各自不同。

实时流处理系统反压机制01

不同的组件可以不同的速度执行(并且每个组件中的处理速度随时间改变)。 例如,考虑一个工作流程,或由于数据倾斜或任务调度而导致数据被处理十分缓慢。
在这种情况下,如果上游阶段不减速,将导致缓冲区建立长队列(队列占用内存、硬盘空间,节点负载加重),或导致系统丢弃元组。
如果元组在中途丢弃,那么效率可能会有损失,因为已经为这些元组产生的计算被浪费了。
并且在一些流处理系统中比如Strom,会将这些丢失的元组重新发送,这样会导致数据的一致性问题(at least once语义),并且还会导致某些Operator状态叠加。
进而整个程序输出结果不准确。第二由于系统接收数据的速率是随着时间改变的,短时负载高峰导致系统接收数据的速率远高于它处理数据的速率的情况,也会导致Tuple在中途丢失。
所以实时流处理系统必须能够解决发送速率远大于系统能处理速率这个问题,大多数实时流处理系统采用反压(BackPressure)机制解决这个问题。

下面我们就来介绍一下不同的实时流处理系统采用的反压机制:

2 Strom 反压机制

2.1 Storm 1.0 以前的反压机制

对于开启了acker机制的storm程序,可以通过设置conf.setMaxSpoutPending参数来实现反压效果,如果下游组件(bolt)处理速度跟不上导致spout发送的tuple没有及时确认的数超过了参数设定的值,spout会停止发送数据,这种方式的缺点是很难调优conf.setMaxSpoutPending参数的设置以达到最好的反压效果,设小了会导致吞吐上不去,设大了会导致worker OOM;有震荡,数据流会处于一个颠簸状态,效果不如逐级反压;另外对于关闭acker机制的程序无效;

2.2 Storm Automatic Backpressure

新的storm自动反压机制(Automatic Back Pressure)通过监控bolt中的接收队列的情况,当超过高水位值时专门的线程会将反压信息写到 Zookeeper ,Zookeeper上的watch会通知该拓扑的所有Worker都进入反压状态,最后Spout降低tuple发送的速度。

实时流处理系统反压机制02

每个Executor都有一个接受队列和发送队列用来接收Tuple和发送Spout或者Bolt生成的Tuple元组。每个Worker进程都有一个单的的接收线程监听接收端口。
它从每个网络上进来的消息发送到Executor的接收队列中。Executor接收队列存放Worker或者Worker内部其他Executor发过来的消息。
Executor工作线程从接收队列中拿出数据,然后调用execute方法,发送Tuple到Executor的发送队列。
Executor的发送线程从发送队列中获取消息,按照消息目的地址选择发送到Worker的传输队列中或者其他Executor的接收队列中。
最后Worker的发送线程从传输队列中读取消息,然后将Tuple元组发送到网络中。

  1. 当Worker进程中的Executor线程发现自己的接收队列满了时,也就是接收队列达到high watermark的阈值后,因此它会发送通知消息到背压线程。
  2. 背压线程将当前worker进程的信息注册到Zookeeper的Znode节点中。具体路径就是 /Backpressure/topo1/wk1
  3. Zookeepre的Znode Watcher监视/Backpreesure/topo1下的节点目录变化情况,如果发现目录增加了znode节点说明或者其他变化。这就说明该Topo1需要反压控制,然后它会通知Topo1所有的Worker进入反压状态。
  4. 最终Spout降低tuple发送的速度。

3 JStorm 反压机制

JStorm做了两级的反压,第一级和Jstorm类似,通过执行队列来监测,但是不会通过ZK来协调,而是通过Topology Master来协调。
在队列中会标记high water mark和low water mark,当执行队列超过high water mark时,就认为bolt来不及处理,则向TM发一条控制消息,上游开始减慢发送速率,直到下游低于low water mark时解除反压。

此外,在Netty层也做了一级反压,由于每个Worker Task都有自己的发送和接收的缓冲区,可以对缓冲区设定限额、控制大小,如果spout数据量特别大,缓冲区填满会导致下游bolt的接收缓冲区填满,造成了反压。

实时流处理系统反压机制03

限流机制:jstorm的限流机制, 当下游bolt发生阻塞时, 并且阻塞task的比例超过某个比例时(现在默认设置为0.1),触发反压

限流方式:计算阻塞Task的地方执行线程执行时间,Spout每发送一个tuple等待相应时间,然后讲这个时间发送给Spout, 于是, spout每发送一个tuple,就会等待这个执行时间。

Task阻塞判断方式:在jstorm 连续4次采样周期中采样,队列情况,当队列超过80%(可以设置)时,即可认为该task处在阻塞状态。

4 SparkStreaming 反压机制

4.1 为什么引入反压机制Backpressure

默认情况下,Spark Streaming通过Receiver以生产者生产数据的速率接收数据,计算过程中会出现batch processing time > batch interval的情况,其中batch processing time 为实际计算一个批次花费时间, batch interval为Streaming应用设置的批处理间隔。
这意味着Spark Streaming的数据接收速率高于Spark从队列中移除数据的速率,也就是数据处理能力低,在设置间隔内不能完全处理当前接收速率接收的数据。如果这种情况持续过长的时间,会造成数据在内存中堆积,导致Receiver所在Executor内存溢出等问题(如果设置StorageLevel包含disk, 则内存存放不下的数据会溢写至disk, 加大延迟)。
Spark 1.5以前版本,用户如果要限制Receiver的数据接收速率,可以通过设置静态配制参数“spark.streaming.receiver.maxRate”的值来实现,此举虽然可以通过限制接收速率,来适配当前的处理能力,防止内存溢出,但也会引入其它问题。比如:producer数据生产高于maxRate,当前集群处理能力也高于maxRate,这就会造成资源利用率下降等问题。为了更好的协调数据接收速率与资源处理能力,Spark Streaming 从v1.5开始引入反压机制(back-pressure),通过动态控制数据接收速率来适配集群数据处理能力。

4.2 反压机制Backpressure

Spark Streaming Backpressure: 根据JobScheduler反馈作业的执行信息来动态调整Receiver数据接收率。通过属性“spark.streaming.backpressure.enabled”来控制是否启用backpressure机制,默认值false,即不启用。

1
sparkConf.set("spark.streaming.backpressure.enabled",”true”)

SparkStreaming 架构图如下所示:

实时流处理系统反压机制04

SparkStreaming 反压过程执行如下图所示:

在原架构的基础上加上一个新的组件RateController,这个组件负责监听“OnBatchCompleted”事件,然后从中抽取processingDelayschedulingDelay信息. Estimator依据这些信息估算出最大处理速度(rate),最后由基于Receiver的Input Stream将rate通过ReceiverTracker与ReceiverSupervisorImpl转发给BlockGenerator(继承自RateLimiter).

实时流处理系统反压机制05

4.3 direct模式-BackPressure(此部分详细说明了direct模式接收:转自-开启Back Pressure使生产环境的Spark Streaming应用更稳定、有效)

当Spark Streaming与Kafka使用Direct API集群时,我们可以很方便的去控制最大数据摄入量–通过一个被称作spark.streaming.kafka.maxRatePerPartition的参数。根据文档描述,他的含义是:Direct API读取每一个Kafka partition数据的最大速率(每秒读取的消息量)。
配置项spark.streaming.kafka.maxRatePerPartition,对防止流式应用在下边两种情况下出现流量过载时尤其重要:
1.Kafka Topic中有大量未处理的消息,并且我们设置是Kafka auto.offset.reset参数值为smallest,他可以防止第一个批次出现数据流量过载情况。
2.当Kafka 生产者突然飙升流量的时候,他可以防止批次处理出现数据流量过载情况。

但是,配置Kafka每个partition每批次最大的摄入量是个静态值,也算是个缺点。随着时间的变化,在生产环境运行了一段时间的Spark Streaming应用,每批次每个Kafka partition摄入数据最大量的最优值也是变化的。有时候,是因为消息的大小会变,导致数据处理时间变化。有时候,是因为流计算所使用的多租户集群会变得非常繁忙,比如在白天时候,一些其他的数据应用(例如Impala/Hive/MR作业)竞争共享的系统资源时(CPU/内存/网络/磁盘IO)。
背压机制可以解决该问题。背压机制是呼声比较高的功能,他允许根据前一批次数据的处理情况,动态、自动的调整后续数据的摄入量,这样的反馈回路使得我们可以应对流式应用流量波动的问题。
Spark Streaming的背压机制是在Spark1.5版本引进的,我们可以添加如下代码启用改功能:

1
2
sparkConf.set("spark.streaming.backpressure.enabled",”true”)

那应用启动后的第一个批次流量怎么控制呢?因为他没有前面批次的数据处理时间,所以没有参考的数据去评估这一批次最优的摄入量。在Spark官方文档中有个被称作spark.streaming.backpressure.initialRate的配置,看起来是控制开启背压机制时初始化的摄入量。其实不然,该参数只对receiver模式起作用,并不适用于direct模式。推荐的方法是使用spark.streaming.kafka.maxRatePerPartition控制背压机制起作用前的第一批次数据的最大摄入量。我通常建议设置spark.streaming.kafka.maxRatePerPartition的值为最优估计值的1.5到2倍,让背压机制的算法去调整后续的值。请注意,spark.streaming.kafka.maxRatePerPartition的值会一直控制最大的摄入量,所以背压机制的算法值不会超过他。
另一个需要注意的是,在第一个批次处理完成前,紧接着的批次都将使用spark.streaming.kafka.maxRatePerPartition的值作为摄入量。通过Spark UI可以看到,批次间隔为5s,当批次调度延迟31秒时候,前7个批次的摄入量是20条记录。直到第八个批次,背压机制起作用时,摄入量变为5条记录。

5 Heron 反压机制

实时流处理系统反压机制06

当下游处理速度跟不上上游发送速度时,一旦StreamManager 发现一个或多个Heron Instance 速度变慢,立刻对本地spout进行降级,降低本地Spout发送速度, 停止从这些spout读取数据。并且受影响的StreamManager 会发送一个特殊的start backpressure message 给其他的StreamManager ,要求他们对spout进行本地降级。 当其他StreamManager 接收到这个特殊消息时,他们通过不读取当地Spout中的Tuple来进行降级。一旦出问题的Heron Instance 恢复速度后,本地的SM 会发送stop backpressure message 解除降级。

很多Socket Channel与应用程序级别的Buffer相关联,该缓冲区由high watermark 和low watermark组成。 当缓冲区大小达到high watermark时触发反压,并保持有效,直到缓冲区大小低于low watermark。 此设计的基本原理是防止拓扑在进入和退出背压缓解模式之间快速振荡。

6 Flink 反压机制

Flink 没有使用任何复杂的机制来解决反压问题,因为根本不需要那样的方案!它利用自身作为纯数据流引擎的优势来优雅地响应反压问题。下面我们会深入分析 Flink 是如何在 Task 之间传输数据的,以及数据流如何实现自然降速的。 Flink 在运行时主要由 operators 和 streams 两大组件构成。每个 operator 会消费中间态的流,并在流上进行转换,然后生成新的流。对于 Flink 的网络机制一种形象的类比是,Flink 使用了高效有界的分布式阻塞队列,就像 Java 通用的阻塞队列(BlockingQueue)一样。还记得经典的线程间通信案例:生产者消费者模型吗?使用 BlockingQueue 的话,一个较慢的接受者会降低发送者的发送速率,因为一旦队列满了(有界队列)发送者会被阻塞。Flink 解决反压的方案就是这种感觉。 在 Flink 中,这些分布式阻塞队列就是这些逻辑流,而队列容量是通过缓冲池来(LocalBufferPool)实现的。每个被生产和被消费的流都会被分配一个缓冲池。缓冲池管理着一组缓冲(Buffer),缓冲在被消费后可以被回收循环利用。这很好理解:你从池子中拿走一个缓冲,填上数据,在数据消费完之后,又把缓冲还给池子,之后你可以再次使用它。

如下图所示展示了 Flink 在网络传输场景下的内存管理。网络上传输的数据会写到 Task 的 InputGate(IG) 中,经过 Task 的处理后,再由 Task 写到 ResultPartition(RS) 中。每个 Task 都包括了输入和输入,输入和输出的数据存在 Buffer 中(都是字节数据)。Buffer 是 MemorySegment 的包装类。

实时流处理系统反压机制07

  1. TaskManager(TM)在启动时,会先初始化NetworkEnvironment对象,TM 中所有与网络相关的东西都由该类来管理(如 Netty 连接),其中就包括NetworkBufferPool。根据配置,Flink 会在 NetworkBufferPool 中生成一定数量(默认2048个)的内存块 MemorySegment(关于 Flink 的内存管理,后续文章会详细谈到),内存块的总数量就代表了网络传输中所有可用的内存。NetworkEnvironment 和 NetworkBufferPool 是 Task 之间共享的,每个 TM 只会实例化一个。
  2. Task 线程启动时,会向 NetworkEnvironment 注册,NetworkEnvironment 会为 Task 的 InputGate(IG)和 ResultPartition(RP) 分别创建一个 LocalBufferPool(缓冲池)并设置可申请的 MemorySegment(内存块)数量。IG 对应的缓冲池初始的内存块数量与 IG 中 InputChannel 数量一致,RP 对应的缓冲池初始的内存块数量与 RP 中的 ResultSubpartition 数量一致。不过,每当创建或销毁缓冲池时,NetworkBufferPool 会计算剩余空闲的内存块数量,并平均分配给已创建的缓冲池。注意,这个过程只是指定了缓冲池所能使用的内存块数量,并没有真正分配内存块,只有当需要时才分配。为什么要动态地为缓冲池扩容呢?因为内存越多,意味着系统可以更轻松地应对瞬时压力(如GC),不会频繁地进入反压状态,所以我们要利用起那部分闲置的内存块。
  3. 在 Task 线程执行过程中,当 Netty 接收端收到数据时,为了将 Netty 中的数据拷贝到 Task 中,InputChannel(实际是 RemoteInputChannel)会向其对应的缓冲池申请内存块(上图中的①)。如果缓冲池中也没有可用的内存块且已申请的数量还没到池子上限,则会向 NetworkBufferPool 申请内存块(上图中的②)并交给 InputChannel 填上数据(上图中的③和④)。如果缓冲池已申请的数量达到上限了呢?或者 NetworkBufferPool 也没有可用内存块了呢?这时候,Task 的 Netty Channel 会暂停读取,上游的发送端会立即响应停止发送,拓扑会进入反压状态。当 Task 线程写数据到 ResultPartition 时,也会向缓冲池请求内存块,如果没有可用内存块时,会阻塞在请求内存块的地方,达到暂停写入的目的。
  4. 当一个内存块被消费完成之后(在输入端是指内存块中的字节被反序列化成对象了,在输出端是指内存块中的字节写入到 Netty Channel 了),会调用 Buffer.recycle() 方法,会将内存块还给 LocalBufferPool (上图中的⑤)。如果LocalBufferPool中当前申请的数量超过了池子容量(由于上文提到的动态容量,由于新注册的 Task 导致该池子容量变小),则LocalBufferPool会将该内存块回收给 NetworkBufferPool(上图中的⑥)。如果没超过池子容量,则会继续留在池子中,减少反复申请的开销。

下面这张图简单展示了两个 Task 之间的数据传输以及 Flink 如何感知到反压的:

实时流处理系统反压机制08

  1. 记录“A”进入了 Flink 并且被 Task 1 处理。(这里省略了 Netty 接收、反序列化等过程)
  2. 记录被序列化到 buffer 中。
  3. 该 buffer 被发送到 Task 2,然后 Task 2 从这个 buffer 中读出记录。

不要忘了:记录能被 Flink 处理的前提是,必须有空闲可用的 Buffer。

结合上面两张图看:Task 1 在输出端有一个相关联的 LocalBufferPool(称缓冲池1),Task 2 在输入端也有一个相关联的 LocalBufferPool(称缓冲池2)。如果缓冲池1中有空闲可用的 buffer 来序列化记录 “A”,我们就序列化并发送该 buffer。

这里我们需要注意两个场景:

  • 本地传输:如果 Task 1 和 Task 2 运行在同一个 worker 节点(TaskManager),该 buffer 可以直接交给下一个 Task。一旦 Task 2 消费了该 buffer,则该 buffer 会被缓冲池1回收。如果 Task 2 的速度比 1 慢,那么 buffer 回收的速度就会赶不上 Task 1 取 buffer 的速度,导致缓冲池1无可用的 buffer,Task 1 等待在可用的 buffer 上。最终形成 Task 1 的降速。
  • 远程传输:如果 Task 1 和 Task 2 运行在不同的 worker 节点上,那么 buffer 会在发送到网络(TCP Channel)后被回收。在接收端,会从 LocalBufferPool 中申请 buffer,然后拷贝网络中的数据到 buffer 中。如果没有可用的 buffer,会停止从 TCP 连接中读取数据。在输出端,通过 Netty 的水位值机制来保证不往网络中写入太多数据(后面会说)。如果网络中的数据(Netty输出缓冲中的字节数)超过了高水位值,我们会等到其降到低水位值以下才继续写入数据。这保证了网络中不会有太多的数据。如果接收端停止消费网络中的数据(由于接收端缓冲池没有可用 buffer),网络中的缓冲数据就会堆积,那么发送端也会暂停发送。另外,这会使得发送端的缓冲池得不到回收,writer 阻塞在向 LocalBufferPool 请求 buffer,阻塞了 writer 往 ResultSubPartition 写数据。

这种固定大小缓冲池就像阻塞队列一样,保证了 Flink 有一套健壮的反压机制,使得 Task 生产数据的速度不会快于消费的速度。我们上面描述的这个方案可以从两个 Task 之间的数据传输自然地扩展到更复杂的 pipeline 中,保证反压机制可以扩散到整个 pipeline。

6.3 反压实验

另外,官方博客中为了展示反压的效果,给出了一个简单的实验。下面这张图显示了:随着时间的改变,生产者(黄色线)和消费者(绿色线)每5秒的平均吞吐与最大吞吐(在单一JVM中每秒达到8百万条记录)的百分比。我们通过衡量task每5秒钟处理的记录数来衡量平均吞吐。该实验运行在单 JVM 中,不过使用了完整的 Flink 功能栈。

实时流处理系统反压机制09

首先,我们运行生产task到它最大生产速度的60%(我们通过Thread.sleep()来模拟降速)。消费者以同样的速度处理数据。然后,我们将消费task的速度降至其最高速度的30%。你就会看到背压问题产生了,正如我们所见,生产者的速度也自然降至其最高速度的30%。接着,停止消费task的人为降速,之后生产者和消费者task都达到了其最大的吞吐。接下来,我们再次将消费者的速度降至30%,pipeline给出了立即响应:生产者的速度也被自动降至30%。最后,我们再次停止限速,两个task也再次恢复100%的速度。总而言之,我们可以看到:生产者和消费者在 pipeline 中的处理都在跟随彼此的吞吐而进行适当的调整,这就是我们希望看到的反压的效果。

在 Storm/JStorm 中,只要监控到队列满了,就可以记录下拓扑进入反压了。但是 Flink 的反压太过于天然了,导致我们无法简单地通过监控队列来监控反压状态。Flink 在这里使用了一个 trick 来实现对反压的监控。如果一个 Task 因为反压而降速了,那么它会卡在向 LocalBufferPool 申请内存块上。那么这时候,该 Task 的 stack trace 就会长下面这样:

1
2
3
4
java.lang.Object.wait(Native Method)
o.a.f.[...].LocalBufferPool.requestBuffer(LocalBufferPool.java:163)
o.a.f.[...].LocalBufferPool.requestBufferBlocking(LocalBufferPool.java:133) <--- BLOCKING request
[...]

那么事情就简单了。通过不断地采样每个 task 的 stack trace 就可以实现反压监控。

实时流处理系统反压机制10

Flink 的实现中,只有当 Web 页面切换到某个 Job 的 Backpressure 页面,才会对这个 Job 触发反压检测,因为反压检测还是挺昂贵的。JobManager 会通过 Akka 给每个 TaskManager 发送TriggerStackTraceSample消息。默认情况下,TaskManager 会触发100次 stack trace 采样,每次间隔 50ms(也就是说一次反压检测至少要等待5秒钟)。并将这 100 次采样的结果返回给 JobManager,由 JobManager 来计算反压比率(反压出现的次数/采样的次数),最终展现在 UI 上。UI 刷新的默认周期是一分钟,目的是不对 TaskManager 造成太大的负担。

7 总结

Flink不需要一种特殊的机制来处理反压,因为Flink 中的数据传输相当于已经提供了应对反压的机制。因此,Flink 所能获得的最大吞吐量由其 pipeline 中最慢的组件决定。相对于 Storm/JStorm 的实现,Flink 的实现更为简洁优雅,源码中也看不见与反压相关的代码,无需 Zookeeper/TopologyMaster 的参与也降低了系统的负载,也利于对反压更迅速的响应。

本文转自 实时流处理系统反压机制(BackPressure)综述
https://blog.csdn.net/qq_21125183/article/details/80708142
开启Back Pressure使生产环境的Spark Streaming应用更稳定、有效

Spark 数据倾斜的成因与治理

一句话定位:数据倾斜是分布式计算里最典型、最高频的性能杀手——少量 key 承载了不成比例的数据量,拖出长尾 Task,让整个 Stage 的吞吐受限于最慢的那一个。
面试里它既是”会不会调优”的试金石,也是把大数据经验迁移到 GPU 推理调优的最佳桥梁(长尾 Task ≈ 静态 Batching 木桶效应)。

1. 成因:为什么会发生倾斜

分布式计算按 key 分区(默认 HashPartitioner:同一 key 必落同一 Task)。当数据本身的 key 分布不均(存在热点 key / 大 key),或分区键 / 哈希函数选得不合适时,少数 Task 会分到远超均值的数据量,成为长尾(straggler),拖慢整个 Stage——因为 Stage 的完成取决于最慢的 Task。

关键结论:倾斜只发生在 Shuffle 过程中。窄依赖(map / filter / mapPartitions)不跨节点重分布、不按 key 汇聚,因此不存在按 key 汇聚的倾斜;只有 Shuffle(按 key 重分区)才会把热点 key 集中到一个 reduce Task 上。

2. 触发 Shuffle 的算子(何时会踩坑)

以下算子会产生 Shuffle,是倾斜的高发区:distinctgroupByKeyreduceByKeyaggregateByKeyjoincogrouprepartition 等。

它们”为什么会倾斜”本质是同一件事——按 key 哈希把数据打进同一个 reduce Task;区别只在于每个算子对 key 的处理方式不同,导致倾斜的烈度不同。

2.1 算子功能与倾斜原理对照

算子 功能 倾斜原因
groupByKey 按 key 分组,把同 key 的 value 收集成 Iterable,不聚合 无 map 端预聚合,热点 key 的全部 value 原样 shuffle 到同一 Task,Shuffle 量最大、最易 OOM
reduceByKey 按 key 用可结合函数聚合(sum / count / max) 自带 map 端 combine,Shuffle 量小;但热点 key 跨分区汇聚到同一 reduce,长尾仍在
aggregateByKey reduceByKey 的通用版:分区内 seqOp + 跨分区 combOp reduceByKey,有 map 端预聚合,热点 key 仍会汇聚
distinct 去重(等价 map(_ -> null).reduceByKey(...) 按值分组去重,热点值(如空串 / 默认值)汇聚到同一 Task;不可加盐,否则破坏去重语义
join 两表按 key 关联配对 热点 key 两侧记录汇聚到同一 Task,且配对结果是 M×N 乘积放大,比聚合倾斜更危险
cogroup 多个 RDD 按 key 联合分组,各取 value 的 Iterable 多个 groupByKey 叠加、无预聚合,多数据源的热点 key 叠加
repartition 重新分区以调整并行度 / 分布 RDD 的 repartition(n) 均匀打散、不按业务 key,本身不引入倾斜;DataFrame 的 repartition(n, col) / repartitionByRange(col) 按列分布,遇热点列会重新引入倾斜

2.2 两条底层规律

  1. map 端 combine 只减”分区内重复”,不减”跨分区热点”reduceByKey / aggregateByKey 把每个 map 分区内的同 key 先合并成一份,Shuffle 数据量下降;但 1000 个 map 各发一份,reduce 端仍要处理这 1000 份——热点 key 的倾斜只是被”缓解”,不会被”消除”。
  2. 聚合是收敛的,join 是发散的:groupBy 类把 M 条收成 1 条(越聚越小);join 则把左表 M 条 × 右表 N 条配对(越配越大)。同一个热点 key,join 的倾斜烈度通常远高于聚合,所以 join 需要单独的治理手段(Broadcast Join / 热键隔离 / AQE skew join,见第 3 节)。

参考:Spark性能优化指南——高级篇

3. 治理手段(四类)

3.1 加盐打散(Salting)

给热点 key 拼接随机前缀/后缀(如 key + "_" + random(N)),把一个热点打散成 N 份落到不同 Task 并行处理。

  • 适用:可拆分的聚合(计数、求和)。
  • Tradeoff:多一次 Shuffle 阶段、逻辑更复杂;不适用于需要精确去重(distinct count)的场景——加盐后同一原始 key 被拆散,distinct 会算错。

3.2 两阶段聚合(局部聚合 + 全局聚合)

也称 Map 端预聚合 / 两遍聚合,专治 groupBy 类倾斜:

  1. 第一阶段(局部):给 key 加随机前缀,按 前缀+key 做局部聚合(reduceByKey / aggregateByKey),先把热点大幅缩小。
  2. 第二阶段(全局):去掉前缀,按原始 key 做全局聚合。
  • 适用:可累加的聚合(sum / count / 可两阶段 avg)。
  • 不适用:不可累加(distinct、median、全局 topN)。
  • 关键配合:优先用 reduceByKey / aggregateByKey(自带 map 端 combine,shuffle 数据量小),避免 groupByKey(无 map 端聚合,shuffle 量最大,最易倾斜)。

3.3 Broadcast Join 消除 Shuffle

小表 join 大表时,把小表广播到所有 Executor,在 map 端完成 join,完全避免大表的 Shuffle——倾斜自然消失。

  • 适用条件:小表能放进 Executor 内存(默认阈值 spark.sql.autoBroadcastJoinThreshold = 10MB,可上调)。
  • Tradeoff:小表过大广播会撑爆内存/网络;大表 join 大表不适用(考虑 Bucket Join / Sort Merge Join)。
  • 若倾斜源于 join 的倾斜 key,且维表(小表)不大,Broadcast Join 是直接”绕开”倾斜的利器。

3.4 倾斜 key 单独处理(隔离热键)

工程上最常用、最可控:先识别热点 key,再把它与非倾斜 key 分开处理

  • 非倾斜 key 走正常 join / 聚合;
  • 倾斜 key 单独用 Broadcast Join,或拆成多个子 key 分别 join 后再 union
  • 例:大表 join 维表但该大表某 key 极热 → 把热 key 从大表过滤出来单独和维表广播 join,其余走正常路径。

3.5 兜底:AQE 自动倾斜处理

Spark 3.0+ 的 AQE(spark.sql.adaptive.enabled=true)默认开启 spark.sql.adaptive.skewJoin.enabled
运行时统计各分区大小,若某分区 > 中位数 N 倍且超过阈值,自动拆分为多个 Task 并行。

  • 主要解决 join 倾斜;对非 join 的 groupBy 倾斜帮助有限(AQE 的 coalesce partitions 仅缓解小分区过多)。
  • 面试要点:AQE 是”白嫖”的兜底,但要知道原理、知道它的边界,不能只依赖它。

4. 面试落点:提前抽样发现,而非事后救火

事后救火(作业跑慢/失败、OOM、长尾 Task 几小时跑不完才去猜)是被动且昂贵的。要做到提前发现

  • 抽样定位:全量跑之前,对 join key / group key 做 sample() + countByValue()(或 groupBy(key).count()),看 key 分布——若某 key 的 count 远高于均值即为倾斜 key。
  • Spark UI 看板:Stage 的 Task 耗时分布、Shuffle Read 大小分布;单个 Task 的 Shuffle Read 特别大、duration 特别长 → 倾斜。
  • 建基线:对作业建立 key 分布基线,新数据接入时校验是否出现新热点。

数字直觉:若某 key 占全量 30%、并行度 100,则该 key 对应的单个 reduce 独吞约 30% 数据,其余 99 个 Task 分摊 70%——这就是典型长尾。

可迁移话术(回指 1.1.1):Spark 长尾 Task ≈ 静态 Batching 木桶效应提前抽样发现倾斜 ≈ 推理前先做输入分布 Profiling / 建立基线,而非 GPU 利用率已呈锯齿才救火。把这套”提前量化基线 → 定位 → 验证”的方法论直接迁移到 GPU 推理调优,是面试里最强的差异化叙事。

1 Spark-Join

参数 默认值 说明
spark.sql.shuffle.partitions 200 Configures the number of partitions to use when shuffling data for joins or aggregations.
spark.default.parallelism For distributed shuffle operations like reduceByKey and join, the largest number of partitions in a parent RDD. For operations like parallelize with no parent RDDs, it depends on the cluster manager:Local mode: number of cores on the local machineMesos fine grained mode: 8 Others: total number of cores on all executor nodes or 2, whichever is larger Default number of partitions in RDDs returned by transformations like join, reduceByKey, and parallelize when not set by user.

上面两个参数都是设置默认的并行度,但是适用的场景不同:

spark.sql.shuffle.partitions是对sparkSQL进行shuffle操作的时候生效,比如 join或者aggregation等操作的时候,之前有个同学设置了spark.default.parallelism 这个并行度为2000,结果还是产生200的stage,排查了很久才发现,是这个原因。
spark.default.parallelism这个参数只是针对rdd的shuffle操作才生效,比如join,reduceByKey。

作者:pcqlegend
链接:https://www.jianshu.com/p/c5914126ef98
来源:简书
著作权归作者所有。商业转载请联系作者获得授权,非商业转载请注明出处。

1.1 shuffle join VS map join

在对大表与大表之间进行连接操作时,通常都会触发 Shuffle Join,两表的所有分区节点会进行 All-to-All 的通讯,这种查询通常比较昂贵,会对网络 IO 会造成比较大的负担。

而对于大表和小表的连接操作,Spark 会在一定程度上进行优化,如果小表的数据量小于 Worker Node 的内存空间,Spark 会考虑将小表的数据广播到每一个 Worker Node,在每个工作节点内部执行连接计算,这可以降低网络的 IO,但会加大每个 Worker Node 的 CPU 负担。

是否采用广播方式进行 Join 取决于程序内部对小表的判断,如果想明确使用广播方式进行 Join,则可以在 DataFrame API 中使用 broadcast 方法指定需要广播的小表:

1
2
empDF.join(broadcast(deptDF), joinExpression).show()
复制代码

作者:heibaiying
链接:https://juejin.im/post/6844903950349500430
来源:掘金
著作权归作者所有。商业转载请联系作者获得授权,非商业转载请注明出处。

1.2 map-join小表误判

1
2
// a是一个几亿行的大表,b是一个只有几十行的小表。a和b都是由hive创建的表
select * from a where id not in (select id from b)

在spark ui中看到了该sql的执行计划,该sql语句执行了Map-side Join操作,但是spark把a表当成了小表,准备把a表broadcast到其他的节点,然后就是一直卡在这步broadcast操作上。 造成上述问题的原因就是spark认为a表是一个小表,但是在spark ui上明显可以看到a表读了很多的行。但是为什么spark还会认为a表是一个小表呢?原因是spark判断一个hive表的大小会用hive的metastore数据来判断,因为我们的a表没有执行过ANALYZE TABLE,自然a表的metastore里面的数据就不准确了。

1.2.1 解决方法

  1. 设置spark.sql.statistics.fallBackToHdfs=True
    该参数能让spark直接读取hdfs的文件大小来判断一个表达大小,从而代替从metastore里面的获取的关于表的信息。这样spark自然能正确的判断出表的大小,从而使用b表来进行broadcast。

  2. 使用hint
    在使用sql语句执行的时候在sql语句里面加上mapjoin的注释,也能够达到相应的效果,比如把上述的sql语句改成:

1
select /*+ BROADCAST (b) */ * from a where id not in (select id from b)

这样spark也会使用b表来进行broadcast。

  1. 使用spark代码的方式
    使用broadcast函数就能达到此效果:
1
2
from pyspark.sql.functions import broadcast
broadcast(spark.table("b")).join(spark.table("a"), "id").show()
  1. 备注
  • 只有当要进行join的表的大小小于spark.sql.autoBroadcastJoinThreshold(默认是10M)的时候,才会进行mapjoin。

  • Impala通过hint和执行表的位置调整也能够优化join操作,通过explain也可以查看sql的执行计划,然后再进行优化。

Join背景

当前SparkSQL支持三种join算法:Shuffle Hash Join、Broadcast Hash Join以及Sort Merge Join。其中前两者归根到底都属于Hash Join,只不过载Hash Join之前需要先Shuffle还是先Broadcast。其实,Hash Join算法来自于传统数据库,而Shuffle和Broadcast是大数据在分布式情况下的概念,两者结合的产物。因此可以说,大数据的根就是传统数据库。Hash Join是内核。

1.2.1.1 Spark Join的分类和实现机制

图片

上图是Spark Join的分类和使用。

1.2.1.1.1 Hash Join

先来看看这样一条SQL语句:select * from order,item where item.id = order.i_id,参与join的两张表是order和item,join key分别是item.id以及order.i_id。现在假设Join采用的是hash join算法,整个过程会经历三步:

  • 确定Build Table以及Probe Table:这个概念比较重要,Build Table会被构建成以join key为key的hash table,而Probe Table使用join key在这张hash table表中寻找符合条件的行,然后进行join链接。Build表和Probe表是Spark决定的。通常情况下,小表会被作为Build Table,较大的表会被作为Probe Table。

  • 构建Hash Table:依次读取Build Table(item)的数据,对于每一条数据根据Join Key(item.id)进行hash,hash到对应的bucket中(类似于HashMap的原理),最后会生成一张HashTable,HashTable会缓存在内存中,如果内存放不下会dump到磁盘中。

  • 匹配:生成Hash Table后,在依次扫描Probe Table(order)的数据,使用相同的hash函数(在spark中,实际上就是要使用相同的partitioner)在Hash Table中寻找hash(join key)相同的值,如果匹配成功就将两者join在一起。

1.2.1.1.2 Broadcast Hash Join

当Join的一张表很小的时候,使用broadcast hash join。

Broadcast Hash Join的条件有以下几个:

  • 被广播的表需要小于spark.sql.autoBroadcastJoinThreshold所配置的信息,默认是10M;

  • 基表不能被广播,比如left outer join时,只能广播右表。

图片

broadcast hash join可以分为两步:

  • broadcast阶段:将小表广播到所有的executor上,广播的算法有很多,最简单的是先发给driver,driver再统一分发给所有的executor,要不就是基于bittorrete的p2p思路;

  • hash join阶段:在每个executor上执行 hash join,小表构建为hash table,大表的分区数据匹配hash table中的数据。

1.2.1.1.3 Sort Merge Join

图片

当两个表都非常大时,SparkSQL采用了一种全新的方案来对表进行Join,即Sort Merge Join。这种方式不用将一侧数据全部加载后再进行hash join,但需要在join前将数据进行排序。

首先将两张表按照join key进行重新shuffle,保证join key值相同的记录会被分在相应的分区,分区后对每个分区内的数据进行排序,排序后再对相应的分区内的记录进行连接。可以看出,无论分区有多大,Sort Merge Join都不用把一侧的数据全部加载到内存中,而是即用即丢;因为两个序列都有有序的,从头遍历,碰到key相同的就输出,如果不同,左边小就继续取左边,反之取右边。从而大大提高了大数据量下sql join的稳定性。

整个过程分为三个步骤:

  • shuffle阶段:将两张大表根据join key进行重新分区,两张表数据会分布到整个集群,以便分布式并行处理

  • sort阶段:对单个分区节点的两表数据,分别进行排序

  • merge阶段:对排好序的两张分区表数据执行join操作。join操作很简单,分别遍历两个有序序列,碰到相同join key就merge输出,否则继续取更小一边的key。

图片

经过上文的分析,很明显可以得出这几种join的代价关系:cost(Broadcast Hash Join)< cost(Shuffle Hash Join) < cost(Sort Merge Join),数据仓库设计时最好避免大表与大表的join查询,SparkSQL也可以根据内存资源、带宽资源适量将参数spark.sql. autoBroadcastJoinThreshold调大,让更多join实际执行为Broadcast Hash Join。

1 Spark-SQL

1.1 join

Spark 中支持多种连接类型:

  • Inner Join : 内连接;
  • Full Outer Join : 全外连接;
  • Left Outer Join : 左外连接;
  • Right Outer Join : 右外连接;
  • Left Semi Join : 左半连接;
  • Left Anti Join : 左反连接;
  • Natural Join : 自然连接;
  • Cross (or Cartesian) Join : 交叉 (或笛卡尔) 连接

SQL JOINS

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
emp 员工表
|-- ENAME: 员工姓名
|-- DEPTNO: 部门编号
|-- EMPNO: 员工编号
|-- HIREDATE: 入职时间
|-- JOB: 职务
|-- MGR: 上级编号
|-- SAL: 薪资
|-- COMM: 奖金

dept 部门表
|-- DEPTNO: 部门编号
|-- DNAME: 部门名称
|-- LOC: 部门所在城市

-- LEFT SEMI JOIN
SELECT * FROM emp LEFT SEMI JOIN dept ON emp.deptno = dept.deptno
-- 等价于如下的 IN 语句
SELECT * FROM emp WHERE deptno IN (SELECT deptno FROM dept)

-- LEFT ANTI JOIN
SELECT * FROM emp LEFT ANTI JOIN dept ON emp.deptno = dept.deptno
-- 等价于如下的 IN 语句
SELECT * FROM emp WHERE deptno NOT IN (SELECT deptno FROM dept)

--CROSS JOIN
SELECT * FROM emp CROSS JOIN dept ON emp.deptno = dept.deptno

--自然连接是在两张表中寻找那些数据类型和列名都相同的字段,然后自动地将他们连接起来,并返回所有符合条件的结果。
SELECT * FROM emp NATURAL JOIN dept

--程序自动推断出使用两张表都存在的 dept 列进行连接
SELECT * FROM emp JOIN dept ON emp.deptno = dept.deptno

1.1.1 内部实现

broadcast join –> hash join –> sort-merge join

在对大表与大表之间进行连接操作时,通常都会触发 Shuffle Join,两表的所有分区节点会进行 All-to-All 的通讯,这种查询通常比较昂贵,会对网络 IO 会造成比较大的负担。

https://github.com/heibaiying

而对于大表和小表的连接操作,Spark 会在一定程度上进行优化,如果小表的数据量小于 Worker Node 的内存空间,Spark 会考虑将小表的数据广播到每一个 Worker Node,在每个工作节点内部执行连接计算,这可以降低网络的 IO,但会加大每个 Worker Node 的 CPU 负担。

是否采用广播方式进行 Join 取决于程序内部对小表的判断,如果想明确使用广播方式进行 Join,则可以在 DataFrame API 中使用 broadcast 方法指定需要广播的小表:

1
empDF.join(broadcast(deptDF), joinExpression).show()

1.2 Driver Collect Exec

数据必须收集到Driver的Exec:

Exec Statement 优化方向
CollectLimitExec 替换为GlobalLimitExec
TakeOrderedAndProjectExec
CollectTailExec

1.3 优化

优化规则 规则名称 简介
列裁剪 column_prune 对于上层算子不需要的列,不在下层算子输出该列,减少计算
子查询去关联 decorrelate 尝试对相关子查询进行改写,将其转换为普通 join 或 aggregation 计算
聚合消除 aggregation_eliminate 尝试消除执行计划中的某些不必要的聚合算子
投影消除 projection_eliminate 消除执行计划中不必要的投影算子
最大最小消除 max_min_eliminate 改写聚合中的 max/min 计算,转化为 order by + limit 1
谓词下推 predicate_push_down 尝试将执行计划中过滤条件下推到离数据源更近的算子上
外连接消除 outer_join_eliminate 尝试消除执行计划中不必要的 left join 或者 right join
分区裁剪 partition_processor 将分区表查询改成为用 union all,并裁剪掉不满足过滤条件的分区
聚合下推 aggregation_push_down 尝试将执行计划中的聚合算子下推到更底层的计算节点
TopN 下推 topn_push_down 尝试将执行计划中的 TopN 算子下推到离数据源更近的算子上
Join 重排序 join_reorder 对多表 join 确定连接顺序

1.4 逻辑优化

1.4.1 子查询相关的优化

关联子查询去关联

1.4.2 列裁剪

1.4.3 关联子查询去关联

1.4.4 Max/Min 消除

1.4.5 谓词下推

1.4.6 分区裁剪

1.4.7 TopN 和 Limit 下推

1.4.8 Join Reorder

1.5 物理优化

1.5.1 选择最优的索引进行表的访问

1.5.2 收集统计信息来获得表的数据分布情况

1.5.3 在错误索引的解决方案中会介绍当发现 TiDB 索引选错时,你应该使用那些手段来让它使用正确的索引

1.5.4 在 Distinct 优化中会介绍在物理优化中会做的一个有关 DISTINCT 关键字的优化,在这一小节中会介绍它的优缺点以及如何使用它。

1.6 [参考文献]

  1. The Business Intelligence for Hadoop Benchmark

1 Spark资源评估

机器机型 内存 硬盘 核数
M10 128G 3.6T 48
BX1 16G×16 256G 4T×12=48T 80
CG3 256G 3.6T 96

1.1 Spark On Yarn 内存计算

在介绍了,spark任务在yarn运行时需要的Continer数量,以及内存大小之后,我们再来看spark on yarn的时候整体任务在yarn中占用资源大小。

Core: yarn中Core指的是Continer数量,所以Core = ContinerNum

而内存的计算则较为复杂了,设单个Continer向集群申请的资源经我们上面公式算出来的需要申请的内存大小为:excutorTotalMemory ,则该Continer在yarn集群上占用的最终资源为continerMemory。
minContiner = yarn.scheduler.minimum-allocation-mb(continer分配资源的最小值,目前是128)
Increment = yarn.scheduler.increment-allocation-mb(yarn分配资源的增量,也叫规整化参数,默认值为1024 mb)
resultMemory的计算方式如下所示:

1
2
3
4
5
If(totalMemory<=minContiner){
continerMemory = minContiner
}else{
continerMemory = minContiner + Math.ceil((excutorTotalMemory - minContiner)/increment) * increment
}

总结

例如某个spark任务的提交参数为,driverMemory=2G,executorMemory=2G,executorNum = 1
minContiner=512m
Increment =1024m
则该任务
executorContinerMemory计算过程如下

申请资源数:executor = Max(executorMemory*0.1,384M)+executorMemory=2432M
ContinerMymory = 512+Math.ceil((2432-512)/1024.0)*1024 = 2.5G

driverContinerMemory计算过程同上:2.5G
最终该任务在yarn消耗资源为5G
可以看出来,spark任务最终消耗资源并非为初始化资源数。

需要join 75张表,每张表的主键分布不同:

  • 直接join会造成数据倾斜,某个节点撑爆
  • 所有的表都shuffle,会造成shuffle数据量太多,撑爆硬盘

申请的资源:

策略一:

  • join后的表,每隔join20次则repartition一次
  • 待join的子表,partition个数超过30,或行数超过1.5亿,则repartition一次

宽表数据量:

1.2 问题点

  1. dag排布的规则是什么?

Spark-task_split_block

梳理一下Spark中关于并发度涉及的几个概念File,Block,Split,Task,Partition,RDD以及节点数、Executor数、core数目的关系。

  1. 用户设置了numSplit,那么goalSize=totalSize/numSplit
  2. minSize=max(1,minSplitSize)
  3. splitSize=max(minSplitSize, min(goalSize,blockSize))
  4. task个数=totalSize除以splitSize

输入可能以多个文件的形式存储在HDFS上,每个File都包含了很多块,称为Block
当Spark读取这些文件作为输入时,会根据具体数据格式对应的InputFormat进行解析,一般是将若干个Block合并成一个输入分片,称为InputSplit,注意InputSplit不能跨越文件。
随后将为这些输入分片生成具体的Task。InputSplit与Task是一一对应的关系。
随后这些具体的Task每个都会被分配到集群上的某个节点的某个Executor去执行。

  • 每个节点可以起一个或多个Executor。
  • 每个Executor由若干core组成,每个Executor的每个core一次只能执行一个Task。
  • 每个Task执行的结果就是生成了目标RDD的一个partiton

注意: 这里的core是虚拟的core而不是机器的物理CPU核,可以理解为就是Executor的一个工作线程。

而 Task被执行的并发度 = Executor数目 * 每个Executor核数。

至于partition的数目:

  • 对于数据读入阶段,例如sc.textFile,输入文件被划分为多少InputSplit就会需要多少初始Task。
  • 在Map阶段partition数目保持不变。
  • 在Reduce阶段,RDD的聚合会触发shuffle操作,聚合后的RDD的partition数目跟具体操作有关,例如repartition操作会聚合成指定分区数,还有一些算子是可配置的。

1,Application

application(应用)其实就是用spark-submit提交的程序。比方说spark examples中的计算pi的SparkPi。一个application通常包含三部分:从数据源(比方说HDFS)取数据形成RDD,通过RDD的transformation和action进行计算,将结果输出到console或者外部存储(比方说collect收集输出到console)。

2,Driver

Spark中的driver感觉其实和yarn中Application Master的功能相类似。主要完成任务的调度以及和executor和cluster manager进行协调。有client和cluster联众模式。client模式driver在任务提交的机器上运行,而cluster模式会随机选择机器中的一台机器启动driver。从spark官网截图的一张图可以大致了解driver的功能。

3,Job

Spark中的Job和MR中Job不一样不一样。MR中Job主要是Map或者Reduce Job。而Spark的Job其实很好区别,一个action算子就算一个Job,比方说count,first等。

4, Task

Task是Spark中最新的执行单元。RDD一般是带有partitions的,每个partition的在一个executor上的执行可以任务是一个Task。

5, Stage

Stage概念是spark中独有的。一般而言一个Job会切换成一定数量的stage。各个stage之间按照顺序执行。至于stage是怎么切分的,首选得知道spark论文中提到的narrow dependency(窄依赖)和wide dependency( 宽依赖)的概念。其实很好区分,看一下父RDD中的数据是否进入不同的子RDD,如果只进入到一个子RDD则是窄依赖,否则就是宽依赖。宽依赖和窄依赖的边界就是stage的划分点