很可惜 T 。T 您現(xiàn)在還不是作者身份,不能自主發(fā)稿哦~
如有投稿需求,請把文章發(fā)送到郵箱tougao@appcpx.com,一經(jīng)錄用會有專人和您聯(lián)系
咨詢?nèi)绾纬蔀榇河鹱髡哒埪?lián)系:鳥哥筆記小羽毛(ngbjxym)
在普通的數(shù)據(jù)處理場景中,處理數(shù)據(jù)很簡單啊,因為數(shù)據(jù)都好好的放在庫里,直接select出來就好了。
但是流式數(shù)據(jù)是一條一條過來的,期間還會因為網(wǎng)絡(luò)延遲,有些數(shù)據(jù)還會遲到。這種“數(shù)據(jù)沒排好隊”的情況,叫做“亂序”。這可讓我們非常麻煩!
大家知道,所有數(shù)據(jù)理論上都應(yīng)該有時間戳,在流式數(shù)據(jù)中,時間戳更重要??梢哉f時間戳就是流式數(shù)據(jù)區(qū)別于離線數(shù)據(jù)的重要標(biāo)志。
在Flink中,我們大多使用EventTime作為時間戳。當(dāng)我們用這個時間來參與計算的時候,由于EventTime是真實世界的時間,那么百分之100可能會發(fā)生亂序數(shù)據(jù)。
那么何為亂序數(shù)據(jù)呢,前面說過了,亂序數(shù)據(jù)就是遲到的數(shù)據(jù)。1分鐘前產(chǎn)生的數(shù)據(jù),1分鐘之后才進(jìn)入到系統(tǒng)中,這就延遲了。
所以那么亂序數(shù)據(jù)就是在正常的時間數(shù)據(jù)流中夾雜著一些非順序的一些數(shù)據(jù)。
亂序是怎么產(chǎn)生的呢?因素太多了,例如某臺機(jī)器的網(wǎng)絡(luò)抖動,或者網(wǎng)卡和系統(tǒng)的延遲,都會導(dǎo)致這臺機(jī)器上報的數(shù)據(jù)延遲到達(dá)。
那么flink在處理的時候,就可能收到了系統(tǒng)在好幾秒之前產(chǎn)生的數(shù)據(jù)。這個一點非常討厭,會直接導(dǎo)致實時Join失敗。
Flink必須得解決這個問題啊,否則怎么保證遲到的數(shù)據(jù)都能用上呢,對吧?watermark就是用來解決這個問題的。
Watermark,就是水位線,用來測量亂序數(shù)據(jù)的進(jìn)度的。
Flink用watermark來確定這條遲到的數(shù)據(jù)如何觸發(fā)計算或者其他操作。嘿嘿,所以Watermark也是一種特殊的數(shù)據(jù)!
單純的從概念上不好理解,我們先假設(shè)一個場景,這樣更容易理解這個事情。你最好有一些流失數(shù)據(jù)的基礎(chǔ),否則不太容易理解這些原理。
假設(shè),我們有一個5s的窗口,并且我們可以容忍的延遲時間為2s,就是說5秒一計算,允許數(shù)據(jù)遲到2秒。
那么也就是說,從0開始,在7s的時候會觸發(fā)一次計算。我畫個圖解釋一下為什么會7s觸發(fā)計算,或者永不觸發(fā)。
為了排除其他影響因素,我們假設(shè)是單task,單分區(qū)的場景。其中的33 是第一條數(shù)據(jù), 2 是他攜帶的時間戳,在右側(cè)有一個5秒的窗口:
那么我們的watermark的計算公式就是 watermark = time - latertime 。那么這個時候我們可以得到這個watermark是0,那么他屬于0-5s的窗口,那么我們就放到窗口里面去。
這個時候又來了一條數(shù)據(jù),就會變成下面這樣對吧,為什么會變成兩個窗口呢?
因為99這條數(shù)據(jù)并不屬于0-5秒這個窗口里面,因為flink窗口的大小是包左不包右的,這點很關(guān)鍵。
這樣你就能明白,為什么33和99應(yīng)該各自進(jìn)到單獨的窗口。所以,數(shù)據(jù)是根據(jù)EventTime來決定應(yīng)該進(jìn)哪個桶或者說窗口的。
現(xiàn)在你能理解為啥EventTime這么重要了吧?
假如,這個時候來了一條亂序數(shù)據(jù),23號(時間戳3S),這條數(shù)據(jù)遲到了,那么我們的watermark怎么更新呢?
現(xiàn)在,請你停下來思考一下,新來的23號數(shù)據(jù)應(yīng)該進(jìn)那個窗口?
案揭曉:我們可以看到我畫的圖:其中,23數(shù)據(jù)攜帶的時間戳是3,watermark也是3,應(yīng)該歸到[0,5)的窗口。
你是不是會奇怪,這遲到的數(shù)據(jù)序號比前面的99號序號要大啊,怎么在后面呢,并且計算出來的watermark是3?這不是違背了我們的公式計算規(guī)則么?
按照前面的公式,watermark = time - latertime,那么23號的watermark應(yīng)該是3-2=1,應(yīng)該排到99號的前面去啊。
其實不是的,watermark首先是時間尺度,然后才是衡量標(biāo)準(zhǔn)。所以watermark 不能倒著走啊,因為他是負(fù)責(zé)測量數(shù)據(jù)的時間進(jìn)度的。
所以他的watermark 并不會按照公式計算,而是采用的上一個數(shù)字的watermark,也就是3。
為了讓你看的更清楚,我們多插幾條數(shù)據(jù)看看。
所有數(shù)據(jù)在watermark上,都是順序排列的,6號數(shù)據(jù)的watermark,按照公式,應(yīng)該是4-2=2,但是很遺憾,前面已經(jīng)有3了,所以只能排在3后面。誰讓你遲到了呢,對吧?
當(dāng)然,這個時候他們只是在排隊,還沒有觸發(fā)窗口的計算操作。
那么窗口計算什么時候觸發(fā)呢?很簡單,當(dāng)watermark大于等于窗口觸發(fā)時間。
第一個窗口觸發(fā)計算操作的時間也就是大于等于5秒的時候。
提問:第二個窗口出發(fā)計算操作的時間是什么時候呢?
答案是10秒。
那么我們講到這里,應(yīng)該大部分人都能夠了解到了watermark的運(yùn)行機(jī)制,以及窗口什么時候計算。
現(xiàn)在還是一條線的情況?,F(xiàn)實情況比這個要復(fù)雜的多的多。
那么我們接下來就來考慮一下我們的多并行度下,我們的watermark如何傳遞?
我們應(yīng)該知道 在一個 task中有很多的subtask,那么這些subtask都有自己的watermark。
所有的數(shù)據(jù)時間上都應(yīng)該同步啊,要不然怎么多并行度計算???就全亂套了。
所以這個時候就會涉及到 watermark的傳遞,因為下游也是依賴這些watermark的。
如上圖所示,我們可以看到Watermark在順序的向下游流動,左側(cè)的向右箭頭,就是這個意思。
那么我們這個時候發(fā)現(xiàn)有一個Partition WM 這個其實就是各個分區(qū)的 SubTask的Watermark。
我這個時候發(fā)現(xiàn),每個subtask的watermark都是不一樣的,并且task會存儲這些watermark,記錄下來各個分區(qū)的watermark,并且把最小的watermark廣播出去。
當(dāng)前需要記錄的是2、4、3、6號watermark,其中2是最小的,好,我們記下來。
這個時候當(dāng)傳遞過來的4號watermark更新了,把原來的2給頂走了。那么現(xiàn)在是4、4、3、6,最小的是3了。這個時候我們就將最小的3作為watermark傳遞出去。
當(dāng)7傳遞過來的時候,我們就會發(fā)現(xiàn),傳遞過去的依然是最小的3,所以不動。
這樣,我們就能解決多并行度下,watermark的傳遞問題。其實就是挑一個最小的watermark放出去。
你現(xiàn)在對watermark的機(jī)制應(yīng)該是比較了解的吧?
注意: WatermarkGenerator 將以前互相獨立的 {@code AssignerWithPunctuatedWatermarks} 調(diào)用此方法生成 watermark 的間隔時間由 {@link ExecutionConfig#getAutoWatermarkInterval()} 決定。 ?在數(shù)據(jù)源直接使用時如果因為數(shù)據(jù)源中的某一個分區(qū)/分片在一段時間內(nèi)未發(fā)送事件數(shù)據(jù),則意味著watermarkStrategy也不會獲得任何數(shù)據(jù)去生成watermark,在這種情況下可以通過設(shè)置有一個空閑時間,當(dāng)超過這個時間則將這個分片或分區(qū)標(biāo)記為空閑狀態(tài)。 WatermarkStrategy . Flink的很多設(shè)計都非常精巧,watermark就是其中之一。我們研究這些實現(xiàn)原理并不是想做源碼級的開發(fā),而是欣賞這種精妙的思想,真是為之嘆息。 如果你覺得有啟發(fā),歡迎留言,一起交流。 -END-
* 和 {@code AssignerWithPeriodicWatermarks} 一同包含了進(jìn)來。
*/
@Public
public interface WatermarkGenerator
/**
* 每來一條事件數(shù)據(jù)調(diào)用一次,可以檢查或者記錄事件的時間戳,或者也可以基于事件數(shù)據(jù)本身去生成 watermark。
*/
void onEvent(T event, long eventTimestamp, WatermarkOutput output);
/**
* 周期性的調(diào)用,也許會生成新的 watermark,也許不會。
*
*
*/
void onPeriodicEmit(WatermarkOutput output);
}3、watermark 分區(qū)數(shù)據(jù)傾斜解決方案
結(jié)語
本文為作者獨立觀點,不代表鳥哥筆記立場,未經(jīng)允許不得轉(zhuǎn)載。
《鳥哥筆記版權(quán)及免責(zé)申明》 如對文章、圖片、字體等版權(quán)有疑問,請點擊 反饋舉報
我們致力于提供一個高質(zhì)量內(nèi)容的交流平臺。為落實國家互聯(lián)網(wǎng)信息辦公室“依法管網(wǎng)、依法辦網(wǎng)、依法上網(wǎng)”的要求,為完善跟帖評論自律管理,為了保護(hù)用戶創(chuàng)造的內(nèi)容、維護(hù)開放、真實、專業(yè)的平臺氛圍,我們團(tuán)隊將依據(jù)本公約中的條款對注冊用戶和發(fā)布在本平臺的內(nèi)容進(jìn)行管理。平臺鼓勵用戶創(chuàng)作、發(fā)布優(yōu)質(zhì)內(nèi)容,同時也將采取必要措施管理違法、侵權(quán)或有其他不良影響的網(wǎng)絡(luò)信息。
一、根據(jù)《網(wǎng)絡(luò)信息內(nèi)容生態(tài)治理規(guī)定》《中華人民共和國未成年人保護(hù)法》等法律法規(guī),對以下違法、不良信息或存在危害的行為進(jìn)行處理。
1. 違反法律法規(guī)的信息,主要表現(xiàn)為:
1)反對憲法所確定的基本原則;
2)危害國家安全,泄露國家秘密,顛覆國家政權(quán),破壞國家統(tǒng)一,損害國家榮譽(yù)和利益;
3)侮辱、濫用英烈形象,歪曲、丑化、褻瀆、否定英雄烈士事跡和精神,以侮辱、誹謗或者其他方式侵害英雄烈士的姓名、肖像、名譽(yù)、榮譽(yù);
4)宣揚(yáng)恐怖主義、極端主義或者煽動實施恐怖活動、極端主義活動;
5)煽動民族仇恨、民族歧視,破壞民族團(tuán)結(jié);
6)破壞國家宗教政策,宣揚(yáng)邪教和封建迷信;
7)散布謠言,擾亂社會秩序,破壞社會穩(wěn)定;
8)宣揚(yáng)淫穢、色情、賭博、暴力、兇殺、恐怖或者教唆犯罪;
9)煽動非法集會、結(jié)社、游行、示威、聚眾擾亂社會秩序;
10)侮辱或者誹謗他人,侵害他人名譽(yù)、隱私和其他合法權(quán)益;
11)通過網(wǎng)絡(luò)以文字、圖片、音視頻等形式,對未成年人實施侮辱、誹謗、威脅或者惡意損害未成年人形象進(jìn)行網(wǎng)絡(luò)欺凌的;
12)危害未成年人身心健康的;
13)含有法律、行政法規(guī)禁止的其他內(nèi)容;
2. 不友善:不尊重用戶及其所貢獻(xiàn)內(nèi)容的信息或行為。主要表現(xiàn)為:
1)輕蔑:貶低、輕視他人及其勞動成果;
2)誹謗:捏造、散布虛假事實,損害他人名譽(yù);
3)嘲諷:以比喻、夸張、侮辱性的手法對他人或其行為進(jìn)行揭露或描述,以此來激怒他人;
4)挑釁:以不友好的方式激怒他人,意圖使對方對自己的言論作出回應(yīng),蓄意制造事端;
5)羞辱:貶低他人的能力、行為、生理或身份特征,讓對方難堪;
6)謾罵:以不文明的語言對他人進(jìn)行負(fù)面評價;
7)歧視:煽動人群歧視、地域歧視等,針對他人的民族、種族、宗教、性取向、性別、年齡、地域、生理特征等身份或者歸類的攻擊;
8)威脅:許諾以不良的后果來迫使他人服從自己的意志;
3. 發(fā)布垃圾廣告信息:以推廣曝光為目的,發(fā)布影響用戶體驗、擾亂本網(wǎng)站秩序的內(nèi)容,或進(jìn)行相關(guān)行為。主要表現(xiàn)為:
1)多次發(fā)布包含售賣產(chǎn)品、提供服務(wù)、宣傳推廣內(nèi)容的垃圾廣告。包括但不限于以下幾種形式:
2)單個帳號多次發(fā)布包含垃圾廣告的內(nèi)容;
3)多個廣告帳號互相配合發(fā)布、傳播包含垃圾廣告的內(nèi)容;
4)多次發(fā)布包含欺騙性外鏈的內(nèi)容,如未注明的淘寶客鏈接、跳轉(zhuǎn)網(wǎng)站等,誘騙用戶點擊鏈接
5)發(fā)布大量包含推廣鏈接、產(chǎn)品、品牌等內(nèi)容獲取搜索引擎中的不正當(dāng)曝光;
6)購買或出售帳號之間虛假地互動,發(fā)布干擾網(wǎng)站秩序的推廣內(nèi)容及相關(guān)交易。
7)發(fā)布包含欺騙性的惡意營銷內(nèi)容,如通過偽造經(jīng)歷、冒充他人等方式進(jìn)行惡意營銷;
8)使用特殊符號、圖片等方式規(guī)避垃圾廣告內(nèi)容審核的廣告內(nèi)容。
4. 色情低俗信息,主要表現(xiàn)為:
1)包含自己或他人性經(jīng)驗的細(xì)節(jié)描述或露骨的感受描述;
2)涉及色情段子、兩性笑話的低俗內(nèi)容;
3)配圖、頭圖中包含庸俗或挑逗性圖片的內(nèi)容;
4)帶有性暗示、性挑逗等易使人產(chǎn)生性聯(lián)想;
5)展現(xiàn)血腥、驚悚、殘忍等致人身心不適;
6)炒作緋聞、丑聞、劣跡等;
7)宣揚(yáng)低俗、庸俗、媚俗內(nèi)容。
5. 不實信息,主要表現(xiàn)為:
1)可能存在事實性錯誤或者造謠等內(nèi)容;
2)存在事實夸大、偽造虛假經(jīng)歷等誤導(dǎo)他人的內(nèi)容;
3)偽造身份、冒充他人,通過頭像、用戶名等個人信息暗示自己具有特定身份,或與特定機(jī)構(gòu)或個人存在關(guān)聯(lián)。
6. 傳播封建迷信,主要表現(xiàn)為:
1)找人算命、測字、占卜、解夢、化解厄運(yùn)、使用迷信方式治??;
2)求推薦算命看相大師;
3)針對具體風(fēng)水等問題進(jìn)行求助或咨詢;
4)問自己或他人的八字、六爻、星盤、手相、面相、五行缺失,包括通過占卜方法問婚姻、前程、運(yùn)勢,東西寵物丟了能不能找回、取名改名等;
7. 文章標(biāo)題黨,主要表現(xiàn)為:
1)以各種夸張、獵奇、不合常理的表現(xiàn)手法等行為來誘導(dǎo)用戶;
2)內(nèi)容與標(biāo)題之間存在嚴(yán)重不實或者原意扭曲;
3)使用夸張標(biāo)題,內(nèi)容與標(biāo)題嚴(yán)重不符的。
8.「飯圈」亂象行為,主要表現(xiàn)為:
1)誘導(dǎo)未成年人應(yīng)援集資、高額消費(fèi)、投票打榜
2)粉絲互撕謾罵、拉踩引戰(zhàn)、造謠攻擊、人肉搜索、侵犯隱私
3)鼓動「飯圈」粉絲攀比炫富、奢靡享樂等行為
4)以號召粉絲、雇用網(wǎng)絡(luò)水軍、「養(yǎng)號」形式刷量控評等行為
5)通過「蹭熱點」、制造話題等形式干擾輿論,影響傳播秩序
9. 其他危害行為或內(nèi)容,主要表現(xiàn)為:
1)可能引發(fā)未成年人模仿不安全行為和違反社會公德行為、誘導(dǎo)未成年人不良嗜好影響未成年人身心健康的;
2)不當(dāng)評述自然災(zāi)害、重大事故等災(zāi)難的;
3)美化、粉飾侵略戰(zhàn)爭行為的;
4)法律、行政法規(guī)禁止,或可能對網(wǎng)絡(luò)生態(tài)造成不良影響的其他內(nèi)容。
二、違規(guī)處罰
本網(wǎng)站通過主動發(fā)現(xiàn)和接受用戶舉報兩種方式收集違規(guī)行為信息。所有有意的降低內(nèi)容質(zhì)量、傷害平臺氛圍及欺凌未成年人或危害未成年人身心健康的行為都是不能容忍的。
當(dāng)一個用戶發(fā)布違規(guī)內(nèi)容時,本網(wǎng)站將依據(jù)相關(guān)用戶違規(guī)情節(jié)嚴(yán)重程度,對帳號進(jìn)行禁言 1 天、7 天、15 天直至永久禁言或封停賬號的處罰。當(dāng)涉及欺凌未成年人、危害未成年人身心健康、通過作弊手段注冊、使用帳號,或者濫用多個帳號發(fā)布違規(guī)內(nèi)容時,本網(wǎng)站將加重處罰。
三、申訴
隨著平臺管理經(jīng)驗的不斷豐富,本網(wǎng)站出于維護(hù)本網(wǎng)站氛圍和秩序的目的,將不斷完善本公約。
如果本網(wǎng)站用戶對本網(wǎng)站基于本公約規(guī)定做出的處理有異議,可以通過「建議反饋」功能向本網(wǎng)站進(jìn)行反饋。
(規(guī)則的最終解釋權(quán)歸屬本網(wǎng)站所有)