Event Time and Watermarks

目錄

文檔解讀

文檔路徑

/Application Development/Streaming (DataStream API)/Event Time

關(guān)于event time和watermark轿塔,內(nèi)容比較多姆吭,本文的順序不一定完全按照原文檔的順序進(jìn)行解讀

Event time

官方的定義

Event time is the time that each individual event occurred on its producing device.

事件時(shí)間就是數(shù)據(jù)流中事件實(shí)際發(fā)生的時(shí)間。事件時(shí)間也是實(shí)際業(yè)務(wù)中需要處理的時(shí)間,由于各種數(shù)據(jù)源的不同特點(diǎn)硬纤,可能在流中計(jì)算的時(shí)候會(huì)遇到事件的延遲或者事件時(shí)間的亂序,為解決這個(gè)問題出現(xiàn)了Watermark敦捧。

Watermark

官方對Watermark概念的說明可以參考兩個(gè)地方袜香,第一就是官網(wǎng)的原文

The mechanism in Flink to measure progress in event time is watermarks. Watermarks flow as part of the data stream and carry a timestamp t. A Watermark(t) declares that event time has reached time t in that stream, meaning that there should be no more elements from the stream with a timestamp t’ <= t (i.e. events with timestamps older or equal to the watermark).

第二就是Watermark類的API說明

A Watermark tells operators that no elements with a timestamp older or equal to the watermark timestamp should arrive at the operator. Watermarks are emitted at the sources and propagate through the operators of the topology. Operators must themselves emit watermarks to downstream operators usingOutput.emitWatermark(Watermark). Operators that do not internally buffer elements can always forward the watermark that they receive. Operators that buffer elements, such as window operators, must forward a watermark after emission of elements that is triggered by the arriving watermark.

In some cases a watermark is only a heuristic and operators should be able to deal with late elements. They can either discard those or update the result and emit updates/retractions to downstream operations.

When a source closes it will emit a final watermark with timestamp Long.MAX_VALUE. When an operator receives this it will know that no more input will be arriving in the future.

本文將針對第二種解釋做進(jìn)一步的說明,通過上文可以總結(jié)2點(diǎn):

  • Watermark中包含了一個(gè)時(shí)間戳乞封,作用就是告訴算子在這個(gè)時(shí)間戳之前的數(shù)據(jù)都已經(jīng)到達(dá)了做裙,不會(huì)再有小于或等于這個(gè)時(shí)間戳的數(shù)據(jù)再到達(dá)這個(gè)算子了。
  • Watermark可以認(rèn)為有一個(gè)“生命周期”肃晚,即出生锚贱,傳播,死亡关串。

Watermark生命周期

Watermark的“出生”

Watermark的產(chǎn)生有兩種方式拧廊,

  1. Source Functions中直接指定,請參考官文Source Functions with Timestamps and Watermarks
  2. 通過DataStream#assignTimestampsAndWatermarks方式指定如何生成Watermark晋修,這種方式又有兩種模式吧碾,時(shí)間驅(qū)動(dòng)的周期水位線(Periodic Watermarks)和數(shù)據(jù)驅(qū)動(dòng)的定點(diǎn)水位線(Punctuated Watermarks),對于周期水位線墓卦,官方提供了兩種實(shí)現(xiàn)可以參考BoundedOutOfOrdernessTimestampExtractor.javaAscendingTimestampExtractor.java倦春。
Watermark的傳播

此部分的內(nèi)容可以參見筆者之前的文章Watermarks in Parallel Streams

Watermark的“死亡”

當(dāng)watermark的時(shí)間戳變成Long.MAX_VALUE的時(shí)候,也就表示告訴算子再也沒有數(shù)據(jù)會(huì)到達(dá)了

擴(kuò)展閱讀

對于以周期模式產(chǎn)生watermark的時(shí)候落剪,官方給出的說明:

The interval (every n milliseconds) in which the watermark will be generated is defined via ExecutionConfig.setAutoWatermarkInterval(...).

對于這句話讀者要知道是:

  • 如果不指定溅漾,默認(rèn)是200毫秒
  • ExecutionConfig.setAutoWatermarkInterval一定要放到env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime)后面才會(huì)生效,因?yàn)樵?code>setStreamTimeCharacteristic里面會(huì)強(qiáng)制設(shè)置周期為200毫秒著榴,如果這個(gè)方法后執(zhí)行添履,就會(huì)覆蓋原有設(shè)置的周期
public void setStreamTimeCharacteristic(TimeCharacteristic characteristic) {
    this.timeCharacteristic = Preconditions.checkNotNull(characteristic);
    if (characteristic == TimeCharacteristic.ProcessingTime) {
        getConfig().setAutoWatermarkInterval(0);
    } else {
        getConfig().setAutoWatermarkInterval(200);
    }
}

這個(gè)參數(shù)是在TimestampsAndPeriodicWatermarksOperator#open中會(huì)拿到設(shè)置的watermarkInterval并將此值傳給timerService

public void open() throws Exception {
  super.open();

  currentWatermark = Long.MIN_VALUE;
  watermarkInterval = getExecutionConfig().getAutoWatermarkInterval();

  if (watermarkInterval > 0) {
    long now = getProcessingTimeService().getCurrentProcessingTime();
    getProcessingTimeService().registerTimer(now + watermarkInterval, this);
  }
}

這里getProcessingTimeService()返回的對象就是StreamTask中的timerService。

最后編輯于
?著作權(quán)歸作者所有,轉(zhuǎn)載或內(nèi)容合作請聯(lián)系作者
  • 序言:七十年代末脑又,一起剝皮案震驚了整個(gè)濱河市暮胧,隨后出現(xiàn)的幾起案子锐借,更是在濱河造成了極大的恐慌,老刑警劉巖往衷,帶你破解...
    沈念sama閱讀 218,682評論 6 507
  • 序言:濱河連續(xù)發(fā)生了三起死亡事件钞翔,死亡現(xiàn)場離奇詭異,居然都是意外死亡席舍,警方通過查閱死者的電腦和手機(jī)布轿,發(fā)現(xiàn)死者居然都...
    沈念sama閱讀 93,277評論 3 395
  • 文/潘曉璐 我一進(jìn)店門,熙熙樓的掌柜王于貴愁眉苦臉地迎上來来颤,“玉大人汰扭,你說我怎么就攤上這事「GΓ” “怎么了萝毛?”我有些...
    開封第一講書人閱讀 165,083評論 0 355
  • 文/不壞的土叔 我叫張陵,是天一觀的道長滑黔。 經(jīng)常有香客問我笆包,道長,這世上最難降的妖魔是什么略荡? 我笑而不...
    開封第一講書人閱讀 58,763評論 1 295
  • 正文 為了忘掉前任庵佣,我火速辦了婚禮,結(jié)果婚禮上汛兜,老公的妹妹穿的比我還像新娘秧了。我一直安慰自己,他們只是感情好序无,可當(dāng)我...
    茶點(diǎn)故事閱讀 67,785評論 6 392
  • 文/花漫 我一把揭開白布验毡。 她就那樣靜靜地躺著,像睡著了一般帝嗡。 火紅的嫁衣襯著肌膚如雪晶通。 梳的紋絲不亂的頭發(fā)上,一...
    開封第一講書人閱讀 51,624評論 1 305
  • 那天哟玷,我揣著相機(jī)與錄音狮辽,去河邊找鬼。 笑死巢寡,一個(gè)胖子當(dāng)著我的面吹牛喉脖,可吹牛的內(nèi)容都是我干的。 我是一名探鬼主播抑月,決...
    沈念sama閱讀 40,358評論 3 418
  • 文/蒼蘭香墨 我猛地睜開眼隧魄,長吁一口氣:“原來是場噩夢啊……” “哼输拇!你這毒婦竟也來了缀雳?” 一聲冷哼從身側(cè)響起,我...
    開封第一講書人閱讀 39,261評論 0 276
  • 序言:老撾萬榮一對情侶失蹤洁仗,失蹤者是張志新(化名)和其女友劉穎,沒想到半個(gè)月后性锭,有當(dāng)?shù)厝嗽跇淞掷锇l(fā)現(xiàn)了一具尸體赠潦,經(jīng)...
    沈念sama閱讀 45,722評論 1 315
  • 正文 獨(dú)居荒郊野嶺守林人離奇死亡,尸身上長有42處帶血的膿包…… 初始之章·張勛 以下內(nèi)容為張勛視角 年9月15日...
    茶點(diǎn)故事閱讀 37,900評論 3 336
  • 正文 我和宋清朗相戀三年草冈,在試婚紗的時(shí)候發(fā)現(xiàn)自己被綠了她奥。 大學(xué)時(shí)的朋友給我發(fā)了我未婚夫和他白月光在一起吃飯的照片。...
    茶點(diǎn)故事閱讀 40,030評論 1 350
  • 序言:一個(gè)原本活蹦亂跳的男人離奇死亡怎棱,死狀恐怖哩俭,靈堂內(nèi)的尸體忽然破棺而出,到底是詐尸還是另有隱情蹄殃,我是刑警寧澤携茂,帶...
    沈念sama閱讀 35,737評論 5 346
  • 正文 年R本政府宣布你踩,位于F島的核電站诅岩,受9級特大地震影響,放射性物質(zhì)發(fā)生泄漏带膜。R本人自食惡果不足惜吩谦,卻給世界環(huán)境...
    茶點(diǎn)故事閱讀 41,360評論 3 330
  • 文/蒙蒙 一、第九天 我趴在偏房一處隱蔽的房頂上張望膝藕。 院中可真熱鬧式廷,春花似錦、人聲如沸芭挽。這莊子的主人今日做“春日...
    開封第一講書人閱讀 31,941評論 0 22
  • 文/蒼蘭香墨 我抬頭看了看天上的太陽袜爪。三九已至蠕趁,卻和暖如春,著一層夾襖步出監(jiān)牢的瞬間辛馆,已是汗流浹背俺陋。 一陣腳步聲響...
    開封第一講書人閱讀 33,057評論 1 270
  • 我被黑心中介騙來泰國打工, 沒想到剛下飛機(jī)就差點(diǎn)兒被人妖公主榨干…… 1. 我叫王不留昙篙,地道東北人腊状。 一個(gè)月前我還...
    沈念sama閱讀 48,237評論 3 371
  • 正文 我出身青樓,卻偏偏與公主長得像苔可,于是被迫代替她去往敵國和親缴挖。 傳聞我的和親對象是個(gè)殘疾皇子,可洞房花燭夜當(dāng)晚...
    茶點(diǎn)故事閱讀 44,976評論 2 355