時(shí)數(shù)據(jù)分析平臺(tái)實(shí)戰(zhàn):從數(shù)據(jù)采集到可視化大屏)
1. 從需求到架構(gòu)先想清楚實(shí)時(shí)到底意味著什么接到一個(gè)實(shí)時(shí)數(shù)據(jù)分析平臺(tái)的需求時(shí)我心里第一反應(yīng)不是寫(xiě)Flink代碼而是先問(wèn)對(duì)方一句話你說(shuō)的實(shí)時(shí)是秒級(jí)、分鐘級(jí)還是小時(shí)級(jí)這個(gè)問(wèn)題問(wèn)出來(lái)很多需求就瞬間清晰了。我見(jiàn)過(guò)太多團(tuán)隊(duì)一上來(lái)就鋪Flink集群、搞大屏結(jié)果做了三個(gè)月發(fā)現(xiàn)核心指標(biāo)延遲五分鐘就能滿足白費(fèi)了一堆功夫。Java工程師做實(shí)時(shí)數(shù)據(jù)平臺(tái)有個(gè)天然優(yōu)勢(shì)Flink本身就是Java/Scala生態(tài)你熟悉的Spring Boot、Maven、JVM調(diào)優(yōu)經(jīng)驗(yàn)全部能復(fù)用。相比Python系或者純SQL系的數(shù)據(jù)棧Java團(tuán)隊(duì)啃Flink的上手成本要低得多。這篇實(shí)戰(zhàn)文章我就圍繞數(shù)據(jù)采集→實(shí)時(shí)計(jì)算→可視化大屏這條完整鏈路把從零搭建一個(gè)實(shí)時(shí)數(shù)據(jù)分析平臺(tái)的架構(gòu)思路、核心代碼、踩坑記錄都講透。1.1 先拆需求你的大屏是不是假實(shí)時(shí)大多數(shù)實(shí)時(shí)數(shù)據(jù)分析平臺(tái)的真實(shí)需求拆開(kāi)來(lái)看無(wú)非三類(lèi)指標(biāo)監(jiān)控類(lèi)比如訂單量、成交額、在線用戶數(shù)要求秒級(jí)或分鐘級(jí)刷新行為分析類(lèi)用戶點(diǎn)擊流、頁(yè)面路徑要求準(zhǔn)實(shí)時(shí)但允許一定延遲預(yù)警通知類(lèi)比如異常流量、交易失敗率飆升要求延遲越低越好這三類(lèi)需求對(duì)技術(shù)選型的影響完全不同。我之前遇到一個(gè)做污水處理可視化大屏的項(xiàng)目客戶說(shuō)實(shí)時(shí)結(jié)果詳細(xì)了解才知道污水?dāng)?shù)據(jù)本身是五分鐘采集一次那你就算用Flink做到毫秒級(jí)計(jì)算也沒(méi)有意義瓶頸在采集端。反過(guò)來(lái)如果是電商大促的實(shí)時(shí)成交大屏每秒鐘都有成千上萬(wàn)條訂單事件那你就需要認(rèn)真設(shè)計(jì)從采集到展示的每一層。所以第一步永遠(yuǎn)是做延遲預(yù)算端到端延遲 采集延遲 傳輸延遲 計(jì)算延遲 存儲(chǔ)延遲 展示刷新延遲。把每一項(xiàng)都列出來(lái)標(biāo)出可接受范圍后續(xù)所有技術(shù)決策都有依據(jù)。1.2 端到端鏈路的分層設(shè)計(jì)我做的實(shí)時(shí)數(shù)據(jù)分析平臺(tái)標(biāo)準(zhǔn)鏈路分五層層級(jí)組件選型職責(zé)采集層Filebeat / Flink CDC / HTTP SDK將日志、數(shù)據(jù)庫(kù)變更、業(yè)務(wù)事件統(tǒng)一送入消息隊(duì)列傳輸層Kafka削峰填谷、緩沖削流解耦采集與計(jì)算計(jì)算層Flink實(shí)時(shí)ETL、窗口聚合、狀態(tài)計(jì)算、規(guī)則匹配存儲(chǔ)層Doris / ClickHouse / Redis結(jié)果表存儲(chǔ)、維度數(shù)據(jù)緩存、大屏查詢加速展示層Vue ECharts / DataV可視化大屏、指標(biāo)卡片、趨勢(shì)圖表這套鏈路跟傳統(tǒng)的離線數(shù)倉(cāng)最大的區(qū)別在于數(shù)據(jù)不是按天批量加工而是以事件流的方式持續(xù)流動(dòng)。Flink跑在Kafka和存儲(chǔ)之間相當(dāng)于一個(gè)永不停止的計(jì)算引擎——上游數(shù)據(jù)來(lái)了就算算完就寫(xiě)寫(xiě)完后端到端延遲通常控制在秒級(jí)。關(guān)于架構(gòu)理念現(xiàn)階段我做項(xiàng)目基本直接采用Kappa架構(gòu)思路不再搭建Lambda架構(gòu)。Lambda那套實(shí)時(shí)鏈路離線鏈路雙跑、最終結(jié)果合并的方案維護(hù)成本太高兩套代碼邏輯要一致本身就是災(zāi)難。現(xiàn)在Flink的流批一體能力已經(jīng)相當(dāng)成熟一套代碼可以同時(shí)跑實(shí)時(shí)和離線Kappa架構(gòu)足夠覆蓋絕大多數(shù)場(chǎng)景。1.3 為什么選Flink而不是Spark Streaming每次做技術(shù)選型都要面對(duì)這個(gè)問(wèn)題。我的答案很直接如果你的場(chǎng)景需要事件時(shí)間處理、精確一次語(yǔ)義、豐富的狀態(tài)管理Flink是當(dāng)前最優(yōu)解。事件時(shí)間處理數(shù)據(jù)在網(wǎng)絡(luò)上傳輸會(huì)有延遲和亂序Flink的Watermark機(jī)制可以基于事件真正發(fā)生的時(shí)間進(jìn)行計(jì)算而不是基于數(shù)據(jù)到達(dá)時(shí)間。這在處理日志類(lèi)數(shù)據(jù)時(shí)尤其重要——用戶點(diǎn)擊發(fā)生在10:00:00但因?yàn)榫W(wǎng)絡(luò)抖動(dòng)這條日志10:00:10才到如果你用處理時(shí)間計(jì)算就把這10秒的誤差算進(jìn)指標(biāo)里了。精確一次語(yǔ)義Exactly-OnceFlink通過(guò)Checkpoint 兩階段提交保證即使任務(wù)崩潰恢復(fù)數(shù)據(jù)也不會(huì)重復(fù)或丟失。做交易類(lèi)指標(biāo)時(shí)這是剛需。狀態(tài)管理Flink可以把中間結(jié)果存在內(nèi)存或RocksDB中實(shí)現(xiàn)跨事件的聚合計(jì)算比如統(tǒng)計(jì)每個(gè)用戶的累計(jì)訪問(wèn)次數(shù)這是純SQL流處理引擎很難做好的。當(dāng)然Spark Streaming在吞吐量上和微批處理也有自己的優(yōu)勢(shì)但說(shuō)實(shí)話真心追求實(shí)時(shí)性的場(chǎng)景Flink的靈活性和生態(tài)完整度更適合。更何況現(xiàn)在Flink CDC已經(jīng)是數(shù)據(jù)庫(kù)實(shí)時(shí)采集的事實(shí)標(biāo)準(zhǔn)配合Java開(kāi)發(fā)效率很高。2. 數(shù)據(jù)采集層的工程落地三種來(lái)源一套規(guī)范數(shù)據(jù)采集是整個(gè)實(shí)時(shí)鏈路的起點(diǎn)也是臟活累活最多的地方。很多同學(xué)把精力都花在Flink計(jì)算邏輯上結(jié)果數(shù)據(jù)源沒(méi)管好后面計(jì)算、展示全是垃圾進(jìn)垃圾出。這里我按來(lái)源類(lèi)型分開(kāi)講。2.1 日志類(lèi)采集Filebeat Kafka是黃金組合服務(wù)端日志是最常見(jiàn)的實(shí)時(shí)數(shù)據(jù)來(lái)源。我通常用Filebeat做日志采集器它比Flume輕量太多部署就是解壓一個(gè)二進(jìn)制文件配置也簡(jiǎn)單filebeat.inputs: - type: filestream enabled: true paths: - /data/logs/*.log fields: app_name: order-service log_type: business output.kafka: hosts: [kafka1:9092, kafka2:9092, kafka3:9092] topic: app-order-log partition.round_robin: reachable_only: true這個(gè)配置看起來(lái)簡(jiǎn)單但有幾個(gè)細(xì)節(jié)務(wù)必注意不要用filestream直接用Kafka producer consumer方式Filebeat自帶背壓機(jī)制Kafka不可用時(shí)會(huì)暫停讀取本地文件不會(huì)丟數(shù)據(jù)。這是它作為采集端的核心理由。fields里打上應(yīng)用名和日志類(lèi)型標(biāo)簽后面Flink消費(fèi)時(shí)可以根據(jù)這些字段路由到不同處理邏輯。每個(gè)應(yīng)用單獨(dú)一個(gè)topic或者至少按業(yè)務(wù)線分topic。我曾經(jīng)見(jiàn)過(guò)所有應(yīng)用混在一個(gè)topic里的架構(gòu)Flink消費(fèi)端要做大量過(guò)濾還會(huì)互相影響消費(fèi)速度非常痛苦。2.2 數(shù)據(jù)庫(kù)變更采集Flink CDC到底怎么部署熱搜詞里flink cdc pipeline部署和flink cdc安裝部署出現(xiàn)頻率很高說(shuō)明這個(gè)方向已經(jīng)成了實(shí)時(shí)數(shù)據(jù)平臺(tái)的主流需求。Flink CDC基于數(shù)據(jù)庫(kù)日志Binlog/Redo Log捕獲變更不打業(yè)務(wù)表對(duì)業(yè)務(wù)系統(tǒng)零侵入。部署上有兩種形態(tài)形態(tài)一Flink CDC作為Source接入Flink作業(yè)DataStreamSourceString stream env .addSource( MySqlSource.Stringbuilder() .hostname(localhost) .port(3306) .databaseList(shop) .tableList(shop.t_order) .username(cdc_user) .password(cdc_pwd) .deserializer(new JsonDebeziumDeserializationSchema()) .build() ) .setParallelism(1);這種形態(tài)適合在Flink作業(yè)里實(shí)時(shí)消費(fèi)數(shù)據(jù)庫(kù)變更。注意setParallelism(1)很關(guān)鍵因?yàn)閱蝹€(gè)MySQL實(shí)例的Binlog讀取是單線程的并行度設(shè)置高了反而會(huì)出問(wèn)題。形態(tài)二Flink CDC Pipeline獨(dú)立部署如果你的目標(biāo)是數(shù)據(jù)庫(kù)實(shí)時(shí)同步到另一個(gè)存儲(chǔ)可以用Flink CDC Pipeline也就是之前的CDAS它基于Yaml配置就能完成整庫(kù)同步不需要寫(xiě)一行Java代碼source: type: mysql hostname: localhost port: 3306 username: cdc_user password: cdc_pwd tables: shop\.* sink: type: doris fenodes: doris:8030 username: admin password: admin123Pipeline形態(tài)適合快速落地但是如果你想在同步過(guò)程中做數(shù)據(jù)加工比如字段映射、類(lèi)型轉(zhuǎn)換、過(guò)濾還是寫(xiě)Java代碼更靈活。我的建議是同步裸數(shù)據(jù)用Pipeline需要加工用源碼。關(guān)于Flink CDC最大的坑是存量數(shù)據(jù)與增量數(shù)據(jù)的一致性問(wèn)題。Flink CDC默認(rèn)會(huì)先做一次全量快照再切換到Binlog增量這個(gè)過(guò)程對(duì)數(shù)據(jù)庫(kù)有一定壓力。建議在業(yè)務(wù)低峰期做首次同步并且監(jiān)控好源庫(kù)的IOPS和連接數(shù)。2.3 業(yè)務(wù)主動(dòng)上報(bào)HTTP SDK Kafka注意采樣與限流有些數(shù)據(jù)源既不是日志也不是數(shù)據(jù)庫(kù)而是客戶端行為埋點(diǎn)前端點(diǎn)擊、APP啟動(dòng)等。這時(shí)候通常是業(yè)務(wù)方直接調(diào)用HTTP接口上報(bào)你在接口里把數(shù)據(jù)寫(xiě)入Kafka。這個(gè)環(huán)節(jié)最常見(jiàn)的坑是突發(fā)流量打垮寫(xiě)入服務(wù)。我在某個(gè)項(xiàng)目中遇到過(guò)前端埋點(diǎn)日志突然暴增導(dǎo)致上報(bào)接口被瞬間打滿Kafka客戶端批量發(fā)送超時(shí)丟了一批數(shù)據(jù)。后來(lái)做了三層保護(hù)SDK端批量發(fā)送不要一條一條發(fā)HTTP請(qǐng)求在SDK內(nèi)攢批比如攢夠100條或500ms顯著降低請(qǐng)求頻率服務(wù)端限流單機(jī)QPS上限設(shè)置好超出部分直接丟棄并記錄日志注意埋點(diǎn)數(shù)據(jù)丟幾條通常不影響大屏指標(biāo)趨勢(shì)但要保證不拖垮服務(wù)Kafka端分區(qū)數(shù)規(guī)劃根據(jù)峰值吞吐預(yù)估分區(qū)數(shù)分區(qū)數(shù) 目標(biāo)吞吐量 / 單分區(qū)吞吐量。例如目標(biāo)10萬(wàn)條/秒單分區(qū)吞吐約2萬(wàn)條/秒分區(qū)數(shù)至少5個(gè)數(shù)據(jù)采集層的通用規(guī)范也很重要。所有上報(bào)數(shù)據(jù)統(tǒng)一JSON格式包含event_id全局唯一、event_time事件發(fā)生時(shí)間、source數(shù)據(jù)來(lái)源、biz_body業(yè)務(wù)字段。有了這個(gè)規(guī)范后續(xù)Flink側(cè)做解析、去重、Watermark定義都有據(jù)可依。3. Flink實(shí)時(shí)計(jì)算核心狀態(tài)、時(shí)間語(yǔ)義與Sink的坑到了計(jì)算層就是Flink的主戰(zhàn)場(chǎng)。這里我把最高頻的三個(gè)技術(shù)點(diǎn)拆開(kāi)講這三個(gè)點(diǎn)也是面試和實(shí)戰(zhàn)中最容易翻車(chē)的狀態(tài)管理、時(shí)間語(yǔ)義、自定義Sink。3.1 狀態(tài)與Checkpoint為什么你的作業(yè)重啟丟數(shù)據(jù)Flink的狀態(tài)State是它區(qū)別于普通流處理引擎的核心能力。簡(jiǎn)單理解狀態(tài)就是算到一半的中間結(jié)果。比如你要統(tǒng)計(jì)每分鐘每個(gè)商品的累計(jì)銷(xiāo)售額這個(gè)累計(jì)值就需要保存下來(lái)這就是State。我見(jiàn)過(guò)很多使用者在應(yīng)用里定義了一個(gè)MapState來(lái)保存用戶維度的累計(jì)數(shù)據(jù)然后把Checkpoint間隔設(shè)置成5分鐘。結(jié)果某個(gè)凌晨Flink作業(yè)因?yàn)镺OM掛掉了恢復(fù)后發(fā)現(xiàn)損失了將近10分鐘的統(tǒng)計(jì)結(jié)果。復(fù)盤(pán)時(shí)發(fā)現(xiàn)Checkpoint間隔太大狀態(tài)恢復(fù)點(diǎn)太靠前中間的數(shù)據(jù)全丟了。這里必須記住一個(gè)基本參數(shù)組合state.backend: rocksdb state.checkpoints.dir: hdfs:///flink/checkpoints execution.checkpointing.interval: 60s execution.checkpointing.mode: EXACTLY_ONCE execution.checkpointing.timeout: 5min我的經(jīng)驗(yàn)是線上作業(yè)至少每分鐘做一次Checkpoint太頻繁會(huì)影響性能但5分鐘就太長(zhǎng)了。另外一定要用RocksDB作為狀態(tài)后端——數(shù)據(jù)量一大純內(nèi)存Heap狀態(tài)分分鐘把JVM堆撐爆。RocksDB是把狀態(tài)寫(xiě)到本地磁盤(pán)內(nèi)存只是緩存可靠性和容量都更好。3.2 事件時(shí)間與Watermark亂序數(shù)據(jù)怎么算Flink的窗口計(jì)算有個(gè)經(jīng)典三選一ProcessingTime、EventTime、IngestionTime。做實(shí)時(shí)大屏我強(qiáng)烈建議用EventTime也就是按業(yè)務(wù)事件發(fā)生的時(shí)間來(lái)劃分窗口。但EventTime帶來(lái)的問(wèn)題是數(shù)據(jù)可能亂序到達(dá)。用戶點(diǎn)擊發(fā)生在10:00:00的日志可能到10:00:30才到Flink。如果你正好在做每分鐘點(diǎn)擊量的滾動(dòng)窗口這條數(shù)據(jù)就會(huì)被算到10:01的窗口里指標(biāo)就錯(cuò)了。解決方案是Watermark水位線它表示事件時(shí)間小于等于這個(gè)值的數(shù)據(jù)都已經(jīng)到達(dá)了。DataStreamOrderEvent withWatermark orders .assignTimestampsAndWatermarks( WatermarkStrategy.OrderEventforBoundedOutOfOrderness( Duration.ofSeconds(30) ) .withTimestampAssigner((event, timestamp) - event.getEventTime()) );forBoundedOutOfOrderness(Duration.ofSeconds(30))的意思是容忍最多30秒的亂序。代價(jià)是窗口結(jié)果會(huì)延遲30秒才輸出。這里就是業(yè)務(wù)延遲和數(shù)據(jù)準(zhǔn)確率的權(quán)衡。如果大屏指標(biāo)允許延遲30秒這個(gè)配置就合理如果要求秒級(jí)延遲那就要接受部分亂序數(shù)據(jù)會(huì)算錯(cuò)窗口。窗口計(jì)算上我做實(shí)時(shí)指標(biāo)統(tǒng)計(jì)會(huì)用TumblingEventTimeWindows滾動(dòng)窗口AllowedLateness的組合stream.keyBy(OrderEvent::getProductId) .window(TumblingEventTimeWindows.of(Time.minutes(1))) .allowedLateness(Time.seconds(30)) .aggregate(new CountAggregate()) .process(new WindowResultFunction());allowedLateness的意思是窗口正常計(jì)算后還會(huì)等30秒的遲到數(shù)據(jù)遲到數(shù)據(jù)到達(dá)時(shí)單獨(dú)觸發(fā)一次計(jì)算輸出更新。這樣既保證了主鏈路結(jié)果快速產(chǎn)出又能修正部分亂序數(shù)據(jù)帶來(lái)的誤差。3.3 自定義DataSource與DataSink從入門(mén)到放棄再入門(mén)熱搜詞里flink 自定義 data source和flink 自定義 data sink出現(xiàn)頻率極高我猜是因?yàn)楣俜轿臋n的示例太簡(jiǎn)單一上生產(chǎn)就漏出各種問(wèn)題。這里我把兩個(gè)痛點(diǎn)講透。自定義DataSource通常是為了從非標(biāo)準(zhǔn)源讀數(shù)據(jù)。核心是繼承RichSourceFunction或?qū)崿F(xiàn)SourceFunctionpublic class MetricSource extends RichSourceFunctionMetricEvent { private volatile boolean running true; private transient KafkaProducer producer; Override public void open(Configuration parameters) { producer new KafkaProducer(...); } Override public void run(SourceContextMetricEvent ctx) throws Exception { while (running) { // 模擬讀取外部數(shù)據(jù)源 MetricEvent event readFromExternalSystem(); synchronized (ctx.getCheckpointLock()) { ctx.collect(event); } } } Override public void cancel() { running false; } }注意兩點(diǎn)一是collect操作必須在ctx.getCheckpointLock()鎖內(nèi)執(zhí)行否則Checkpoint時(shí)的狀態(tài)一致性會(huì)出問(wèn)題數(shù)據(jù)可能重復(fù)或丟失。二是cancel()方法里要釋放外部連接資源否則作業(yè)取消時(shí)連接泄漏時(shí)間長(zhǎng)了會(huì)把源系統(tǒng)連接池打滿。自定義DataSink的坑就更多了。我之前寫(xiě)過(guò)自定義Sink寫(xiě)入某個(gè)內(nèi)部監(jiān)控平臺(tái)代碼如下public class MonitorSink extends RichSinkFunctionMetricEvent { private MonitorClient client; Override public void open(Configuration parameters) { client MonitorClient.connect(monitor-server:8080); } Override public void invoke(MetricEvent value, Context context) throws Exception { boolean success client.send(value); if (!success) { throw new RuntimeException(send metric failed: value); } } Override public void close() { client.close(); } }這段代碼看起來(lái)沒(méi)問(wèn)題生產(chǎn)上卻出過(guò)事故監(jiān)控平臺(tái)的單機(jī)處理能力有限Flink端并發(fā)寫(xiě)入量一大client.send就頻繁超時(shí)我讓invoke直接拋異常結(jié)果Flink作業(yè)一直在重啟上游Kafka消費(fèi)被阻滯整個(gè)實(shí)時(shí)鏈路癱瘓。3.4 從事故學(xué)到的Sink設(shè)計(jì)原則那次事故之后我給自己定了幾條Sink設(shè)計(jì)的鐵律也分享給你寫(xiě)外部系統(tǒng)必須做重試和熔斷不能一失敗就拋異常重啟作業(yè)。應(yīng)該捕獲異常做有限次數(shù)重試重試仍失敗就寫(xiě)本地容災(zāi)文件或者發(fā)告警跳過(guò)保證主鏈路不中斷區(qū)分業(yè)務(wù)錯(cuò)誤和系統(tǒng)錯(cuò)誤數(shù)據(jù)格式錯(cuò)誤比如字段缺失屬于業(yè)務(wù)錯(cuò)誤直接throw沒(méi)問(wèn)題因?yàn)橹卦囈蝗f(wàn)次也還是會(huì)失敗外部系統(tǒng)不可用屬于系統(tǒng)錯(cuò)誤應(yīng)該讓作業(yè)保留現(xiàn)場(chǎng)繼續(xù)運(yùn)行等待外部系統(tǒng)恢復(fù)批量寫(xiě)入優(yōu)先于逐條寫(xiě)入能批量就別單條單條寫(xiě)的性能開(kāi)銷(xiāo)太大了。Flink提供了JdbcBatchingOutputFormat支持?jǐn)€批提交但要注意攢批參數(shù)batchSize和batchInterval要配合好另外熱搜詞里flink的jdbc連接器異常是個(gè)高頻問(wèn)題。我遇到過(guò)的大部分情況是連接池耗盡和連接空閑超時(shí)。Flink JDBC Sink的每個(gè)并發(fā)Task都會(huì)建自己的連接你在連接池里配置了最大連接數(shù)10結(jié)果Flink作業(yè)并行度是20直接就有10個(gè)Task拿不到連接報(bào)錯(cuò)。解決辦法很簡(jiǎn)單要么把連接池最大連接數(shù)設(shè)成大于等于Flink并行度要么給連接設(shè)置合理的maxRetryTimes和connectionTimeout。4. 高頻事故復(fù)盤(pán)Flink Sink到Hive表數(shù)據(jù)不落盤(pán)的根因這一節(jié)我要重點(diǎn)復(fù)盤(pán)一個(gè)幾乎每個(gè)做Flink接數(shù)倉(cāng)的人都會(huì)踩的坑——Flink sink Hive表數(shù)據(jù)不入表。這個(gè)熱搜詞出現(xiàn)得如此頻繁說(shuō)明大家都在這上面栽過(guò)跟頭。我把排查鏈路完整還原出來(lái)你以后遇到可以直接照著查。4.1 現(xiàn)象與第一反應(yīng)當(dāng)時(shí)的情況是Flink作業(yè)運(yùn)行狀態(tài)正常沒(méi)有報(bào)錯(cuò)但查詢Hive表時(shí)發(fā)現(xiàn)數(shù)據(jù)一直是空的或者只有很久以前的一部分?jǐn)?shù)據(jù)。我第一反應(yīng)是是不是SQL寫(xiě)錯(cuò)了結(jié)果檢查Flink SQL和Table Schema都對(duì)得上Kafka source也在正常消費(fèi)。于是開(kāi)始逐步排查。4.2 排查鏈路四個(gè)層面逐個(gè)擊破第一層看Flink作業(yè)日志別被正常騙了打開(kāi)TaskManager日志結(jié)果發(fā)現(xiàn)了端倪日志里出現(xiàn)了大量Need to partition the files into Hives format和Abortable相關(guān)的詞。這個(gè)信息很關(guān)鍵——Flink寫(xiě)Hive是按照分區(qū)來(lái)管理的如果你沒(méi)有開(kāi)啟自動(dòng)提交分區(qū)數(shù)據(jù)寫(xiě)入的是臨時(shí)目錄永遠(yuǎn)不會(huì)變成Hive的正式分區(qū)。第二層確認(rèn)Hive表的分區(qū)提交機(jī)制Flink寫(xiě)Hive表默認(rèn)配置涉及兩個(gè)核心參數(shù)。如果你的Hive表是分區(qū)表必須顯式開(kāi)啟分區(qū)提交并且設(shè)置正確的提交觸發(fā)策略CREATE TABLE hive_orders ( order_id BIGINT, product_id BIGINT, amount DOUBLE ) PARTITIONED BY (dt STRING, hour STRING) WITH ( connector hive, sink.partition-commit.trigger partition-time, sink.partition-commit.delay 0s, sink.partition-commit.policy.kind metastore,success-file );sink.partition-commit.trigger如果沒(méi)配或者配成process-time意味著Flink按數(shù)據(jù)到達(dá)時(shí)間來(lái)決定提交分區(qū)不是按數(shù)據(jù)的事件時(shí)間。我當(dāng)時(shí)就是用了process-time結(jié)果業(yè)務(wù)上凌晨的數(shù)據(jù)被算到了早上的分區(qū)等了一早上沒(méi)看到該有的數(shù)據(jù)。第三層檢查寫(xiě)入文件格式與可見(jiàn)性很多剛用Flink寫(xiě)Hive的同學(xué)不知道Flink寫(xiě)Hive默認(rèn)是寫(xiě)ORC或Parquet格式文件到分區(qū)的臨時(shí)目錄然后通過(guò)Table Metastore注冊(cè)分區(qū)。但文件從寫(xiě)入中到可見(jiàn)之間有一個(gè)提交環(huán)節(jié)。如果你看到HDFS上分區(qū)目錄下已經(jīng)有Parquet文件但查詢不到數(shù)據(jù)大概率就是分區(qū)提交沒(méi)有正確執(zhí)行。還有一種可能是你寫(xiě)的是非分區(qū)表Flink寫(xiě)非分區(qū)表會(huì)把數(shù)據(jù)直接寫(xiě)到表的目錄下。但我見(jiàn)過(guò)一個(gè)案例表本身是分區(qū)表Flink SQL里卻只指定了分區(qū)字段的部分值導(dǎo)致Sink端認(rèn)為這是一個(gè)不可寫(xiě)分區(qū)就一直默默丟數(shù)據(jù)。排查方法是用SHOW PARTITIONS hive_orders看分區(qū)元數(shù)據(jù)是否存在。第四層Hive Streaming協(xié)議與Metastore對(duì)接Flink寫(xiě)Hive底層有兩種協(xié)議一種是通用的Hive Streaming API通過(guò)HiveTableSink另一種是直接寫(xiě)文件然后調(diào)用Metastore注冊(cè)分區(qū)。前者需要開(kāi)啟hive.streaming.enabled老版本。如果你用的是較老版本的Flink和Hive建議用hive-streaming-client包并顯式開(kāi)啟dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-hive_2.12/artifactId version你的版本/version /dependency注意Flink和Hive的版本兼容矩陣很挑剔Flink 1.14之前和Hive 3.1.0有兼容問(wèn)題Hive Streaming API在部分版本上是broken的。最穩(wěn)妥的做法是Flink 1.15配Hive 3.1.2同時(shí)把Metastore從嵌入式切換為獨(dú)立部署避免并發(fā)寫(xiě)入時(shí)Metastore鎖沖突。4.3 問(wèn)題根因與修復(fù)最終我們定位到了根因Flink作業(yè)里寫(xiě)Hive表的并行度設(shè)置過(guò)高默認(rèn)等于Kafka分區(qū)數(shù)每個(gè)并發(fā)Task都在嘗試寫(xiě)同一個(gè)分區(qū)的臨時(shí)文件而且Flink內(nèi)部會(huì)出現(xiàn)文件內(nèi)容不完整的競(jìng)態(tài)。配合開(kāi)啟分區(qū)提交后問(wèn)題消失。修復(fù)后的經(jīng)驗(yàn)總結(jié)成一條Flink寫(xiě)Hive表并行度不建議大于1。因?yàn)镠ive表Sink的文件寫(xiě)入不是天然按key分區(qū)的多個(gè)并行度同時(shí)寫(xiě)一個(gè)分區(qū)文件合并和提交的復(fù)雜度會(huì)指數(shù)上升。要么做rebalance并設(shè)置并行度1要么就用bucket功能把數(shù)據(jù)按字段散列到多個(gè)文件讓Flink自己去整理。還有一個(gè)很小的點(diǎn)容易被忽略檢查你的作業(yè)是在本地IDEA跑還是集群跑。本地跑Flink時(shí)HDFS路徑如果寫(xiě)的是hdfs://...會(huì)直接連不上集群如果寫(xiě)的是本地路徑file://...那數(shù)據(jù)其實(shí)是寫(xiě)到你個(gè)人電腦磁盤(pán)上Hive當(dāng)然查不到。這種環(huán)境不一致問(wèn)題我遇到過(guò)不止一次。5. 可視化大屏的實(shí)時(shí)感數(shù)據(jù)刷新頻率、聚合策略與接口設(shè)計(jì)實(shí)時(shí)計(jì)算做完數(shù)據(jù)源源不斷寫(xiě)入結(jié)果表了最后一步是可視化大屏。很多團(tuán)隊(duì)在這里其實(shí)沒(méi)有技術(shù)問(wèn)題但做出來(lái)的大屏看起來(lái)不實(shí)時(shí)——有實(shí)時(shí)數(shù)據(jù)卻沒(méi)有實(shí)時(shí)感。這節(jié)講講大屏之后端數(shù)據(jù)接口設(shè)計(jì)。5.1 大屏的數(shù)據(jù)不要直連數(shù)據(jù)庫(kù)查詢我剛做第一個(gè)實(shí)時(shí)大屏項(xiàng)目的時(shí)候犯過(guò)低級(jí)錯(cuò)誤大屏前端每5秒輪詢一次數(shù)據(jù)庫(kù)原始明細(xì)表SQL里現(xiàn)場(chǎng)做SUM和GROUP BY。這種做法的結(jié)果是數(shù)據(jù)庫(kù)CPU飆升、查詢?cè)絹?lái)越慢、大屏的實(shí)時(shí)變成了每10秒才刷新一次因?yàn)椴樵兒臅r(shí)就占了5秒。正確的思路是打一層結(jié)果表。Flink實(shí)時(shí)計(jì)算出來(lái)的指標(biāo)本來(lái)就已經(jīng)按分鐘/小時(shí)粒度聚合好了直接寫(xiě)入結(jié)果表比如dashboard_metrics大屏端每5秒查詢的就只有幾行聚合好的數(shù)據(jù)查詢耗時(shí)基本在毫秒級(jí)。這才是實(shí)時(shí)大屏該有的性能。5.2 大屏接口的三種刷新模式輪詢模式前端每N秒調(diào)一次后端接口適合指標(biāo)值更新不頻繁、需要簡(jiǎn)單穩(wěn)定的場(chǎng)景。N一般設(shè)為5秒或10秒。WebSocket推送Flink側(cè)結(jié)果更新時(shí)后端主動(dòng)向已連接的大屏客戶端推送數(shù)據(jù)適合大屏數(shù)量多、希望即時(shí)刷新、減少無(wú)效請(qǐng)求的場(chǎng)景。SSE流推送如果你只做單向數(shù)據(jù)推送SSE比WebSocket更簡(jiǎn)單基于HTTP協(xié)議兼容性和調(diào)試成本都低很多。這三個(gè)模式可以混合使用關(guān)鍵指標(biāo)用WebSocket推送次要指標(biāo)用輪詢兜底。前端技術(shù)棧我用得比較多的是Vue ECharts大屏布局用Grid實(shí)現(xiàn)自適應(yīng)。如果你不想花太多時(shí)間調(diào)布局可以直接用現(xiàn)成的DataV或者大屏編輯器但要注意編輯器的數(shù)據(jù)接入?yún)f(xié)議是否支持實(shí)時(shí)推送。5.3 減少大屏刷新壓力聚合結(jié)果表 緩存策略大屏本身有幾十個(gè)圖表如果每個(gè)圖表都單獨(dú)去查一次結(jié)果表也是壓力。我的做法是按業(yè)務(wù)場(chǎng)景把大屏所需的指標(biāo)打包成一個(gè)JSON大接口一次查詢返回所有圖表的數(shù)據(jù)。比如實(shí)時(shí)成交大屏這個(gè)場(chǎng)景接口返回的數(shù)據(jù)結(jié)構(gòu)大致是{ timestamp: 1715673600000, gmv: 102400.5, orderCount: 1287, userCount: 846, trend: [...], rankList: [...], geoDistribution: [...] }大屏端拿到這個(gè)JSON各自渲染對(duì)應(yīng)的圖表組件。這樣一個(gè)接口的查詢時(shí)間通常能控制在50ms以內(nèi)刷新頻率甚至可以提到1秒。還有一點(diǎn)給結(jié)果數(shù)據(jù)加Redis緩存。Flink寫(xiě)入結(jié)果表的同時(shí)把熱數(shù)據(jù)同步一份到Redis大屏接口優(yōu)先讀Redis而不是查Doris或ClickHouse。Redis查詢是純內(nèi)存操作性能遠(yuǎn)高于OLAP數(shù)據(jù)庫(kù)。但要注意最終一致性——如果Flink寫(xiě)入結(jié)果表成功但寫(xiě)Redis失敗緩存里就是舊數(shù)據(jù)。我的解決方式是結(jié)果表帶一個(gè)update_time大屏接口拿數(shù)據(jù)時(shí)會(huì)比對(duì)Redis緩存時(shí)間和本地時(shí)間超過(guò)5秒則強(qiáng)制回源查結(jié)果表。5.4 大屏可視化的實(shí)時(shí)感還有視覺(jué)層面說(shuō)實(shí)話大屏的實(shí)時(shí)感一半靠數(shù)據(jù)一半靠視覺(jué)設(shè)計(jì)。有幾位項(xiàng)目里的前端同學(xué)總結(jié)過(guò)一些經(jīng)驗(yàn)非常有效數(shù)字跳動(dòng)效果關(guān)鍵指標(biāo)成交額、訂單數(shù)用滾動(dòng)數(shù)字代替靜態(tài)數(shù)字視覺(jué)上強(qiáng)化正在變化的感受刷新閃光提示每次刷新成功后給指標(biāo)卡片加一個(gè)淡入的閃爍效果說(shuō)明我更新了時(shí)序圖的時(shí)間軸ECharts的時(shí)間軸坐標(biāo)保持固定寬度數(shù)據(jù)向右側(cè)推進(jìn)給人一種趨勢(shì)正在流動(dòng)的感覺(jué)最后更新時(shí)間顯示大屏角落永遠(yuǎn)顯示數(shù)據(jù)截至 HH:mm:ss讓使用者知道數(shù)據(jù)有多新這些都是細(xì)節(jié)但對(duì)于不懂技術(shù)的領(lǐng)導(dǎo)來(lái)說(shuō)看起來(lái)實(shí)時(shí)和數(shù)據(jù)實(shí)時(shí)同等重要。6. 全鏈路延遲測(cè)量與容災(zāi)沒(méi)有指標(biāo)就沒(méi)有發(fā)言權(quán)實(shí)時(shí)平臺(tái)上線只是開(kāi)始真正難的是讓它穩(wěn)定運(yùn)行、出了問(wèn)題能快速定位。這一節(jié)講講我怎么給實(shí)時(shí)鏈路做體檢和急救。6.1 延遲指標(biāo)每個(gè)環(huán)節(jié)都要有鐘表我構(gòu)建的任何實(shí)時(shí)平臺(tái)都會(huì)在數(shù)據(jù)流里埋一個(gè)端到端延遲衡量機(jī)制。思路很簡(jiǎn)單在數(shù)據(jù)入口打上時(shí)間戳在每個(gè)關(guān)鍵節(jié)點(diǎn)記錄觀察時(shí)間。我在采集端會(huì)在每個(gè)事件的頭部塞一個(gè)ingest_time然后Flink計(jì)算層、存儲(chǔ)層、接口層分別在日志里記下當(dāng)前時(shí)間。通過(guò)一條測(cè)試數(shù)據(jù)就能算出延遲環(huán)節(jié)計(jì)算方式常見(jiàn)瓶頸采集延遲Kafka收到時(shí)間 - 事件發(fā)生時(shí)間日志攢批時(shí)間過(guò)長(zhǎng)、Filebeat端阻塞傳輸延遲Flink收到時(shí)間 - Kafka收到時(shí)間Kafka broker配置、網(wǎng)絡(luò)帶寬計(jì)算延遲Flink輸出時(shí)間 - Flink收到時(shí)間窗口尺寸、狀態(tài)大小、反壓存儲(chǔ)延遲數(shù)據(jù)庫(kù)落庫(kù)時(shí)間 - Flink輸出時(shí)間Sink并行度、批量提交間隔展示延遲大屏收到時(shí)間 - 數(shù)據(jù)庫(kù)返回時(shí)間前端輪詢周期、接口查詢耗時(shí)實(shí)操里我會(huì)寫(xiě)一個(gè)LatencyMonitor的Flink作業(yè)專(zhuān)門(mén)消費(fèi)Kafka的監(jiān)控topic解析每個(gè)事件的ingest_time并計(jì)算延遲分布P50/P95/P99再寫(xiě)入監(jiān)控面板。延遲一旦超過(guò)閾值就觸發(fā)告警。6.2 每個(gè)環(huán)節(jié)的容災(zāi)機(jī)制實(shí)時(shí)鏈路比離線鏈路脆弱得多任何一個(gè)環(huán)節(jié)抖動(dòng)都會(huì)波及到后面。我的容災(zāi)設(shè)計(jì)分三層數(shù)據(jù)源頭采集端必須保證數(shù)據(jù)不丟。Filebeat有本地backlog機(jī)制Flink CDC有Binlog位點(diǎn)記錄Kafka有多副本。這三層可以保證即使整個(gè)實(shí)時(shí)平臺(tái)崩潰數(shù)據(jù)還在源端或Kafka里躺著。Flink作業(yè)打開(kāi)Checkpoint配合RestartStrategy自動(dòng)恢復(fù)。我常用的策略是fixed-delay3次重試間隔10秒。如果3次都失敗就不盲目重啟了發(fā)告警讓人工介入避免無(wú)限重啟導(dǎo)致?tīng)顟B(tài)反復(fù)加載、Kafka消費(fèi)位點(diǎn)反復(fù)跳躍的惡性循環(huán)。存儲(chǔ)與展示結(jié)果表要設(shè)計(jì)冪等寫(xiě)入Flink重啟后重放數(shù)據(jù)不會(huì)產(chǎn)生重復(fù)數(shù)據(jù)。大屏端接口要做降級(jí)——如果結(jié)果表查詢失敗至少返回緩存數(shù)據(jù)或者數(shù)據(jù)暫不可用的明確提示而不是白屏。6.3 數(shù)據(jù)積壓是最大的坑三招止損實(shí)時(shí)鏈路最怕的故障就是數(shù)據(jù)積壓——Kafka里堆積了大量未消費(fèi)的數(shù)據(jù)Flink作業(yè)無(wú)論如何都追不上這時(shí)你看到的大屏是越來(lái)越舊的數(shù)據(jù)實(shí)時(shí)性徹底丟失。數(shù)據(jù)積壓的典型原因和處理方式Flink作業(yè)遇到瓶頸看Flink UI的Backpressure指標(biāo)如果Source端顯示High/Medium說(shuō)明是下游處理不過(guò)來(lái)需要增加并行度或優(yōu)化算子邏輯如果Sink端顯示High說(shuō)明寫(xiě)外部系統(tǒng)慢了需要檢查外部系統(tǒng)的連接池、批量參數(shù)。上游突然峰值流量比如大促秒殺采集量瞬間漲10倍。這時(shí)候Flink集群如果沒(méi)有彈性擴(kuò)縮容只能硬扛。我的建議是Kafka的topic保留時(shí)間設(shè)長(zhǎng)一點(diǎn)7天等峰值過(guò)去后Flink作業(yè)自動(dòng)追趕消費(fèi)。Sink端故障比如ClickHouse或Doris暫時(shí)不可用Flink的Sink會(huì)積壓數(shù)據(jù)在算子內(nèi)部。如果積壓太嚴(yán)重我在Doris不可用期間會(huì)臨時(shí)把結(jié)果寫(xiě)到Kafka的另一個(gè)備份topic等Doris恢復(fù)后重放。數(shù)據(jù)積壓其實(shí)是實(shí)時(shí)平臺(tái)的急性病處理原則是先止損、后排查——先通過(guò)擴(kuò)容或調(diào)整并行度把消費(fèi)速度提上來(lái)再來(lái)分析瓶頸根因。反之如果先停下來(lái)查根因積壓只會(huì)越來(lái)越多雪上加霜。6.4 關(guān)于實(shí)時(shí)平臺(tái)到底需要多實(shí)時(shí)的一些個(gè)人體會(huì)做完整條鏈路我的體會(huì)是實(shí)時(shí)平臺(tái)的技術(shù)難點(diǎn)從來(lái)不是某個(gè)單一組件而是整個(gè)鏈路的平衡工程。很多時(shí)候你不需要追求極致的秒級(jí)延遲只要端到端控制在10秒內(nèi)大屏的體驗(yàn)已經(jīng)相當(dāng)好了。為了那個(gè)極致實(shí)時(shí)你付出的代價(jià)可能是系統(tǒng)復(fù)雜度翻倍、穩(wěn)定性踩坑無(wú)數(shù)。給新手的最實(shí)際建議是先用最簡(jiǎn)單的方案跑通全鏈路——Kafka Flink 結(jié)果表 大屏接口把延遲指標(biāo)測(cè)出來(lái)再針對(duì)瓶頸做優(yōu)化。不要一上來(lái)就上CDC、上高級(jí)狀態(tài)、上復(fù)雜窗口先把骨架立起來(lái)。數(shù)據(jù)和可視化這條路上跑通的那一刻獲得的成就感比任何理論推演都來(lái)得實(shí)在。希望這篇實(shí)戰(zhàn)經(jīng)驗(yàn)貼能幫你少走幾步彎路。