0%

Flink:时间

Flink Time

Event Time与Process Time

如果你对数据的准确性要求比较高的话,采用 Event time 能保障 exactly-once。Processing Time 一般用于实时消费、精准性要求略低的场景,主要是因为时间生成不是 deterministic。

我们可以看下面的关系图, X 轴是 Event time,Y 轴是 Processing time。理想情况下 Event timeProcessing Time 是相同的,就是说只要有一个事件发生,就可以立刻处理。但是实际场景中,事件发生后往往会经过一定延时才会被处理,这样就会导致我们系统的时间往往会滞后于事件时间。这里它们两个的差 Processing-time lag 表示我们处理事件的延时。

事件时间常用在窗口中,使用 watermark 来确保数据完备性,比如说 watermarker 值大于 window 末尾时间时,我们就可以认为 window 窗口所有数据都已经到达了,就可以触发计算了。

比如上面 [0-10] 的窗口,现在 watermark 走到了 10,已经到达了窗口的结束,触发计算 SUM=21。如果要是想对迟到的数据再进行触发,可以再定义一下后面 late data 的触发,比如说后面来了个 9,我们的 SUM 就等于 30。

北京时间与格林尼治时间

北京时间

划定时间窗口和时间戳以毫秒为单位

1
2
3
4
5
eventStream.assignTimestampsAndWatermarks(assigner = new AscendingTimestampExtractor[Event]() {
override def extractAscendingTimestamp(t: Event): Long = {
t.getTime * 1000
}
})

processFunction中的时间

1
2
3
4
new ProcessFunction[Event, (String, Long)] {

override def processElement(event: Event, context: ProcessFunction[Event,
(String, Long)]#Context, collector: Collector[(String, Long)]): Unit = {
1
2
3
4
5
context.timestamp=1565064599999
context.timerService().currentWatermark()=-9223372036854775808
context.timerService().currentProcessingTime()=1565064607343

Event(time=1565064350, count=9, url=/apply/main, channelId=1057)

以事件事件,窗口大小5min

一条消息,事件的时间是 1565064350(2019-08-06 12:05:50)
Context.timestamp是1565064599999(2019-08-06 12:09:59), 可以理解为窗口结束的时间
context.timerService().currentWatermark() 莫名其妙,这么大
context.timerService().currentProcessingTime()是当前处理事件 2019-08-06 12:10:07

Context.timestamp

1
2
Timestamp of the element currently being processed or timestamp of a firing timer.
This might be {@code null}, for example if the time characteristic of your program

时间戳timestamp与水印watermark

提取Timestamp与生成Watermark一般步骤

  1. 设置时间特性为Event Time。StreamExecutionEnvironment#setStreamTimeCharacteristic(TimeCharacteristic.EventTime)
  2. 在Source后Window前用DataStream#assignTimestampsAndWatermarks方法(AssignerWithPeriodicWatermarksAssignerWithPunctuatedWatermarks)提取时间戳并生成水印。
  3. 重写extractTimestamp方法提取Timestamp,重写getCurrentWatermark方法或checkAndGetNextWatermark方法生成水印。

数据源中指定

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
public static class ExampleSourceFunction implements SourceFunction<Tuple4<String,Long,String,Integer>>{

@Override
public void run(SourceContext<Tuple4<String,Long,String,Integer>> ctx) throws Exception {
while (isRunning){
// 构造测试数据
....

// 发出一条数据以及数据对应的Timestamp
ctx.collectWithTimestamp(record,eventTime);

// 发出一条Watermark
ctx.emitWatermark(new Watermark(eventTime - maxOutOfOrderness));

Thread.sleep(1000);
}
}

@Override
public void cancel() {
isRunning = false;
}
}
}

时间戳和水位线分别通过两个方法实现的

时间戳的抽取方法

1
2
3
4
/**
* 传入了两个参数,第二个参数是上一个时间戳,元素的当前内部时间戳,如果尚未分配时间戳,则为负值。
*/
long extractTimestamp(T element, long recordTimestamp);

周期性水印

分配器为元素分配事件时间时间戳,并生成表示流中事件时间进度的低水位线。这些时间戳和水印由对事件时间进行操作的函数和运算符使用,例如事件时间窗口。

使用此类以周期性间隔生成水印。最多每隔I毫秒(通过ExecutionConfig.getautowerMarkinterval())配置一次,系统将调用getCurrentWatermark()方法来探测下一个水印值。
如果探测到的值为非空,并且时间戳大于前一个水印的时间戳,系统将生成新的水印(以保留升序水印的约定)。 如果自上次调用getCurrentWatermark()方法后没有新元素到达,系统调用该方法的频率可能会低于每I毫秒一次。
时间戳和水印被定义为代表自纪元(世界协调时1970年1月1日午夜)以来的毫秒的长度。某个值为t的水印表示不会再出现带有事件时间戳x的元素,其中x小于或等于t。

1
2
3
4
5
6
/**
* 返回当前水印。系统定期调用此方法来检索当前水印。该方法可能返回空值,表示没有新的水印可用。
* 只有当返回的水印为非空且其时间戳大于之前发出的水印的时间戳时,才会发出该水印(以保留升序水印的约定)。
* 如果当前水印仍然与前一个相同,则自上一次调用此方法以来,事件时间没有任何进展。如果返回空值,或者返回的水印的时间戳小于上次发出的时间戳,则不会生成新的水印。
*/
Watermark AssignerWithPeriodicWatermarks.getCurrentWatermark()

非连续性水印

1
2
3
4
5
6
/**
* 询问此实现是否要发出水印。这个方法是在extractTimestamp(对象,长)方法之后调用的。
* 只有当返回的水印为非空且其时间戳大于之前发出的水印的时间戳时,才会发出该水印(以保留升序水印的约定)。
* 如果返回空值,或者返回的水印的时间戳小于上次发出的时间戳,则不会生成新的水印。 有关如何使用此方法的示例,请参见此类的分配器带标点水印的文档。
*/
Watermark AssignerWithPunctuatedWatermarks.checkAndGetNextWatermark(T lastElement, long extractedTimestamp)

升序水印 AscendingTimestampExtractor

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
org.apache.flink.streaming.api.functions.timestamps.AscendingTimestampExtractor

public abstract class AscendingTimestampExtractor<T> implements AssignerWithPeriodicWatermarks<T> {
/** The current timestamp. */
private long currentTimestamp = Long.MIN_VALUE;
/** Handler that is called when timestamp monotony is violated. */
private MonotonyViolationHandler violationHandler = new LoggingHandler();
/**
* Extracts the timestamp from the given element. The timestamp must be monotonically increasing.
*
* @param element The element that the timestamp is extracted from.
* @return The new timestamp.
*/
public abstract long extractAscendingTimestamp(T element);
/**
* Sets the handler for violations to the ascending timestamp order.
*
* @param handler The violation handler to use.
* @return This extractor.
*/
public AscendingTimestampExtractor<T> withViolationHandler(MonotonyViolationHandler handler) {
this.violationHandler = requireNonNull(handler);
return this;
}

// ------------------------------------------------------------------------
// 调用
@Override
public final long extractTimestamp(T element, long elementPrevTimestamp) {
final long newTimestamp = extractAscendingTimestamp(element);
if (newTimestamp >= this.currentTimestamp) {
this.currentTimestamp = newTimestamp;
return newTimestamp;
} else {
violationHandler.handleViolation(newTimestamp, this.currentTimestamp);
return newTimestamp;
}
}

@Override
public final Watermark getCurrentWatermark() {
return new Watermark(currentTimestamp == Long.MIN_VALUE ? Long.MIN_VALUE : currentTimestamp - 1);
}
...
}

AscendingTimestampExtractor抽象类实现了ssignerWithPeriodicWatermarks接口的extractTimestamp及getCurrentWatermark方法,同时声明抽象方法extractAscendingTimestamp供子类实现
水印的生产策略为:getCurrentWatermark 方法在currentTimestamp不为Long.MIN_VALUE时返回Watermark(currentTimestamp - 1)

综上:
发现AscendingTimestampExtractor适用于elements的时间在每个并行task里面事单调递增的,(timestamp monotony)数据场景。

无序水印 BoundedOutOfOrdernessTimestampExtractor

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
public abstract class BoundedOutOfOrdernessTimestampExtractor<T> implements AssignerWithPeriodicWatermarks<T> {

private static final long serialVersionUID = 1L;

/** The current maximum timestamp seen so far. */
private long currentMaxTimestamp;

/** The timestamp of the last emitted watermark. */
private long lastEmittedWatermark = Long.MIN_VALUE;

/**
* The (fixed) interval between the maximum seen timestamp seen in the records
* and that of the watermark to be emitted.
*/
private final long maxOutOfOrderness;

public BoundedOutOfOrdernessTimestampExtractor(Time maxOutOfOrderness) {
if (maxOutOfOrderness.toMilliseconds() < 0) {
throw new RuntimeException("Tried to set the maximum allowed " +
"lateness to " + maxOutOfOrderness + ". This parameter cannot be negative.");
}
this.maxOutOfOrderness = maxOutOfOrderness.toMilliseconds();
this.currentMaxTimestamp = Long.MIN_VALUE + this.maxOutOfOrderness;
}

public long getMaxOutOfOrdernessInMillis() {
return maxOutOfOrderness;
}

/**
* Extracts the timestamp from the given element.
*
* @param element The element that the timestamp is extracted from.
* @return The new timestamp.
*/
public abstract long extractTimestamp(T element);

@Override
public final Watermark getCurrentWatermark() {
// this guarantees that the watermark never goes backwards.
long potentialWM = currentMaxTimestamp - maxOutOfOrderness;
if (potentialWM >= lastEmittedWatermark) {
lastEmittedWatermark = potentialWM;
}
return new Watermark(lastEmittedWatermark);
}

@Override
public final long extractTimestamp(T element, long previousElementTimestamp) {
long timestamp = extractTimestamp(element);
if (timestamp > currentMaxTimestamp) {
currentMaxTimestamp = timestamp;
}
return timestamp;
}
}

BoundedOutOfOrdernessTimestampExtractor 抽象类实现AssignerWithPeriodicWatermarks接口的extractTimestamp及getCurrentWatermark方法,同时声明抽象方法extractAscendingTimestamp供子类实现
BoundedOutOfOrdernessTimestampExtractor 的构造器接收maxOutOfOrderness参数用于指定element允许滞后(t~`t_wt为element的eventTime,t_w为前一次watermark的时间)的最大时间,在计算窗口数据时,如果超过该值则会被忽略。 BoundedOutOfOrdernessTimestampExtractorextractTimestamp方法会调用子类的extractTimestamp方法抽取时间,如果该时间大于currentMaxTimestamp,则更新currentMaxTimestampgetCurrentWatermark先计算potentialWM,如果potentialWM大于等于lastEmittedWatermark则更新lastEmittedWatemakr(currentMaxTimestamp - lastEmittedWatermark >= maxOutOfOrderness, 这里表示lastEmittedWatermark太小了,所以差值超过了maxOutOfOrderness,因此会调大lastEmittedWatermark),最后返回watermark`

具体逻辑参考:就是判断如果新的水印时间大于上次最后一次发出水印时间,选择新的水印时间,否则选择上次最后发送水印时间戳

TimeStamp分配器和Watermark生成器(Timestamp Assigners / Watermark Generators)

数据延迟

主要的办法是给定一个允许延迟的时间,在该时间范围内仍可以接受处理延迟数据

设置允许延迟的时间是通过allowedLateness(lateness: Time)设置

保存延迟数据则是通过sideOutputLateData(outputTag: OutputTag[T])保存

获取延迟数据是通过DataStream.getSideOutput(tag: OutputTag[X])获取

allowedLateness(lateness: Time)

1
2
3
4
def allowedLateness(lateness: Time): WindowedStream[T, K, W] = {
javaStream.allowedLateness(lateness)
this
}

该方法传入一个Time值,设置允许数据迟到的时间,这个时间和waterMark中的时间概念不同。再来回顾一下,

waterMark=数据的事件时间-允许乱序时间值

随着新数据的到来,waterMark的值会更新为最新数据事件时间-允许乱序时间值,但是如果这时候来了一条历史数据,waterMark值则不会更新。总的来说,waterMark是为了能接收到尽可能多的乱序数据。

那这里的Time值呢?主要是为了等待迟到的数据,在一定时间范围内,如果属于该窗口的数据到来,仍会进行计算,后面会对计算方式仔细说明

allowedLateness 是等待数据进入窗口的最大时间,未进入的数据,可以通过偏流导出

注意:该方法只针对于基于event-time的窗口,如果是基于processing-time,并且指定了非零的time值则会抛出异常

sideOutputLateData(outputTag: OutputTag[T])

1
2
3
4
def sideOutputLateData(outputTag: OutputTag[T]): WindowedStream[T, K, W] = {
javaStream.sideOutputLateData(outputTag)
this
}

该方法是将迟来的数据保存至给定的outputTag参数,而OutputTag则是用来标记延迟数据的一个对象。

DataStream.getSideOutput(tag: OutputTag[X])

通过window等操作返回的DataStream调用该方法,传入标记延迟数据的对象来获取延迟的数据

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
waterStream.keyBy(0)
.window(TumblingEventTimeWindows.of(Time.seconds(5L)))

.allowedLateness(Time.seconds(2L))
.sideOutputLateData(lateData)
.apply(new WindowFunction[(String, Long), String, Tuple, TimeWindow] {
override def apply(key: Tuple, window: TimeWindow, input: Iterable[(String, Long)], out: Collector[String]): Unit = {
val timeArr = ArrayBuffer[String]()
val iterator = input.iterator
while (iterator.hasNext) {
val tup2 = iterator.next()
timeArr.append(sdf.format(tup2._2))
}
val outData = String.format("key: %s data: %s startTime: %s endTime: %s",
key.toString,
timeArr.mkString("-"),
sdf.format(window.getStart),
sdf.format(window.getEnd))
out.collect(outData)
}
})
result.print("window计算结果:")

val late = result.getSideOutput(lateData)
late.print("迟到的数据:")

【参考文献】

  1. Flink时间系列-EventTime下数据延迟处理
  2. Flink: 时间属性深度解析