0%

转载自https://zhuanlan.zhihu.com/p/87131964

作者:邱从贤(山智)

在 Flink 中,状态可靠性保证由 Checkpoint 支持,当作业出现 failover 的情况下,Flink 会从最近成功的 Checkpoint 恢复。在实际情况中,我们可能会遇到 Checkpoint 失败,或者 Checkpoint 慢的情况,本文会统一聊一聊 Flink 中 Checkpoint 异常的情况(包括失败和慢),以及可能的原因和排查思路。

1. Checkpoint 流程简介

首先我们需要了解 Flink 中 Checkpoint 的整个流程是怎样的,在了解整个流程之后,我们才能在出问题的时候,更好的进行定位分析。

img

从上图我们可以知道,Flink 的 Checkpoint 包括如下几个部分:

  • JM trigger checkpoint
  • Source 收到 trigger checkpoint 的 PRC,自己开始做 snapshot,并往下游发送 barrier
  • 下游接收 barrier(需要 barrier 都到齐才会开始做 checkpoint)
  • Task 开始同步阶段 snapshot
  • Task 开始异步阶段 snapshot
  • Task snapshot 完成,汇报给 JM

上面的任何一个步骤不成功,整个 checkpoint 都会失败。

2 Checkpoint 异常情况排查

2.1 Checkpoint 失败

可以在 Checkpoint 界面看到如下图所示,下图中 Checkpoint 10423 失败了。

img

点击 Checkpoint 10423 的详情,我们可以看到类系下图所示的表格(下图中将 operator 名字截取掉了)。

img

上图中我们看到三行,表示三个 operator,其中每一列的含义分别如下:

  • 其中 Acknowledged 一列表示有多少个 subtask 对这个 Checkpoint 进行了 ack,从图中我们可以知道第三个 operator 总共有 5 个 subtask,但是只有 4 个进行了 ack;
  • 第二列 Latest Acknowledgement 表示该 operator 的所有 subtask 最后 ack 的时间;
  • End to End Duration 表示整个 operator 的所有 subtask 中完成 snapshot 的最长时间;
  • State Size 表示当前 Checkpoint 的 state 大小 – 主要这里如果是增量 checkpoint 的话,则表示增量大小;
  • Buffered During Alignment 表示在 barrier 对齐阶段积攒了多少数据,如果这个数据过大也间接表示对齐比较慢);

Checkpoint 失败大致分为两种情况:Checkpoint Decline 和 Checkpoint Expire。

2.1.1 Checkpoint Decline

我们能从 jobmanager.log 中看到类似下面的日志
Decline checkpoint 10423 by task 0b60f08bf8984085b59f8d9bc74ce2e1 of job 85d268e6fbc19411185f7e4868a44178. 其中
10423 是 checkpointID,0b60f08bf8984085b59f8d9bc74ce2e1 是 execution id,85d268e6fbc19411185f7e4868a44178 是 job id,我们可以在 jobmanager.log 中查找 execution id,找到被调度到哪个 taskmanager 上,类似如下所示:

1
2
2019-09-02 16:26:20,972 INFO  [jobmanager-future-thread-61] org.apache.flink.runtime.executiongraph.ExecutionGraph        - XXXXXXXXXXX (100/289) (87b751b1fd90e32af55f02bb2f9a9892) switched from SCHEDULED to DEPLOYING.
2019-09-02 16:26:20,972 INFO [jobmanager-future-thread-61] org.apache.flink.runtime.executiongraph.ExecutionGraph - Deploying XXXXXXXXXXX (100/289) (attempt #0) to slot container_e24_1566836790522_8088_04_013155_1 on hostnameABCDE

从上面的日志我们知道该 execution 被调度到 hostnameABCDEcontainer_e24_1566836790522_8088_04_013155_1 slot 上,接下来我们就可以到 container container_e24_1566836790522_8088_04_013155 的 taskmanager.log 中查找 Checkpoint 失败的具体原因了。

另外对于 Checkpoint Decline 的情况,有一种情况我们在这里单独抽取出来进行介绍:Checkpoint Cancel。

当前 Flink 中如果较小的 Checkpoint 还没有对齐的情况下,收到了更大的 Checkpoint,则会把较小的 Checkpoint 给取消掉。我们可以看到类似下面的日志:

1
$taskNameWithSubTaskAndID: Received checkpoint barrier for checkpoint 20 before completing current checkpoint 19. Skipping current checkpoint.

这个日志表示,当前 Checkpoint 19 还在对齐阶段,我们收到了 Checkpoint 20 的 barrier。然后会逐级通知到下游的 task checkpoint 19 被取消了,同时也会通知 JM 当前 Checkpoint 被 decline 掉了。

在下游 task 收到被 cancelBarrier 的时候,会打印类似如下的日志:

1
2
3
4
5
6
7
8
9
10
11
12
DEBUG
$taskNameWithSubTaskAndID: Checkpoint 19 canceled, aborting alignment.

或者

DEBUG
$taskNameWithSubTaskAndID: Checkpoint 19 canceled, skipping alignment.

或者

WARN
$taskNameWithSubTaskAndID: Received cancellation barrier for checkpoint 20 before completing current checkpoint 19. Skipping current checkpoint.

上面三种日志都表示当前 task 接收到上游发送过来的 barrierCancel 消息,从而取消了对应的 Checkpoint。

2.1.2 Checkpoint Expire

如果 Checkpoint 做的非常慢,超过了 timeout 还没有完成,则整个 Checkpoint 也会失败。当一个 Checkpoint 由于超时而失败是,会在 jobmanager.log 中看到如下的日志:

1
Checkpoint 1 of job 85d268e6fbc19411185f7e4868a44178  expired before completing.

表示 Chekpoint 1 由于超时而失败,这个时候可以可以看这个日志后面是否有类似下面的日志:

1
Received late message for now expired checkpoint attempt 1 from 0b60f08bf8984085b59f8d9bc74ce2e1 of job 85d268e6fbc19411185f7e4868a44178.

可以按照 2.1.1 中的方法找到对应的 taskmanager.log 查看具体信息。

下面的日志如果是 DEBUG 的话,我们会在开始处标记 DEBUG

我们按照下面的日志把 TM 端的 snapshot 分为三个阶段,开始做 snapshot 前,同步阶段,异步阶段:

1
2
DEBUG
Starting checkpoint (6751) CHECKPOINT on task taskNameWithSubtasks (4/4)

这个日志表示 TM 端 barrier 对齐后,准备开始做 Checkpoint。

1
2
3
DEBUG
2019-08-06 13:43:02,613 DEBUG org.apache.flink.runtime.state.AbstractSnapshotStrategy - DefaultOperatorStateBackend snapshot (FsCheckpointStorageLocation {fileSystem=org.apache.flink.core.fs.SafetyNetWrapperFileSystem@70442baf, checkpointDirectory=xxxxxxxx, sharedStateDirectory=xxxxxxxx, taskOwnedStateDirectory=xxxxxx, metadataFilePath=xxxxxx, reference=(default), fileStateSizeThreshold=1024}, synchronous part) in thread Thread[Async calls on Source: xxxxxx
_source -> Filter (27/70),5,Flink Task Threads] took 0 ms.

上面的日志表示当前这个 backend 的同步阶段完成,共使用了 0 ms。

1
2
DEBUG
DefaultOperatorStateBackend snapshot (FsCheckpointStorageLocation {fileSystem=org.apache.flink.core.fs.SafetyNetWrapperFileSystem@7908affe, checkpointDirectory=xxxxxx, sharedStateDirectory=xxxxx, taskOwnedStateDirectory=xxxxx, metadataFilePath=xxxxxx, reference=(default), fileStateSizeThreshold=1024}, asynchronous part) in thread Thread[pool-48-thread-14,5,Flink Task Threads] took 369 ms

上面的日志表示异步阶段完成,异步阶段使用了 369 ms

在现有的日志情况下,我们通过上面三个日志,定位 snapshot 是开始晚,同步阶段做的慢,还是异步阶段做的慢。然后再按照情况继续进一步排查问题。

2.2 Checkpoint 慢

在 2.1 节中,我们介绍了 Checkpoint 失败的排查思路,本节会分情况介绍 Checkpoint 慢的情况。

Checkpoint 慢的情况如下:比如 Checkpoint interval 1 分钟,超时 10 分钟,Checkpoint 经常需要做 9 分钟(我们希望 1 分钟左右就能够做完),而且我们预期 state size 不是非常大。

对于 Checkpoint 慢的情况,我们可以按照下面的顺序逐一检查。

2.2.0 Source Trigger Checkpoint 慢

这个一般发生较少,但是也有可能,因为 source 做 snapshot 并往下游发送 barrier 的时候,需要抢锁(这个现在社区正在进行用 mailBox 的方式替代当前抢锁的方式,详情参考[1])。如果一直抢不到锁的话,则可能导致 Checkpoint 一直得不到机会进行。如果在 Source 所在的 taskmanager.log 中找不到开始做 Checkpoint 的 log,则可以考虑是否属于这种情况,可以通过 jstack 进行进一步确认锁的持有情况。

2.2.1 使用增量 Checkpoint

现在 Flink 中 Checkpoint 有两种模式,全量 Checkpoint 和 增量 Checkpoint,其中全量 Checkpoint 会把当前的 state 全部备份一次到持久化存储,而增量 Checkpoint,则只备份上一次 Checkpoint 中不存在的 state,因此增量 Checkpoint 每次上传的内容会相对更好,在速度上会有更大的优势。

现在 Flink 中仅在 RocksDBStateBackend 中支持增量 Checkpoint,如果你已经使用 RocksDBStateBackend,可以通过开启增量 Checkpoint 来加速,具体的可以参考 [2]。

2.2.2 作业存在反压或者数据倾斜

我们知道 task 仅在接受到所有的 barrier 之后才会进行 snapshot,如果作业存在反压,或者有数据倾斜,则会导致全部的 channel 或者某些 channel 的 barrier 发送慢,从而整体影响 Checkpoint 的时间,这两个可以通过如下的页面进行检查:

img

上图中我们选择了一个 task,查看所有 subtask 的反压情况,发现都是 high,表示反压情况严重,这种情况下会导致下游接收 barrier 比较晚。

img

上图中我们选择其中一个 operator,点击所有的 subtask,然后按照 Records Received/Bytes Received/TPS 从大到小进行排序,能看到前面几个 subtask 会比其他的 subtask 要处理的数据多。

如果存在反压或者数据倾斜的情况,我们需要首先解决反压或者数据倾斜问题之后,再查看 Checkpoint 的时间是否符合预期。

2.2.2 Barrier 对齐慢

从前面我们知道 Checkpoint 在 task 端分为 barrier 对齐(收齐所有上游发送过来的 barrier),然后开始同步阶段,再做异步阶段。如果 barrier 一直对不齐的话,就不会开始做 snapshot。

barrier 对齐之后会有如下日志打印:

1
2
DEBUG
Starting checkpoint (6751) CHECKPOINT on task taskNameWithSubtasks (4/4)

如果 taskmanager.log 中没有这个日志,则表示 barrier 一直没有对齐,接下来我们需要了解哪些上游的 barrier 没有发送下来,如果你使用 At Least Once 的话,可以观察下面的日志:

1
2
DEBUG
Received barrier for checkpoint 96508 from channel 5

表示该 task 收到了 channel 5 来的 barrier,然后看对应 Checkpoint,再查看还剩哪些上游的 barrier 没有接受到,对于 ExactlyOnce 暂时没有类似的日志,可以考虑自己添加,或者 jmap 查看。

2.2.3 主线程太忙,导致没机会做 snapshot

在 task 端,所有的处理都是单线程的,数据处理和 barrier 处理都由主线程处理,如果主线程在处理太慢(比如使用 RocksDBBackend,state 操作慢导致整体处理慢),导致 barrier 处理的慢,也会影响整体 Checkpoint 的进度,在这一步我们需要能够查看某个 PID 对应 hotmethod,这里推荐两个方法:

  1. 多次连续 jstack,查看一直处于 RUNNABLE 状态的线程有哪些;
  2. 使用工具 AsyncProfile dump 一份火焰图,查看占用 CPU 最多的栈;

如果有其他更方便的方法当然更好,也欢迎推荐。

2.2.4 同步阶段做的慢

同步阶段一般不会太慢,但是如果我们通过日志发现同步阶段比较慢的话,对于非 RocksDBBackend 我们可以考虑查看是否开启了异步 snapshot,如果开启了异步 snapshot 还是慢,需要看整个 JVM 在干嘛,也可以使用前一节中的工具。对于 RocksDBBackend 来说,我们可以用 iostate 查看磁盘的压力如何,另外可以查看 tm 端 RocksDB 的 log 的日志如何,查看其中 SNAPSHOT 的时间总共开销多少。

RocksDB 开始 snapshot 的日志如下:

1
2019/09/10-14:22:55.734684 7fef66ffd700 [utilities/checkpoint/checkpoint_impl.cc:83] Started the snapshot process -- creating snapshot in directory /tmp/flink-io-87c360ce-0b98-48f4-9629-2cf0528d5d53/XXXXXXXXXXX/chk-92729

snapshot 结束的日志如下:

1
2019/09/10-14:22:56.001275 7fef66ffd700 [utilities/checkpoint/checkpoint_impl.cc:145] Snapshot DONE. All is good

2.2.6 异步阶段做的慢

对于异步阶段来说,tm 端主要将 state 备份到持久化存储上,对于非 RocksDBBackend 来说,主要瓶颈来自于网络,这个阶段可以考虑观察网络的 metric,或者对应机器上能够观察到网络流量的情况(比如 iftop)。

对于 RocksDB 来说,则需要从本地读取文件,写入到远程的持久化存储上,所以不仅需要考虑网络的瓶颈,还需要考虑本地磁盘的性能。另外对于 RocksDBBackend 来说,如果觉得网络流量不是瓶颈,但是上传比较慢的话,还可以尝试考虑开启多线程上传功能[3]。

3 总结

在第二部分内容中,我们介绍了官方编译的包的情况下排查一些 Checkpoint 异常情况的主要场景,以及相应的排查方法,如果排查了上面所有的情况,还是没有发现瓶颈所在,则可以考虑添加更详细的日志,逐步将范围缩小,然后最终定位原因。

上文提到的一些 DEBUG 日志,如果 flink dist 包是自己编译的话,则建议将 Checkpoint 整个步骤内的一些 DEBUG 改为 INFO,能够通过日志了解整个 Checkpoint 的整体阶段,什么时候完成了什么阶段,也在 Checkpoint 异常的时候,快速知道每个阶段都消耗了多少时间。

参考内容

[1] Change threading-model in StreamTask to a mailbox-based approach
[2] 增量 checkpoint 原理介绍
[3] RocksDBStateBackend 多线程上传 State

Flink cep

CEP的处理范例引起了人们的极大兴趣,并在各种用例中得到了应用。 最值得注意的是,CEP现在用于诸如股票市场趋势和信用卡欺诈检测等金融应用

模式,从流中查找符合某个pattern的个体事件。
可以将一个pattern sequence视为pattern组成的图, 基于用户定义的条件,从一个pattern传递到下一个pattern
一个match是事件必须流过复杂pattern图的所有的pattern。

注意

  • 每一个pattern必须具有唯一的名称,用于标示符合条件的事件
  • pattern名称不能包含:

Flink join

Flink DataStream API为用户提供了3个算子来实现双流join,分别是:

  • join(): inner join,on window
  • coGroup(): custom join, on window
  • intervalJoin(): inner join, on time range, keyed stream

另外,还提供了broadcast join来关联较小的

准备数据

从Kafka分别接入点击流和订单流,并转化为POJO。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
DataStream<String> clickSourceStream = env
.addSource(new FlinkKafkaConsumer011<>(
"ods_analytics_access_log",
new SimpleStringSchema(),
kafkaProps
).setStartFromLatest());
DataStream<String> orderSourceStream = env
.addSource(new FlinkKafkaConsumer011<>(
"ods_ms_order_done",
new SimpleStringSchema(),
kafkaProps
).setStartFromLatest());

DataStream<AnalyticsAccessLogRecord> clickRecordStream = clickSourceStream
.map(message -> JSON.parseObject(message, AnalyticsAccessLogRecord.class));
DataStream<OrderDoneLogRecord> orderRecordStream = orderSourceStream
.map(message -> JSON.parseObject(message, OrderDoneLogRecord.class));

join()

join()算子提供的语义为”Window join“,即按照指定字段和(滚动/滑动/会话)窗口进行inner join,支持处理时间和事件时间两种时间特征。

以下示例以10秒滚动窗口,将两个流通过商品ID关联,取得订单流中的售价相关字段。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
clickRecordStream
.join(orderRecordStream)
.where(record -> record.getMerchandiseId())
.equalTo(record -> record.getMerchandiseId())
.window(TumblingProcessingTimeWindows.of(Time.seconds(10)))
.apply(new JoinFunction<AnalyticsAccessLogRecord, OrderDoneLogRecord, String>() {
@Override
public String join(AnalyticsAccessLogRecord accessRecord, OrderDoneLogRecord orderRecord) throws Exception {
return StringUtils.join(Arrays.asList(
accessRecord.getMerchandiseId(),
orderRecord.getPrice(),
orderRecord.getCouponMoney(),
orderRecord.getRebateAmount()
), '\t');
}
})
.print().setParallelism(1);

简单易用。

coGroup()

只有inner join肯定还不够,如何实现left/right outer join呢?答案就是利用coGroup()算子。它的调用方式类似于join()算子,也需要开窗,但是CoGroupFunction比JoinFunction更加灵活,可以按照用户指定的逻辑匹配左流和/或右流的数据并输出。

以下的例子就实现了点击流left join订单流的功能,是很朴素的nested loop join思想(二重循环)。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
clickRecordStream
.coGroup(orderRecordStream)
.where(record -> record.getMerchandiseId())
.equalTo(record -> record.getMerchandiseId())
.window(TumblingProcessingTimeWindows.of(Time.seconds(10)))
.apply(new CoGroupFunction<AnalyticsAccessLogRecord, OrderDoneLogRecord, Tuple2<String, Long>>() {
@Override
public void coGroup(Iterable<AnalyticsAccessLogRecord> accessRecords, Iterable<OrderDoneLogRecord> orderRecords, Collector<Tuple2<String, Long>> collector) throws Exception {
for (AnalyticsAccessLogRecord accessRecord : accessRecords) {
boolean isMatched = false;
for (OrderDoneLogRecord orderRecord : orderRecords) {
// 右流中有对应的记录
collector.collect(new Tuple2<>(accessRecord.getMerchandiseName(), orderRecord.getPrice()));
isMatched = true;
}
if (!isMatched) {
// 右流中没有对应的记录
collector.collect(new Tuple2<>(accessRecord.getMerchandiseName(), null));
}
}
}
})
.print().setParallelism(1);

CoGroupFunction中会返回所有数据,不管有没有匹配上

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
DataStream<Tuple3<Long, String, String>> input1 = ...;
input1 = input1.assignTimestampsAndWatermarks(new AscendingTimestampExtractor<Tuple3<Long, String, String>>() {

@Override
public long extractAscendingTimestamp(Tuple3<Long, String, String> arg0) {
return arg0.f0;
}

});

DataStream<Tuple2<Long, String>> input2 = ...;
input2 = input2.assignTimestampsAndWatermarks(new AscendingTimestampExtractor<Tuple2<Long, String>>() {
@Override
public long extractAscendingTimestamp(Tuple2<Long, String> stringStringTuple2) {
return stringStringTuple2.f0;
}
});

input1.coGroup(input2).where(new KeySelector<Tuple3<Long, String, String>, String>() {
@Override
public String getKey(Tuple3<Long, String, String> itemEntity) throws Exception {
return itemEntity.f1;
}
})
.equalTo(new KeySelector<Tuple2<Long, String>, String>() {
@Override
public String getKey(Tuple2<Long, String> value) throws Exception {
return value.f1;
}
})
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.apply(new CoGroupFunction<Tuple3<Long, String, String>, Tuple2<Long, String>, String>() {
@Override
public void coGroup(Iterable<Tuple3<Long, String, String>> first,
Iterable<Tuple2<Long, String>> second, Collector<String> collector) throws Exception {
StringBuilder buffer = new StringBuilder();
buffer.append("DataStream first:\n");
for (Tuple3<Long, String, String> value : first) {
buffer.append(value).append("\n");
}
buffer.append("DataStream second:\n");
for (Tuple2<Long, String> value : second) {
buffer.append(value.f0).append("=>").append(value.f1).append("\n");
}
collector.collect(buffer.toString());
}
})
.print();

上面的例子,左流有三个元素 Tuple3<String,String,String>,右流有两个元素Tuple2<String,String>
两个流第一个元素相互关联。分别指定两个流的事件时间字段。
两个流关联后,按照EventTime划分窗口。与单流类似。
不管元素是否可以关联上,都会输出

用户可以定义CoGroupFunction函数, 可以实现在窗口内,任意组合,如笛卡尔积

intervalJoin()

join()和coGroup()都是基于窗口做关联的。但是在某些情况下,两条流的数据步调未必一致。例如,订单流的数据有可能在点击流的购买动作发生之后很久才被写入,如果用窗口来圈定,很容易join不上。所以Flink又提供了”Interval join“的语义,按照指定字段以及右流相对左流偏移的时间区间进行关联,即:

$right.timestamp ∈ [left.timestamp + lowerBound; left.timestamp + upperBound]$

interval join也是inner join,虽然不需要开窗,但是需要用户指定偏移区间的上下界,并且只支持事件时间

示例代码如下。注意在运行之前,需要分别在两个流上应用assignTimestampsAndWatermarks()方法获取事件时间戳和水印。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
clickRecordStream
.keyBy(record -> record.getMerchandiseId())
.intervalJoin(orderRecordStream.keyBy(record -> record.getMerchandiseId()))
.between(Time.seconds(-30), Time.seconds(30))
.process(new ProcessJoinFunction<AnalyticsAccessLogRecord, OrderDoneLogRecord, String>() {
@Override
public void processElement(AnalyticsAccessLogRecord accessRecord, OrderDoneLogRecord orderRecord, Context context, Collector<String> collector) throws Exception {
collector.collect(StringUtils.join(Arrays.asList(
accessRecord.getMerchandiseId(),
orderRecord.getPrice(),
orderRecord.getCouponMoney(),
orderRecord.getRebateAmount()
), '\t'));
}
})
.print().setParallelism(1);

由上可见,interval join与window join不同,是两个KeyedStream之上的操作,并且需要调用between()方法指定偏移区间的上下界。如果想令上下界是开区间,可以调用upperBoundExclusive()/lowerBoundExclusive()方法。

interval join的实现原理

以下是KeyedStream.process(ProcessJoinFunction)方法调用的重载方法的逻辑。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
public <OUT> SingleOutputStreamOperator<OUT> process(
ProcessJoinFunction<IN1, IN2, OUT> processJoinFunction,
TypeInformation<OUT> outputType) {
Preconditions.checkNotNull(processJoinFunction);
Preconditions.checkNotNull(outputType);
final ProcessJoinFunction<IN1, IN2, OUT> cleanedUdf = left.getExecutionEnvironment().clean(processJoinFunction);
final IntervalJoinOperator<KEY, IN1, IN2, OUT> operator =
new IntervalJoinOperator<>(
lowerBound,
upperBound,
lowerBoundInclusive,
upperBoundInclusive,
left.getType().createSerializer(left.getExecutionConfig()),
right.getType().createSerializer(right.getExecutionConfig()),
cleanedUdf
);
return left
.connect(right)
.keyBy(keySelector1, keySelector2)
.transform("Interval Join", outputType, operator);
}

可见是先对两条流执行connect()和keyBy()操作,然后利用IntervalJoinOperator算子进行转换。在IntervalJoinOperator中,会利用两个MapState分别缓存左流和右流的数据。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
private transient MapState<Long, List<BufferEntry<T1>>> leftBuffer;
private transient MapState<Long, List<BufferEntry<T2>>> rightBuffer;

@Override
public void initializeState(StateInitializationContext context) throws Exception {
super.initializeState(context);
this.leftBuffer = context.getKeyedStateStore().getMapState(new MapStateDescriptor<>(
LEFT_BUFFER,
LongSerializer.INSTANCE,
new ListSerializer<>(new BufferEntrySerializer<>(leftTypeSerializer))
));
this.rightBuffer = context.getKeyedStateStore().getMapState(new MapStateDescriptor<>(
RIGHT_BUFFER,
LongSerializer.INSTANCE,
new ListSerializer<>(new BufferEntrySerializer<>(rightTypeSerializer))
));
}

其中Long表示事件时间戳,List<BufferEntry<T>>表示该时刻到来的数据记录。

当左流和右流有数据到达时,会分别调用processElement1()processElement2()方法,它们都调用了processElement()方法,代码如下。

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
@Override
public void processElement1(StreamRecord<T1> record) throws Exception {
processElement(record, leftBuffer, rightBuffer, lowerBound, upperBound, true);
}

@Override
public void processElement2(StreamRecord<T2> record) throws Exception {
processElement(record, rightBuffer, leftBuffer, -upperBound, -lowerBound, false);
}

@SuppressWarnings("unchecked")
private <THIS, OTHER> void processElement(
final StreamRecord<THIS> record,
final MapState<Long, List<IntervalJoinOperator.BufferEntry<THIS>>> ourBuffer,
final MapState<Long, List<IntervalJoinOperator.BufferEntry<OTHER>>> otherBuffer,
final long relativeLowerBound,
final long relativeUpperBound,
final boolean isLeft) throws Exception {
final THIS ourValue = record.getValue();
final long ourTimestamp = record.getTimestamp();
if (ourTimestamp == Long.MIN_VALUE) {
throw new FlinkException("Long.MIN_VALUE timestamp: Elements used in " +
"interval stream joins need to have timestamps meaningful timestamps.");
}
if (isLate(ourTimestamp)) {
return;
}
addToBuffer(ourBuffer, ourValue, ourTimestamp);
for (Map.Entry<Long, List<BufferEntry<OTHER>>> bucket: otherBuffer.entries()) {
final long timestamp = bucket.getKey();
if (timestamp < ourTimestamp + relativeLowerBound ||
timestamp > ourTimestamp + relativeUpperBound) {
continue;
}
for (BufferEntry<OTHER> entry: bucket.getValue()) {
if (isLeft) {
collect((T1) ourValue, (T2) entry.element, ourTimestamp, timestamp);
} else {
collect((T1) entry.element, (T2) ourValue, timestamp, ourTimestamp);
}
}
}
long cleanupTime = (relativeUpperBound > 0L) ? ourTimestamp + relativeUpperBound : ourTimestamp;
if (isLeft) {
internalTimerService.registerEventTimeTimer(CLEANUP_NAMESPACE_LEFT, cleanupTime);
} else {
internalTimerService.registerEventTimeTimer(CLEANUP_NAMESPACE_RIGHT, cleanupTime);
}
}

这段代码的思路是:

  1. 取得当前流StreamRecord的时间戳,调用isLate()方法判断它是否是迟到数据(即时间戳小于当前水印值),如是则丢弃。
  2. 调用addToBuffer()方法,将时间戳和数据一起插入当前流对应的MapState
  3. 遍历另外一个流的MapState,如果数据满足前述的时间区间条件,则调用collect()方法将该条数据投递给用户定义的ProcessJoinFunction进行处理。
    collect()方法的代码如下,注意结果对应的时间戳是左右流时间戳里较大的那个。
1
2
3
4
5
6
private void collect(T1 left, T2 right, long leftTimestamp, long rightTimestamp) throws Exception {
final long resultTimestamp = Math.max(leftTimestamp, rightTimestamp);
collector.setAbsoluteTimestamp(resultTimestamp);
context.updateTimestamps(leftTimestamp, rightTimestamp, resultTimestamp);
userFunction.processElement(left, right, context, collector);
}
  1. 调用TimerService.registerEventTimeTimer()注册时间戳为timestamp + relativeUpperBound的定时器,该定时器负责在水印超过区间的上界时执行状态的清理逻辑,防止数据堆积。注意左右流的定时器所属的namespace是不同的,具体逻辑则位于onEventTime()方法中。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
@Override
public void onEventTime(InternalTimer<K, String> timer) throws Exception {
long timerTimestamp = timer.getTimestamp();
String namespace = timer.getNamespace();
logger.trace("onEventTime @ {}", timerTimestamp);
switch (namespace) {
case CLEANUP_NAMESPACE_LEFT: {
long timestamp = (upperBound <= 0L) ? timerTimestamp : timerTimestamp - upperBound;
logger.trace("Removing from left buffer @ {}", timestamp);
leftBuffer.remove(timestamp);
break;
}
case CLEANUP_NAMESPACE_RIGHT: {
long timestamp = (lowerBound <= 0L) ? timerTimestamp + lowerBound : timerTimestamp;
logger.trace("Removing from right buffer @ {}", timestamp);
rightBuffer.remove(timestamp);
break;
}
default:
throw new RuntimeException("Invalid namespace " + namespace);
}
}

broadcast join

interval join

Flink Connector之插件机制

Connector是Flink与外部存储连接的桥梁,Flink在对接外部存储和数据格式上

对于一个外部存储系统的库表来说,需要处理三方面的定义:

  1. connector: 定义了流、批、读、写的实现
  2. Format: 定义了数据的解析
  3. Schema: 定义了表的字段信息

Flink使用Java SPI机制实现插件,方便扩展和管理。Flink更专注于计算,

而各种不同的Connector专注于实现批量读写、流式读写、数据一致性协议、分区、offset和时间戳生成等逻辑

本文主要介绍Flink Connector插件机制

Java SPI

(Service provider interface,服务提供者接口),是Java提供的一套实现扩展的API。

Java SPI是符合接口隔离原则的设计,实现调用者和实现者的解耦,不同的实现方实现了可插拔:

  1. 接口定义了交互的”协议”
  2. 根据不同的实现方案,定义该接口的不同实现类
  3. 调用者根据实际的使用需要,启用、扩展或者替换框架的实现策略

常见的SPI的例子:

  • 数据库驱动加载: JDBC加载不同类型的数据库的驱动
  • 日志门面实现类的加载: Slf4J加载不同提供商的日志实现类

使用介绍:

  1. 服务提供者提供接口的实现类,实现类必须含义不带参数的改造方法
  2. 在Jar包的META-INF/services目录下创建一个以“接口全限定名”为命名的文件,内容为实现类的全限定名;
  3. 将接口的实现类所在的Jar包放在classpath
  4. java.util.ServiceLoder通过扫描META-INF/services目录下的配置文件找到实现类,加载类到JVM

ServiceLoader

  1. ServiceLoader可以跨越jar包获取META-INF下的配置文件
  2. 通过反射方法Class.forName()加载类对象,并用instance()方法将类实例化
  3. 把实例化后的类缓存到providers对象中,(LinkedHashMap<String,S>类型)
    然后返回实例对象

性能损失

  1. 延迟加载,需要遍历全部和加载全部实现类
  2. 非线程安全

TableFactory的定义

flink-table-common模块中定义了TableFactory的接口、调用者和工具类,

TableFactory

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
package org.apache.flink.table.factories;
/**
用于从基于字符串的属性中创建与表相关的不同实例的一个工厂。 该工厂使用Java的服务提供商接口(SPI)做服务发现,使用一组描述所需配置的规范化属性来调用工厂。 工厂允许匹配给定的属性集。
将实现此接口的类添加到要找到的当前类路径中的JAR文件的"META_INF/services/org.apache.flink.table.factories.TableFactory"文件中
*/
public interface TableFactory {

/**
指定为此工厂实现的Context。 该框架保证仅在满足指定的属性和值集的情况下才与此工厂匹配。
典型的属性可能是:
-connector.type
-format.type
指定的属性版本允许框架在字符串格式更改的情况下提供向后兼容的属性:
-connector.property
-version-format.property-version
空上下文意味着工厂匹配所有请求。
*/
Map<String, String> requiredContext();

/**
返回该工厂可以处理的属性键列表。 此方法将用于验证。 如果传递了该工厂无法处理的属性,则将引发异常。 该列表不得包含上下文指定的键。
示例属性可能是:
- schema.#.type
- schema.#.name
- connector.topic
- format.line-delimiter
- format.ignore-parse-errors
- format.fields.#.type
- format.fields.#.name

注意:使用“#”表示值的数组,其中“#”表示一个或多个数字。 诸如“ format.property-version”之类的属性版本不得成为受支持属性的一部分。
在某些情况下,将通配符声明为“ *”可能会很有用。 通配符只能在属性键的末尾声明。
例如,如果应支持任意格式:
- format.*
注意:应谨慎使用通配符,因为它们会吞下不受支持的属性,从而可能导致不良行为
*/
List<String> supportedProperties();
}

以flink-hbase为例, 在META-INF/services/中定义了文件org.apache.flink.table.factories.TableFactory

文件的内容为:

1
org.apache.flink.addons.hbase.HBaseTableFactory

这个类就是TableFactory的实现

TableFactory查找

Flink 在 Connector的实现上,也是采用SPI机制

graph LR
A[SPI查找实现类]-->B[是否TableFactory的子类]
B-->C[是否满足必要属性]
C-->D[是否满足属性匹配]
D-->E[创建TableFactory实现类]

入口方法: org.apache.flink.table.factories.TableFactoryService#find(java.lang.Class<T>, java.util.Map<java.lang.String,java.lang.String>)

调用样例:

1
TableFactoryService.find(classOf[TableFactory[_]],sinkProperties)

最终调用查找方法

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
	private static <T extends TableFactory> T findSingleInternal(
Class<T> factoryClass,
Map<String, String> properties,
Optional<ClassLoader> classLoader) {
// 查找Factory
List<TableFactory> tableFactories = discoverFactories(classLoader);
// 过滤Factory
List<T> filtered = filter(tableFactories, factoryClass, properties);

if (filtered.size() > 1) {
// 查找到超过1个,抛出异常
throw new AmbiguousTableFactoryException(
filtered,
factoryClass,
tableFactories,
properties);
} else {
return filtered.get(0);
}
}

查找实现类

1
2
3
4
5
6
7
8
9
10
11
12
13
14
private static List<TableFactory> discoverFactories(Optional<ClassLoader> classLoader) {
try {
List<TableFactory> result = new LinkedList<>();
ClassLoader cl = classLoader.orElse(Thread.currentThread().getContextClassLoader());
ServiceLoader
.load(TableFactory.class, cl)
.iterator()
.forEachRemaining(result::add);
return result;
} catch (ServiceConfigurationError e) {
LOG.error("Could not load service provider for table factories.", e);
throw new TableException("Could not load service provider for table factories.", e);
}
}

过滤

有三层过滤:

  • filterByFactoryClass: 判断是否为TableFactory的实现类
  • filterByContext: 判断必要属性是否匹配: 来源于TableFactory.requiredContext
  • filterBySupportedProperties: 判断属性是否支持:来源于TableFactory.supportedProperties

TableFactory的分类

OPPO's Real-time Data Warehouse and Development of Flink SQL | Medium

对于一个外部存储系统的库表来说,需要处理三方面的定义:

  1. connector: 定义了流、批、读、写的实现
  2. Format: 定义了数据的解析
  3. Schema: 定义了表的字段信息

Schema信息可以在注册库表的时候直接增加,connector、format信息通过TableFactory实现可插拔

classDiagram
      TableFactory <|-- StreamTableSourceFactory~T~
      TableFactory <|-- StreamTableSinkFactory~T~
      TableFactory <|-- BatchTableSourceFactory~T~
      TableFactory <|-- BatchTableSinkFactory~T~
      TableFactory <|-- TableFormatFactory~T~
	  StreamTableSourceFactory <|-- HBaseTableFactory
	  StreamTableSinkFactory <|-- HBaseTableFactory
      class TableFactory{
         +requiredConext() Map~StringString~
         +supportedProperties() List~String~
      }
      class StreamTableSourceFactory{
      	 +createStreamTableSink(Map) StreamTableSink~T~
      	 +createTableSink(Map) TableSink~T~
      }
	  class StreamTableSinkFactory{
		  +createStreamTableSink(Map) StreamTableSink~T~
		  +createTableSink(Map) TableSink~T~

	  }
      class BatchTableSourceFactory{
      	+createBatchTableSource(Map) BatchTableSource~T~
      	+createTableSource(Map) TableSource~T~
      }
      class BatchTableSinkFactory{
      	+createBatchTableSink(Map) BatchTableSink~T~
      	+createTableSink(Map) TableSink~T~
      }
	  class TableFormatFactory{
		  +supportsSchemaDerivation()
		  +supportedProperties() List~String~
	  }
	  class HBaseTableFactory{
	  }

TableFactory可以分为:

  • StreamTableSourceFactory
  • StreamTableSinkFactory
  • BatchTableSourceFactory
  • BatchTableSinkFactory
  • TableFormatFactory

常见的Connector样例

援引一张1.10的优化设计图

img

kafka-DDL定义

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
CREATE TABLE orders_kafka (
...
etime timestamp(3),
watermark for etime as etime - interval '5' second
) WITH (
'connector.type' = 'kafka', -- connector信息
'connector.version' = 'universal',
'connector.topic' = 'kafka-2-mysql-window',
'connector.properties.zookeeper.connect' = 'localhost:2181',
'connector.properties.bootstrap.servers' = 'localhost:9092',
'connector.properties.group.id' = 'order-streaming-test',
'connector.startup-mode' = 'earliest-offset',
'update-mode' = 'append',
'format.type' = 'json', -- format信息
'format.derive-schema' = 'true'
)

MySQL-DDL定义

1
2
3
4
5
6
7
8
9
10
11
CREATE TABLE agg_result (
...
time_id STRING
) WITH (
'connector.type' = 'jdbc',
'connector.driver' = 'com.mysql.cj.jdbc.Driver',
'connector.url' = 'jdbc:mysql://localhost:3306/fdata',
'connector.table' = 't_agg_result',
'connector.username' = 'root',
'connector.password' = '---'
)

Connector设计要点

  1. 自定义Factory,根据需要实现StreamTableSourceFactory和StreamTableSinkFactory
  2. 根据需要继承ConnectorDescriptorValidator,定义自己的connector参数(with 后面跟的那些)
  3. Factory中的requiredContext、supportedProperties都比较重要,框架中对Factory的过滤和检查需要他们
  4. 需要自定义个TableSink,根据你需要连接的中间件选择是AppendStreamTableSink、Upsert、Retract重写consumeDataStream方法
  5. 自定义一个SinkFunction,在invoke方法中实现将数据写入到外部中间件。

SourceFunction

SourceFunction是定义Flink Source的根接口,其源码如下。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
@Public   
public interface SourceFunction<T> extends Function, Serializable {
void run(SourceContext<T> ctx) throws Exception;
void cancel();
interface SourceContext<T> {
void collect(T element);
@PublicEvolving
void collectWithTimestamp(T element, long timestamp);
void emitWatermark(Watermark mark);
@PublicEvolving
void markAsTemporarilyIdle();
Object getCheckpointLock();
void close();
}
}

SourceFunction接口定义了run()方法,该方法用于源源不断地产生源数据,因此重写的时候一般都写成循环,用标志位控制是否结束。cancel()方法则用来打断run()方法中的循环,终止产生数据的过程。

SourceFunction中还嵌套定义了SourceContext接口,它表示这个Source对应的上下文,用来发射数据。其中起主要作用的是前三个方法:

  • collect():发射一个不带自定义时间戳的元素。如果流程序的时间特征(TimeCharacteristic)是处理时间(ProcessingTime),元素没有时间戳;如果是摄入时间(IngestionTime),元素会附带系统时间;如果是事件时间(EventTime),那么初始没有时间戳,但一旦要做与时间戳相关的操作(如窗口)时,就必须用TimestampAssigner设定一个。

  • collectWithTimestamp():发射一个带有自定义时间戳的元素。该方法对于时间特征为事件时间的程序是绝对必须的,如果为处理时间就会被直接忽略,如果为摄入时间就会被系统时间覆盖。

  • emitWatermark():发射一个水印,仅对于事件时间有效。一个带有时间戳t的水印表示不会有任何t’ <= t的事件再发生,如果发生,会被当做迟到事件忽略掉。

SourceFunction还有一些其他实现,如:

  • ParallelSourceFunction,表示该Source可以按照设置的并行度并发执行。

  • RichSourceFunction,继承自富函数RichFunction,表示该Source可以感知到运行时上下文(RuntimeContext,如Task、State、并行度的信息),以及可以自定义初始化和销毁逻辑(通过open()/close()方法)。

  • RichParallelSourceFunction,以上两者的综合。

SinkFunction

SinkFunction是自定义Sink的根接口,其源码如下。

1
2
3
4
5
6
7
8
9
10
11
12
13
public interface SinkFunction<IN> extends Function, Serializable {       
@Deprecated
default void invoke(IN value) throws Exception {}
default void invoke(IN value, Context context) throws Exception {
invoke(value);
}
@Public
interface Context<T> {
long currentProcessingTime();
long currentWatermark();
Long timestamp();
}
}

它的定义比SourceFunction要简单,只有一个invoke()方法,对收集来的每条数据都会调用它来处理。SinkFunction也有对应的上下文对象Context,可以从中获得当前处理时间、当前水印和时间戳。它也有衍生出来的富函数版本RichSinkFunction。

Flink内部提供了一个最简单的实现DiscardingSink。顾名思义,就是将所有汇集的数据全部丢弃。

1
2
3
4
5
6
7
8
9
@Public   
public class DiscardingSink<T> implements SinkFunction<T> {
private static final long serialVersionUID = 1L;

@Override
public void invoke(T value) {

}
}

参考文献

  1. Flink深入浅出:JDBC Connector源码分析
  2. Flink Source/Sink探究与实践:RocketMQ数据写入HBase

Flink savepoint恢复

运行日志分析

正常的恢复

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
[flink-akka.actor.default-dispatcher-2] org.apache.flink.runtime.jobmaster.JobMaster                 :236 - Initializing job t_org_trade_day_dal_statistics_job (000000000000eace000000000000001c).
[flink-akka.actor.default-dispatcher-2] org.apache.flink.runtime.jobmaster.JobMaster :164 - Using restart strategy FixedDelayRestartStrategy(maxNumberRestartAttempts=10000, delayBetweenRestartAttempts=20000) for t_org_trade_day_dal_statistics_job (000000000000eace000000000000001c).
[flink-akka.actor.default-dispatcher-4] org.apache.flink.yarn.YarnResourceManager :238 - Recovered 0 containers from previous attempts ([]).
[flink-akka.actor.default-dispatcher-4] org.apache.hadoop.yarn.client.api.impl.ContainerManagementProtocolProxy :81 - yarn.client.max-cached-nodemanagers-proxies : 0
[flink-akka.actor.default-dispatcher-2] org.apache.flink.runtime.executiongraph.ExecutionGraph :528 - Job recovers via failover strategy: full graph restart
[flink-akka.actor.default-dispatcher-2] org.apache.flink.runtime.jobmaster.JobMaster :203 - Running initialization on master for job t_org_trade_day_dal_statistics_job (000000000000eace000000000000001c).
[flink-akka.actor.default-dispatcher-2] org.apache.flink.runtime.jobmaster.JobMaster :221 - Successfully ran initialization on master in 0 ms.
[flink-akka.actor.default-dispatcher-4] org.apache.flink.yarn.YarnResourceManager :1029 - ResourceManager akka.tcp://flink@9.44.10.157:35473/user/resourcemanager was granted leadership with fencing token 92d59c139e046177dbed78af07b744e8
[main-EventThread] org.apache.flink.runtime.leaderservice.zookeeper.ZooKeeperLeaderElectionService :262 - org.apache.flink.yarn.YarnResourceManager@83ab9bf has been granted leadership.
[flink-akka.actor.default-dispatcher-4] org.apache.flink.runtime.resourcemanager.slotmanager.SlotManagerImpl :215 - Starting the SlotManager.
[flink-akka.actor.default-dispatcher-2] org.apache.flink.runtime.util.ZooKeeperUtils :301 - Initialized ZooKeeperCompletedCheckpointStore in '/checkpoints/000000000000eace000000000000001c'.
[flink-akka.actor.default-dispatcher-2] org.apache.flink.runtime.jobmaster.JobMaster :200 - Using application-defined state backend: RocksDBStateBackend{checkpointStreamBackend=File State Backend (checkpoints: 'hdfs://qy-flink-1-v3/user/u_teg_tdbank/oceanus2/lj_cft_flink/fit-oceanus-prod/project-10001/job-60110/snapshots', savepoints: 'null', asynchronous: UNDEFINED, fileStateThreshold: -1), localRocksDbDirectories=null, enableIncrementalCheckpointing=UNDEFINED, numberOfTransferingThreads=-1}
[flink-akka.actor.default-dispatcher-2] org.apache.flink.runtime.jobmaster.JobMaster :207 - Configuring application-defined state backend with job/cluster config
[flink-akka.actor.default-dispatcher-2] org.apache.flink.contrib.streaming.state.RocksDBStateBackend :348 - Using predefined options: DEFAULT.
[flink-akka.actor.default-dispatcher-2] org.apache.flink.contrib.streaming.state.RocksDBStateBackend :563 - Using default options factory: DefaultConfigurableOptionsFactory{configuredOptions={}}.
[flink-akka.actor.default-dispatcher-2] org.apache.flink.runtime.checkpoint.ZooKeeperCompletedCheckpointStore :128 - Recovering checkpoints from ZooKeeper.
[flink-akka.actor.default-dispatcher-2] org.apache.flink.runtime.checkpoint.ZooKeeperCompletedCheckpointStore :146 - Found 0 checkpoints in ZooKeeper.
[flink-akka.actor.default-dispatcher-2] org.apache.flink.runtime.checkpoint.ZooKeeperCompletedCheckpointStore :163 - Trying to fetch 0 checkpoints from storage.
[flink-akka.actor.default-dispatcher-2] org.apache.flink.runtime.checkpoint.CheckpointCoordinator :1124 - Starting job 000000000000eace000000000000001c from savepoint hdfs://qy-flink-1-v3/user/u_teg_tdbank/oceanus2/lj_cft_flink/fit-oceanus-prod/project-10001/job-60110/snapshots/savepoint-000000-a5c820c4c53a (allowing non restored state)
[flink-akka.actor.default-dispatcher-2] org.apache.flink.runtime.checkpoint.CheckpointCoordinator :1139 - Reset the checkpoint ID of job 000000000000eace000000000000001c to 370.
[flink-akka.actor.default-dispatcher-2] org.apache.flink.runtime.checkpoint.ZooKeeperCompletedCheckpointStore :128 - Recovering checkpoints from ZooKeeper.
[flink-akka.actor.default-dispatcher-2] org.apache.flink.runtime.checkpoint.ZooKeeperCompletedCheckpointStore :146 - Found 1 checkpoints in ZooKeeper.
[flink-akka.actor.default-dispatcher-2] org.apache.flink.runtime.checkpoint.ZooKeeperCompletedCheckpointStore :163 - Trying to fetch 1 checkpoints from storage.
[flink-akka.actor.default-dispatcher-2] org.apache.flink.runtime.checkpoint.ZooKeeperCompletedCheckpointStore :336 - Trying to retrieve checkpoint 369.
[flink-akka.actor.default-dispatcher-2] org.apache.flink.runtime.checkpoint.CheckpointCoordinator :1079 - Restoring job 000000000000eace000000000000001c from latest valid checkpoint: Checkpoint 369 @ 0 for 000000000000eace000000000000001c.
[flink-akka.actor.default-dispatcher-2] org.apache.flink.runtime.checkpoint.CheckpointCoordinator :234 - No master state to restore
[main-EventThread] org.apache.flink.runtime.jobmaster.JobManagerRunner :298 - JobManager runner for job t_org_trade_day_dal_statistics_job (000000000000eace000000000000001c) was granted leadership with session id 3beb8d55-f7c2-4175-9c7e-ac256f03e8bb at akka.tcp://flink@9.44.10.157:35473/user/jobmanager_0.
[main-EventThread] org.apache.flink.runtime.leaderservice.zookeeper.ZooKeeperLeaderElectionService :262 - org.apache.flink.runtime.jobmaster.JobManagerRunner@548d10c7 has been granted leadership.

JobGraph修改后恢复失败

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
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.jobmaster.JobMaster                 :236 - Initializing job flake-sql-job-runner (000000000000eafb0000000000000011).
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.jobmaster.JobMaster :164 - Using restart strategy FixedDelayRestartStrategy(maxNumberRestartAttempts=10000, delayBetweenRestartAttempts=20000) for flake-sql-job-runner (000000000000eafb0000000000000011).
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.executiongraph.ExecutionGraph :528 - Job recovers via failover strategy: full graph restart
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.jobmaster.JobMaster :203 - Running initialization on master for job flake-sql-job-runner (000000000000eafb0000000000000011).
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.jobmaster.JobMaster :221 - Successfully ran initialization on master in 0 ms.
[flink-akka.actor.default-dispatcher-4] org.apache.flink.yarn.YarnResourceManager :238 - Recovered 0 containers from previous attempts ([]).
[flink-akka.actor.default-dispatcher-4] org.apache.hadoop.yarn.client.api.impl.ContainerManagementProtocolProxy :81 - yarn.client.max-cached-nodemanagers-proxies : 0
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.util.ZooKeeperUtils :301 - Initialized ZooKeeperCompletedCheckpointStore in '/checkpoints/000000000000eafb0000000000000011'.
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.jobmaster.JobMaster :200 - Using application-defined state backend:
RocksDBStateBackend{
checkpointStreamBackend=
File State Backend (
checkpoints: 'hdfs://qy-flink-1-v3/user/u_teg_tdbank/oceanus2/lj_cft_flink/fit-oceanus-prod/project-10001/job-60155/snapshots',
savepoints: 'null',
asynchronous: UNDEFINED,
fileStateThreshold: -1
),
localRocksDbDirectories=null,
enableIncrementalCheckpointing=UNDEFINED,
numberOfTransferingThreads=-1
}
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.jobmaster.JobMaster :207 - Configuring application-defined state backend with job/cluster config
[flink-akka.actor.default-dispatcher-3] org.apache.flink.contrib.streaming.state.RocksDBStateBackend :348 - Using predefined options: DEFAULT.
[flink-akka.actor.default-dispatcher-3] org.apache.flink.contrib.streaming.state.RocksDBStateBackend :563 - Using default options factory: DefaultConfigurableOptionsFactory{configuredOptions={}}.
[flink-akka.actor.default-dispatcher-4] org.apache.flink.yarn.YarnResourceManager :1029 - ResourceManager akka.tcp://flink@9.44.34.108:40622/user/resourcemanager was granted leadership with fencing token 9e26d17ff1120b75114052d13be242e0
[main-EventThread] org.apache.flink.runtime.leaderservice.zookeeper.ZooKeeperLeaderElectionService :262 - org.apache.flink.yarn.YarnResourceManager@1a74bd90 has been granted leadership.
[flink-akka.actor.default-dispatcher-4] org.apache.flink.runtime.resourcemanager.slotmanager.SlotManagerImpl :215 - Starting the SlotManager.
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.checkpoint.ZooKeeperCompletedCheckpointStore :128 - Recovering checkpoints from ZooKeeper.
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.checkpoint.ZooKeeperCompletedCheckpointStore :146 - Found 0 checkpoints in ZooKeeper.
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.checkpoint.ZooKeeperCompletedCheckpointStore :163 - Trying to fetch 0 checkpoints from storage.
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.checkpoint.CheckpointCoordinator :1124 - Starting job 000000000000eafb0000000000000011 from savepoint hdfs://qy-flink-1-v3/user/u_teg_tdbank/oceanus2/lj_cft_flink/fit-oceanus-prod/project-10001/job-60155/snapshots/savepoint-000000-630ca6fb8b9c (allowing non restored state)
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.checkpoint.Checkpoints :172 - Could not find ExecutionJobVertex. Including user-defined OperatorIDs in search.
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.checkpoint.Checkpoints :194 - Skipping savepoint state for operator 749b2cad71aaee5be25025871420d4b2.
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.checkpoint.Checkpoints :194 - Skipping savepoint state for operator f7ad6dac6ff3d8ada68fde148d943256.
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.checkpoint.Checkpoints :194 - Skipping savepoint state for operator 5c72ff74f7b814cc2b15db0311cc0aaf.
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.checkpoint.Checkpoints :194 - Skipping savepoint state for operator d6883e647f8d8dd07dbe048d4037223c.
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.checkpoint.Checkpoints :194 - Skipping savepoint state for operator bb6ff7b6d8a06dd49e1764b59667cf62.
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.checkpoint.CheckpointCoordinator :1139 - Reset the checkpoint ID of job 000000000000eafb0000000000000011 to 5.
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.checkpoint.ZooKeeperCompletedCheckpointStore :128 - Recovering checkpoints from ZooKeeper.
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.checkpoint.ZooKeeperCompletedCheckpointStore :146 - Found 1 checkpoints in ZooKeeper.
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.checkpoint.ZooKeeperCompletedCheckpointStore :163 - Trying to fetch 1 checkpoints from storage.
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.checkpoint.ZooKeeperCompletedCheckpointStore :336 - Trying to retrieve checkpoint 4.
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.checkpoint.CheckpointCoordinator :1079 - Restoring job 000000000000eafb0000000000000011 from latest valid checkpoint: Checkpoint 4 @ 0 for 000000000000eafb0000000000000011.
[flink-akka.actor.default-dispatcher-3] org.apache.flink.runtime.checkpoint.CheckpointCoordinator :234 - No master state to restore
[main-EventThread] org.apache.flink.runtime.jobmaster.JobManagerRunner :298 - JobManager runner for job flake-sql-job-runner (000000000000eafb0000000000000011) was granted leadership with session id 9ecd11e9-590a-4186-8270-edff06b7a810 at akka.tcp://flink@9.44.34.108:40622/user/jobmanager_0.
[main-EventThread] org.apache.flink.runtime.leaderservice.zookeeper.ZooKeeperLeaderElectionService :262 - org.apache.flink.runtime.jobmaster.JobManagerRunner@3da1e67a has been granted leadership.
[flink-akka.actor.default-dispatcher-2] org.apache.flink.runtime.jobmaster.JobMaster :691 - Starting execution of job flake-sql-job-runner (000000000000eafb0000000000000011) under job master id 8270edff06b7a8109ecd11e9590a4186.

1 watermark

a window operator registers a timer for every active window, which cleans up the window’s state when the event time passes the window’s ending time.

1.1 处理空闲数据源

在某些情况下,由于数据产生的比较少,导致一段时间内没有数据产生,进而就没有水印的生成,导致下游依赖水印的一些操作就会出现问题,
比如某一个算子的上游有多个算子,这种情况下,水印是取其上游两个算子的较小值,如果上游某一个算子因为缺少数据迟迟没有生成水印,就会出现eventtime倾斜问题,导致下游没法触发计算。

所以filnk通过WatermarkStrategy.withIdleness()方法允许用户在配置的时间内(即超时时间内)没有记录到达时将一个流标记为空闲。这样就意味着下游的数据不需要等待水印的到来。

当下次有水印生成并发射到下游的时候,这个数据流重新变成活跃状态。

通过下面的代码来实现对于空闲数据流的处理

1
2
3
WatermarkStrategy
.<Tuple2<Long, String>>forBoundedOutOfOrderness(Duration.ofSeconds(20))
.withIdleness(Duration.ofMinutes(1));

在 Flink SQL 使用 Kafka 作为 Source 的场景中,即使已配置 40 秒乱序 Watermark,仍可能出现 Watermark 滞后数分钟甚至 10 分钟以上 的现象;这在 多并行度单并行度 情况下本质原因一致,均不属于乱序参数失效。Watermark 的推进规则是:算子 Watermark 由其输入 Channel 的最小值决定,且只能在消费到新数据时基于最新事件时间推进,不会随系统时间自动前进。多并行度时,低频或空闲的 Kafka partition 会拖慢整体 Watermark(需通过 scan.idle-timeout 处理 idle 分区);单并行度时,则更应重点关注 Kafka 是否存在历史或批量数据、当前是否持续有新数据到达、event_time 字段本身是否整体晚于当前时间。此外,任务在 failover / savepoint 恢复 后会继承历史 Watermark 状态,短时间内表现为明显滞后;下游算子或 Sink 出现 反压、阻塞、慢 IO 或重 UDF 计算,也可能导致 Source 无法持续 poll 数据,从而间接延缓 Watermark 发送。需要明确的是,乱序时间只用于容忍事件顺序的乱序,并不等价于限制 Watermark 最大延迟;当 Watermark 长时间落后时,应系统性地从 Kafka 分区与入流模式、事件时间字段的真实性与时效性、任务恢复历史以及上下游执行链路是否阻塞等方面进行排查并沉淀为工程经验。

向量化执行引擎

向量化执行引擎

解决的问题

  1. 传统的一次一个tuple的pipeline模式,CPU的大部分处理时在遍历操作数,而不是在处理数据,CPU利用率低,还导致缓存性能低和频繁跳转。
  2. 寻址性能是存储的瓶颈,顺序读写性能比随机读取效率更高,但大部分查询都是随机读取,此外磁盘读写速度远远落后于CPU数据执行速度。列存储可以最大化利用磁盘读写能力

列存储的优势:

  • 压缩能力提升: 列类型统一,存放一起,易于压缩
  • 减少IO读写总量:仅仅需要读取需要的列
  • 减少查询过程中的节点函数的调用次数:计算过程中,列存以数据块的形式返回上层节点,减少函数调用次数
  • 向量化执行:计算的过程中,相同的列执行相同的操作,可以使用SIMD提升计算效率;不支持SIMD,也可以通过循环提升效率
  • 延迟物化:减少查询计划树之间传递数据总量

向量化引擎的适用条件

  • 列存储
  • OLAP

OLAP适合列存储
OLTP点查询适合行存储

实现方法

  • 火山模型: 修改成一次返回一组列
  • 层次型执行模式: 将优化好的执行计划数转换为编译执行:一次调用下来后,每一层都完成后才向上返回数据,减少各层次节点间的调用次数

在数据量比较大的情况下,内存可能放不下这些数据,需要写盘,这样会造成额外的开销。

优势

  • 向量化执行引擎可以减少节点间的调度,提高CPU的利用率。
  • 因为列存数据,同一列的数据放在一起,导致向量化执行引擎在执行的时候拥有了更多的机会能够利用的当前硬件与编译的新优化特征。
  • 因为列存数据存储将同类型的类似数据放在一起使得压缩比能够达到更高,这样可以拉近一些磁盘IO能力与计算能力的差距。

注意的问题

  • 通信库的效率: MPP架构上是share nothing架构的,所以它的集群各执行节点是有通信需要的,通信效率的高低也是决定了查询执行效率。另外就是大集群情况下,如果使用tcp方式连接,连接数会受限。
  • 数据读写争抢问题 这个问题本身不是向量化执行引擎的,而是列存带来的,因为列存储表每一列单独存储为一个文件,这样在写盘的时候有优化与没有优化的差距还是非常明显的。
  • 列存数据过滤效率问题 列存数据中的一个处理单元是由连续的N个值放在一起组成的一个Col(数组),然后再由多个Col的数组组成了一个处理单元。在进行过率的时候如何能够更加紧凑的放置数据是需要我们考虑列存在过滤掉效率和存放之间如何优化的问题。
  • 表达式计算问题(LLVM) LLVM优化可以将表达式计算由遍历树多层调用模式变为,只调用一个函数的扁平式执行方式。这样可以极大的提高表达式的执行性能。值得一提的是LLVM技术的优势也可以应用在执行计划编译执行模型的构建上面。

火山模型

例如 SQL:

1
2
3
SELECT Id, Name, Age, (Age - 30) * 50 AS Bonus
FROM People
WHERE Age > 30

对应火山模型如下:

其中——

User:客户端;

Project:垂直分割(投影),选择字段;

Select(或 Filter):水平分割(选择),用于过滤行,也称为谓词;

Scan:扫描数据。

这里包含了 3 个 Operator,首先 User 调用最上方的 Operator(Project)希望得到 next tuple,Project 调用子节点(Select),而 Select 又调用子节点(Scan),Scan 获得表中的 tuple 返回给 Select,Select 会检查是否满足过滤条件,如果满足则返回给 Project,如果不满足则请求 Scan 获取 next tuple。Project 会对每一个 tuple 选择需要的字段或者计算新字段并返回新的 tuple 给 User。当 Scan 发现没有数据可以获取时,则返回一个结束标记告诉上游已结束。

为了更好地理解一个 Operator 中发生了什么,下面通过伪代码来理解 Select Operator:

1
2
3
4
5
6
7
8
9
Tuple Select::next() {
while (true) {
Tuple candidate = child->next(); // 从子节点中获取 next tuple
if (candidate == EndOfStream) // 是否得到结束标记
return EndOfStream;
if (condition->check(candidate)) // 是否满足过滤条件
return candidate; // 返回 tuple
}
}

火山模型的优缺点

可以看出火山模型的优点在于:简单,每个 Operator 可以单独抽象实现、不需要关心其他 Operator 的逻辑。

那么缺点呢?也够明显吧?每次都是计算一个 tuple(Tuple-at-a-time),这样会造成多次调用 next ,也就是造成大量的虚函数调用,这样会造成 CPU 的利用率不高。

知识点补习——虚函数
C++ 中用 virtual 标记的函数,而在 Java 中没有 final 修饰的普通方法(没有标记为 static、native)都是虚函数。
虚函数的重要特性是支持在子类中进行 override(重写),从而实现面向对象的重要特性之一:多态。

但是为什么之前的数据库设计者没有去优化这方面呢?是他们没想到吗?怎么可能?这个时候我们可能要考虑到 30 年前的硬件水平了,当时的 IO 速度是远远小于 CPU 的计算速度的,那么 SQL 查询引擎的优化则会被 IO 开销所遮蔽(毕竟花费很多精力只带来 1% 场景下的速度提升意义并不大)。

可是随着近些年来存储越来越快,这个时候我们再思考如何让计算更快可能就有点意思了。

虚函数造成cpu利用率不高

  • 空间开销
    首先,由于需要为每一个包含虚函数的类生成一个虚函数表,所以程序的二进制文件大小会相应的增大;其次,对于包含虚函数的类的实例来说,每个实例都包含一个虚函数表指针用于指向对应的虚函数表,所以每个实例的空间占用都增加一个指针大小(32位系统4字节,64位系统8字节)。这些空间开销可能会造成缓存的不友好,在一定程度上影响程序性能。
  • 时间开销
    虚函数的时间开销主要是增加了一次内存寻址,通过虚函数表指针找到虚函数表,虽对程序性能有一些影响,但是影响并不大。

影响到虚函数调用性能的背后原因是流水线和分支预测,由于虚函数调用需要间接跳转,所以会导致虚函数调用比普通函数调用多了分支预测的过程,产生性能差距的原因主要是分支预测失败导致的流水线冲刷性能开销。

优化方向

代码生成

向量化

按照行加载到CPU Cache,如果只访问个别列就丢弃了,CPU Cache利用率不高。
按列存储,因为输入是同列的一组数据,面对的是相同的操作,并且如果每次只取一列的部分数据,返回一个可以放到 CPU Cache 的向量,那么又可以利用到 CPU Cache。

知识点补习:向量化
向量化计算就是将一个循环处理一个数组的时候每次处理 1 个数据共处理 N 次,转化为向量化——每次同时处理 8 个数据共处理 N/8 次,其中依赖的技术就是 SIMD(Single Instruction Multiple Data,单指令流多数据流),SIMD 可以在一条 CPU 指令上处理 2、4、8 或者更多份的数据。

代码生成是否因为SQL复杂过慢

循环优化

share nothing架构

表达式计算问题(LLVM)

列式存储

SIMD

火山模型

延迟物化

谓词

Greenplum

upsert: on duplicate key update

MySQL upsert 语法

不能加 where 条件

因为这是个插入语句,所以不能加 where 条件。

影响的行数

  • 如果是插入操作,受到影响行的值为 1;
  • 如果更新操作,受到影响行的值为 2;
  • 如果更新的数据和已有的数据一样(就相当于没变,所有值保持不变),受到影响的行的值为 0。

样例一

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
53
54
55
56
57
58
59
60
61
> desc t1;
|-------|------|------|-----|---------|----------------|
| Field | Type | Null | Key | Default | Extra |
|-------|------|------|-----|---------|----------------|
| f1 | int | NO | PRI | <null> | auto_increment |
| f2 | int | YES | | <null> | |
| f3 | int | YES | | <null> | |
|-------|------|------|-----|---------|----------------|

> select * from t1;
|----|--------|----|
| f1 | f2 | f3 |
|----|--------|----|

> insert into t1(f1,f3) values (1,3),(2,7) on duplicate key update f3=f3|1;
Query OK, 2 rows affected

> select * from t1;
|----|--------|----|
| f1 | f2 | f3 |
|----|--------|----|
| 1 | <null> | 3 |
| 2 | <null> | 7 |
|----|--------|----|
2 rows in set

> insert into t1(f1,f3) values (1,3),(2,7) on duplicate key update f3=f3|1;
Query OK, 4 rows affected

> select * from t1;
|----|--------|----|
| f1 | f2 | f3 |
|----|--------|----|
| 1 | <null> | 4 |
| 2 | <null> | 8 |
|----|--------|----|
2 rows in set


> insert into t1(f1,f3) values (1,3),(1,7) on duplicate key update f3=f3|1;
-- 同一行冲突两次,会执行两次更新语句
> select * from t1;
|----|--------|----|
| f1 | f2 | f3 |
|----|--------|----|
| 1 | <null> | 6 |
| 2 | <null> | 8 |
|----|--------|----|
2 rows in set


> select * from t1;
|----|--------|----|
| f1 | f2 | f3 |
|----|--------|----|
| 1 | <null> | 6 |
| 2 | <null> | 8 |
|----|--------|----|
2 rows in set


VALUES 引用

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
> truncate t1;
> INSERT INTO t1 (f1,f2,f3) VALUES (1,2,3),(4,5,6) ON DUPLICATE KEY UPDATE f3=VALUES(f1)|VALUES(f2);
Query OK, 2 rows affected
Time: 0.003s

> select * from t1;
-- 没有出现冲突,直接赋值
|----|----|----|
| f1 | f2 | f3 |
|----|----|----|
| 1 | 2 | 3 |
| 4 | 5 | 6 |
|----|----|----|
2 rows in set
Time: 0.008s

-- 再执行一次,冲突了.执行冲突后的处理语句

> INSERT INTO t1 (f1,f2,f3) VALUES (1,2,3),(4,5,6) ON DUPLICATE KEY UPDATE f3=VALUES(f1)|VALUES(f2);
Query OK, 2 rows affected
Time: 0.002s

> select * from t1;
|----|----|----|
| f1 | f2 | f3 |
|----|----|----|
| 1 | 2 | 3 |
| 4 | 5 | 9 |
|----|----|----|
2 rows in set
Time: 0.008s

ALAIS 引用

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
> truncate table t1;
> INSERT INTO t1 (f1,f2,f3) VALUES (1,2,3),(4,5,6) AS new ON DUPLICATE KEY UPDATE f3 = new.f1|new.f2;
-- 空表没有冲突,直接添加了
> select * from t1;
|----|----|----|
| f1 | f2 | f3 |
|----|----|----|
| 1 | 2 | 3 |
| 4 | 5 | 6 |
|----|----|----|
2 rows in set
Time: 0.008s
-- 再来
> INSERT INTO t1 (f1,f2,f3) VALUES (1,2,3),(4,5,6) AS new ON DUPLICATE KEY UPDATE f3 = new.f1|new.f2;
-- 出现冲突了, 两行更新
> select * from t1;
|----|----|----|
| f1 | f2 | f3 |
|----|----|----|
| 1 | 2 | 3 |
| 4 | 5 | 9 |
|----|----|----|
2 rows in set

等价写法

1
2
3
4
5
6
7
INSERT INTO t1 (f1,f2,f3) VALUES (1,2,3),(4,5,6) AS new(m,n,p)
ON DUPLICATE KEY UPDATE c = m|n;
INSERT INTO t1 SET f1=1,f2=2,f3=3 AS new
ON DUPLICATE KEY UPDATE c = new.a|new.b;

INSERT INTO t1 SET f1=1,f2=2,f3=3 AS new(m,n,p)
ON DUPLICATE KEY UPDATE c = m|n;

UNION 引用

1
2
3
4
5
6
INSERT INTO t1 (f1, f2)
SELECT * FROM
(SELECT c, d FROM t2
UNION
SELECT e, f FROM t3) AS dt
ON DUPLICATE KEY UPDATE f2 = f2 | f3;

key 自增情况

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
> truncate table t1;
> select * from t1;
|----|----|----|
| f1 | f2 | f3 |
|----|----|----|

> insert into t1(f2,f3) values (1,3),(1,7) on duplicate key update f3=f2|f3;
Query OK, 2 rows affected
Time: 0.002s

> insert into t1(f2,f3) values (1,3),(1,7) on duplicate key update f3=f2|f3;
Query OK, 2 rows affected
Time: 0.002s

-- 不存在主键冲突,连续可以插入两次

> select * from t1;
|----|----|----|
| f1 | f2 | f3 |
|----|----|----|
| 1 | 1 | 3 |
| 2 | 1 | 7 |
| 3 | 1 | 3 |
| 4 | 1 | 7 |
|----|----|----|
4 rows in set

同一个表有多个unique index会出现冲突

on duplicate key 可以设置多余字段

1
2
3
4
5
6
7
8
CREATE TABLE `duplication_key` (
`f1` int NOT NULL,
`f2` varchar(10) DEFAULT NULL,
`f3` bigint DEFAULT NULL,
`f4` decimal(10, 0) DEFAULT NULL,
PRIMARY KEY (`f1`),
UNIQUE KEY `f2_f3` (`f2`, `f3`)
);
f1 f2 f3 f4
1 1 1 1
2 2 2 2

验证primary key冲突情况

1
2
3
4
5
6
insert into duplication_key
values('2','3','3','20') on duplicate key update
`f1`=VALUES(`f1`),
`f2`=VALUES(`f2`),
`f3`=VALUES(`f3`),
`f4`=VALUES(`f4`);
f1 f2 f3 f4
1 1 1 1
2 3 3 20

f1是主键,出现冲突,更新数据

验证唯一index冲突情况

1
2
3
4
5
6
insert into duplication_key
values('3','3','3','30') on duplicate key update
`f1`=VALUES(`f1`),
`f2`=VALUES(`f2`),
`f3`=VALUES(`f3`),
`f4`=VALUES(`f4`);
f1 f2 f3 f4
1 1 1 1
3 3 3 30

只要有一个冲突就会出现更新

如果同一张表使用不同的unique key来更新数据,那么会出现错误更新问题

慢查询优化

1
mysqldumpslow -s t -t 10  ~/Downloads/慢日志/mysql-slow.log
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
Count: 2  Time=1238.48s (2476s)  Lock=0.00s (0s)  Rows=90750.0 (181500), cdw[cdw]@2hosts
select t3.PRODUCT_CHANNEL_CODE,t3.MAIN_PROJECT_NAME ,
t3.MAIN_PROJECT_CODE ,t3.SUB_PROJECT_NAME,t3.SUB_PROJECT_CODE,
t3.CHANNEL_TYPE,t3.CHANNEL_NAME,t3.CHANNEL_ID,t3.CHANNEL_SUB_CODE,
t3.CHANNEL_SUB_TYPE,
t4.FRIST_CHANNEL,
t4.SECOND_CHANNEL,
t4.thrid_channel
from (SELECT t1.PRODUCT_CHANNEL_CODE,t1.MAIN_PROJECT_NAME ,
t1.MAIN_PROJECT_CODE ,t1.SUB_PROJECT_NAME,t1.SUB_PROJECT_CODE,
t1.CHANNEL_TYPE,t1.CHANNEL_NAME,t1.CHANNEL_ID,t2.CHANNEL_SUB_CODE,
t2.CHANNEL_SUB_TYPE FROM qd_dz_db.qd_prochannel t1 left join qd_channel t2
on t1.CHANNEL_ID=t2.CHANNEL_ID ) t3 join
(
SELECT
a.PRODUCT_CHANNEL_CODE,
b.FRIST_CHANNEL,
b.SECOND_CHANNEL,
b.THIRD_CHANNEL AS thrid_channel
FROM
qd_dz_db.qd_prochannel a
JOIN
qd_dz_db.qd_channel b ON a.CHANNEL_ID = b.CHANNEL_ID
)t4 on t3.PRODUCT_CHANNEL_CODE=t4.PRODUCT_CHANNEL_CODE
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
explain
select
t3.PRODUCT_CHANNEL_CODE,
t3.MAIN_PROJECT_NAME ,
t3.MAIN_PROJECT_CODE ,
t3.SUB_PROJECT_NAME,
t3.SUB_PROJECT_CODE,
t3.CHANNEL_TYPE,
t3.CHANNEL_NAME,
t3.CHANNEL_ID,
t3.CHANNEL_SUB_CODE,
t3.CHANNEL_SUB_TYPE,
t4.FRIST_CHANNEL,
t4.SECOND_CHANNEL,
t4.thrid_channel
from (
SELECT
t1.PRODUCT_CHANNEL_CODE,
t1.MAIN_PROJECT_NAME ,
t1.MAIN_PROJECT_CODE ,
t1.SUB_PROJECT_NAME,
t1.SUB_PROJECT_CODE,
t1.CHANNEL_TYPE,
t1.CHANNEL_NAME,
t1.CHANNEL_ID,
t2.CHANNEL_SUB_CODE,
t2.CHANNEL_SUB_TYPE
FROM
qd_dz_db.qd_prochannel t1
left join
qd_channel t2
on t1.CHANNEL_ID=t2.CHANNEL_ID
) t3
join
(
SELECT
a.PRODUCT_CHANNEL_CODE,
b.FRIST_CHANNEL,
b.SECOND_CHANNEL,
b.THIRD_CHANNEL AS thrid_channel
FROM
qd_dz_db.qd_prochannel a
JOIN
qd_dz_db.qd_channel b
ON a.CHANNEL_ID = b.CHANNEL_ID
)t4
on t3.PRODUCT_CHANNEL_CODE=t4.PRODUCT_CHANNEL_CODE
id select_type table partitions type possible_keys key key_len ref rows filtered Extra
1 SIMPLE b null ALL PRIMARY,index_qd_channel null null null 212 100.0 null
1 SIMPLE a null ref index_qd_prochannel index_qd_prochannel 4 qd_dz_db.b.CHANNEL_ID 43 100.0 null
1 SIMPLE t1 null ALL null null null null 56367 10.0 Using where; Using join buffer (Block Nested Loop)
1 SIMPLE t2 null eq_ref PRIMARY,index_qd_channel PRIMARY 4 qd_dz_db.t1.CHANNEL_ID 1 100.0 null
1
show index from qd_prochannel
Table Non_unique Key_name Seq_in_index Column_name Collation Cardinality Sub_part Packed Null Index_type CommentIndex_comment
qd_prochannel 0 PRIMARY 1 PR_CHANNEL_ID A 56367 null null BTREE
qd_prochannel 1 index_qd_prochannel 1 CHANNEL_ID A 1311 null null BTREE
qd_prochannel 1 index_qd_prochannel 2 MAIN_PROJECT_CODE A 3316 null null YES BTREE
qd_prochannel 1 index_qd_prochannel 3 CREATE_TIME A 14092 null null YES BTREE
qd_prochannel 1 index_qd_prochannel_MAIN_PROJECT_CODE 1 MAIN_PROJECT_CODE A 613 null null YES BTREE
1
show index from qd_channel
Table Non_unique Key_name Seq_in_index Column_name Collation Cardinality Sub_part Packed Null Index_type Comment Index_comment
qd_channel 0 PRIMARY 1 CHANNEL_ID A 212 null null BTREE
qd_channel 1 index_qd_channel 1 CHANNEL_ID A 212 null null BTREE
qd_channel 1 index_qd_channel 2 CHANNEL_CODE A 212 null null YES BTREE
1