资讯详情 天气爬虫到Kafka Flume HBase大数据实时处理链路实战
📅 2026/10/3 2:45:54
简介围绕天气数据采集到存储的完整链路资源包内整理了一个可运行示例项目适合正在学习大数据实时处理、希望打通数据采集与入库流程的初中级开发者。内容覆盖天气爬虫编写、Kafka消息队列分发、Flume消费并导入HBase之后还通过Hive建立外部映射并用Superset做可视化分析相当于一套小型大数据平台实践参考。包体共17个文件核心是11个Java源码文件对应爬虫、Kafka生产者/消费者、Flume配置等主体逻辑另含1个XML配置、1个Markdown说明文档和3张流程图/效果图可辅助理解整体架构与运行结果压缩包整体约985KB体积小、便于下载后直接阅读。当前已有152人学习适合作为课堂项目或毕业设计的起步模板。读者可借此明确从数据源到存储、分析展示的实现路径获得Kafka与Flume组合使用的关键配置思路以及Hive与HBase映射操作和Superset对接的实际写法相比零散教程包内文件把主干代码、配置和说明集中在一起能显著节省环境搭建和代码调试时间。1. 天氣爬蟲採集Kafka 實時分發Flume 收集數據導入 HBase這條鏈路能落地嗎很多人第一眼看到「天氣爬蟲採集kafka實時分發flume 收集數據導入到 Hbase」這個資源包以為重點是爬蟲代碼。其實真正難的是讓一條從 HTTP 請求發出的天氣數據穩定地流進 Kafka、寫入 HBase、再被 Hive 查出來、最後在 Superset 上畫成圖。這份 weather-mrs-master 項目恰好把整條鏈路串起來了爬蟲抓數據Kafka 做實時中轉Flume 消費 Kafka 並寫入 HBaseHive 與 HBase 建映射Superset 做分析展示。你想在本地復現一個大數據實時處理項目拿它當骨架是比較順手的選擇。適合正在學 Kafka、Flume、HBase 的開發者也適合給簡歷補一個完整項目經驗的從業者。2. 爬蟲採集模塊拆解從 HTTP 請求到 Kafka 生產者先解決數據從哪來2.1 先看項目的 Maven 骨架採集、配置、啟動三塊怎麼分工打開資源包根部是pom.xml、README.md、src/main。這說明整個項目用 Maven 管理主體代碼在 Java 裡不是 Python 腳本拼出來的。我一般會先看pom.xml裡依賴了哪些庫能快速判斷作者意圖有 httpclient 或 okhttp說明採集用 HTTP 客戶端有 jackson-databind說明 JSON 解析走 Jackson有 kafka-clients說明生產者和消費者都走原聲 Kafka API。src/main/java下通常會拆成collector、producer、config、util幾個包resources裡放接口配置、topic 名稱這些可變參數。這個拆法比較樸素但對於一個工程項目來說夠用。天氣數據的來源常見的是各地天氣服務商提供的 JSON 接口比如按城市名查詢實時天氣返回溫度、濕度、風向、天氣現象。有的免費接口每天有調用次數限制有的需要申請 key。項目代碼一般把接口地址和 key 放在配置裡不建議寫死在類裡否則換數據源時要改代碼重新編譯。我在復現時會先在 README 裡確認作者默認用的哪個接口再看resources下的配置文件。如果原項目沒有給具體接口自己換一個能訪問的天氣 API 就行爬蟲模塊的核心邏輯不受影響。選數據源有一個原則優先選返回結構簡單、字段穩定的 JSON 接口不要選那種返回 HTML 還需要正則清洗的頁面。原因是後面鏈路很長只要採集端多一個解析依賴整個數據管道就多一個脆弱點。天氣接口的字段不多溫度、濕度、天氣描述、更新時間足以構成 HBase 一張表的列。如果要採集多個城市還需要一個城市列表配置代碼裡循環調用接口即可。項目的入口類一般就是一個帶main方法的啟動類定時或手動觸發全量採集。2.2 爬蟲抓取代碼HttpClient 拉數據Jackson 解析 JSON採集模塊的核心是一個 HTTP 客戶端加一個 JSON 解析。用 HttpClient 發 GET 請求拿到字符串後用 Jackson 轉成JsonNode再逐字段取值。下面是一段我在類似項目裡常用的寫法讀者可以直接對照資源裡的WeatherFetcher類看。package com.weather.collector; import org.apache.http.client.methods.HttpGet; import org.apache.http.impl.client.CloseableHttpClient; import org.apache.http.impl.client.HttpClients; import org.apache.http.util.EntityUtils; import org.jetbrains.annotations.NotNull; import java.nio.charset.StandardCharsets; import java.util.concurrent.TimeUnit; public class WeatherFetcher { private final CloseableHttpClient httpClient; public WeatherFetcher() { this.httpClient HttpClients.custom() .setConnectionTimeToLive(30, TimeUnit.SECONDS) .setMaxConnPerRoute(20) .build(); } public String fetch(NotNull String apiUrl) throws Exception { HttpGet request new HttpGet(apiUrl); request.setHeader(User-Agent, Mozilla/5.0 (weather-collector)); request.setHeader(Accept, application/json); // 接口通常需要 1~3 秒響應超時設太短會被對方誤判為攻擊 request.setConfig(RequestConfig.custom() .setConnectTimeout(5000) .setSocketTimeout(5000) .build()); try (CloseableHttpResponse response httpClient.execute(request)) { int code response.getStatusLine().getStatusCode(); if (code ! 200) { throw new RuntimeException(HTTP request failed, code code); } return EntityUtils.toString(response.getEntity(), StandardCharsets.UTF_8); } } }這段代碼裡有兩個容易忽略的點一個是CloseableHttpClient要複用不要在每個請求裡new一個實例否則連接池等於沒用高頻採集時會出現大量 TIME_WAIT 連接最終把本地端口耗盡。另一個是超時設置天氣接口在早晚高峰時會抖動5 秒超時比較合適。如果設置成 1 秒稍微一慢就報異常整個循環直接中斷如果設置成 30 秒一個接口掛了會把後續城市全部阻塞。這個「超時設多少」本身就是個坑後面避坑章節會再提。拿到 JSON 字符串後用 Jackson 解析時要養成用path()而不是get()的習慣。get()遇到字段不存在會拋異常path()則返回一個空節點這樣某個字段偶爾缺失不至於讓整個採集程序崩掉。import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import java.util.HashMap; import java.util.Map; public class WeatherParser { private final ObjectMapper mapper new ObjectMapper(); public MapString, String parse(String json, String city) throws Exception { JsonNode root mapper.readTree(json); MapString, String record new HashMap(); record.put(city, city); // path 取值缺失時返回 null不會拋異常 record.put(temperature, root.path(now).path(temp).asText()); record.put(humidity, root.path(now).path(humidity).asText()); record.put(weather, root.path(now).path(text).asText()); record.put(wind, root.path(now).path(windDir).asText()); record.put(updateTime, root.path(updateTime).asText()); return record; } }ObjectMapper也要複用它是線程安全的每次 new 會增加無謂的開銷。parse方法把 JSON 轉成MapString, String這樣後面無論是發 Kafka 還是寫 HBase字段結構都統一了。這裡沒有把數據直接寫到文件或庫裡因為這套項目的設計思路是「採集只負責產出數據」後續交給 Kafka 和 Flume。這樣做的好處是採集模塊變得很純粹換接口、加城市都不會影響下游鏈路。2.3 發送 Kafka 的生產者參數acks、batch.size、linger.ms 背後的取捨採集到數據後下一步是發給 Kafka。天氣數據量不算誇張一個城市幾百字節但如果你採 300 個城市、每分鐘一輪壓力還是有的。生產者的參數能不能調好直接決定 Kafka 會不會成為鏈路瓶頸。下面是一份生產者配置參考的是KafkaProducer的標準寫法。import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.serialization.StringSerializer; import java.util.Properties; public class WeatherProducer { private final KafkaProducerString, String producer; public WeatherProducer(String bootstrapServers) { Properties props new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); // 可靠性參數acks1 表示 leader 寫成功就返回想要更高的可靠性用 all props.put(ProducerConfig.ACKS_CONFIG, 1); props.put(ProducerConfig.RETRIES_CONFIG, 3); // 吞吐參數積累 16KB 或等待 5ms 再批量發送 props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384); props.put(ProducerConfig.LINGER_MS_CONFIG, 5); props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432L); producer new KafkaProducer(props); } public void send(String topic, String city, String json) { producer.send(new ProducerRecord(topic, city, json), (metadata, e) - { if (e ! null) { // 異步回調裡必須留日誌不然發送失敗完全無感知 System.err.println(send failed, topic topic , city city, e); } }); } }參數要理解著調不要照抄。acks1是性能和可靠性的折中Kafka 的 leader 把消息寫進日誌就返回給生產者不需要等所有副本。如果數據不能丟改成all但延遲會增加。retries設 3 是防止網絡抖動但重試可能導致消息亂序所以最好配合enable.idempotencetrue使用這裡沒開因為天氣數據對順序要求不高。batch.size和linger.ms是成對出現的前者表示批次大小後者表示等待時間目的都是攢一批再發。對於這種小字段 JSON5ms 延遲完全可接受。這裡要提一個常見誤區很多人把發送動作寫成同步阻塞producer.send(...).get()這樣每請求一次就等一次網絡往返採集 300 個城市會明顯變慢。正確做法是用異步回調讓生產者內部自己批量發送。上面的send方法就是異步的回調裡只打日誌。對天氣採集這種實時性要求不極端的場景異步加 5ms 積累吞吐完全夠用。項目裡如果把生產者封裝成單例記得在程序退出時調用producer.close()否則緩衝區裡的數據可能沒發完就被強制結束這也算一個低級但常見的翻車點。3. Kafka 實時分發鏈路主題、分區與消費者 offset 的設計3.1 Kafka 集群與主題創建分區數和副本數怎麼定數據從生產者出來第一站是 Kafka topic。資源包裡一般不會帶集群環境需要自己有一套 Kafka 可用。這一步我建議先在單機或三節點虛擬機上把 Kafka 起起來然後用命令行工具創建主題。下面是最基礎的命令。kafka-topics.sh --create \ --bootstrap-server node01:9092,node02:9092,node03:9092 \ --topic weather-topic \ --partitions 6 \ --replication-factor 1 \ --config retention.ms86400000分區數設 6是考慮到後面 Flume 可能開多個 agent 併發消費。如果只開一個 Flume agent分區設 3 或 6 差別不大如果打算橫向擴展多個 Consumer分區數就是並發上限。副本數在生產環境至少 2 或 3本地單機只能設 1設大了會報Replication factor: 2 larger than available brokers: 1這是新手最容易撞到的報錯。retention.ms86400000表示保留一天天氣數據的時效性很強沒必要留七天保留時間太長反而佔磁盤。創建完主題後可以用kafka-topics.sh --describe檢查分區和副本狀態。確認Isr沒有異常Leader 沒有傾斜再繼續往下走。很多時候後面的 Flume 消費不到數據不是 Flume 配錯而是 topic 本身沒建好或者生產者發到了另一個集群。先用命令把叢集狀況摸清能省掉大量排查時間。3.2 生產者與消費者的配置對照別讓數據卡在上游Kafka 的吞吐瓶頸往往不在 broker而在生產者和消費者的配置不對。生產者要看batch.size、linger.ms、buffer.memory消費者要看fetch.min.bytes、max.poll.records、session.timeout.ms。我把這個項目裡最值得調整的幾個參數列成對照方便配置時參考。環節參數推薦值作用生產者acks1 或 all控制消息確認級別天氣數據可容忍少量重試生產者linger.ms5等待 5ms 攢批減少請求次數生產者batch.size16384批次大小字節為單位生產者buffer.memory33554432生產者緩衝區總大小爆了會阻塞 send消費者fetch.min.bytes1消費者拉取的最小數據量調大可減少請求消費者max.poll.records500單次 poll 最大記錄數調太大容易觸發 rebalance消費者session.timeout.ms10000消費者失聯判定時間太短容易誤判消費者enable.auto.commitfalse關閉自動提交手動確認處理完再提交 offset這個表裡最關鍵的是enable.auto.commit。很多項目默認開自動提交處理邏輯還沒跑完就提交了 offset進程一掛就丟數據。Flume 消費 Kafka 時也有對應配置後面會講。生產者的buffer.memory容易被忽略如果發送速度遠大於 Kafka 處理速度緩衝區撐滿後send會阻塞表現就是採集線程卡住不動。解決方法不是無限調大 buffer而是找上游生產速率或下游消費瓶頸。3.3 group.id 和 offset 策略重啟不丟數的基本功Kafka 消費者必須有 group.id同一個 group 裡的消費者在邏輯上共享分區。Flume 作為 Kafka Source 消費時如果不指定 group.id默認會生成一個隨機 group重啟後 resume 從最新的位置開始讀導致漏數據。所以要養成固定 group.id 的習慣比如flume-weather-group。檢查消費位置有兩個常用命令。一個看 consumer group 當前的 lag一個手動模擬消費驗證數據格式。kafka-consumer-groups.sh --bootstrap-server node01:9092 \ --group flume-weather-group --describe kafka-console-consumer.sh --bootstrap-server node01:9092 \ --topic weather-topic --from-beginning --max-messages 5--describe會輸出每個分區的current-offset、log-end-offset、lag三列。如果 lag 一直漲說明消費者處理不過來如果 lag 為 0說明消費速度跟得上生產速度。這個命令是鏈路調試的必備工具。kafka-console-consumer驗證數據格式時記得加--max-messages不然在生產環境會一直刷屏。如果這裡能看到爬蟲採集到的 JSON説明 Kafka 這一段沒問題問題大概率出在 Flume 消費端。offset 策略裡還有一個參數auto.offset.reset。新 group 第一次消費一個沒有 offset 的 topic 時earliest表示從頭開始latest表示從最新開始。對天氣數據項目我建議用earliest這樣 Flume 重啟或建新 group 時能把積壓數據補上。如果業務允許丟數據再考慮latest。這個選擇本身要寫清楚不然重啟一次就會發現丟了一堆「看起來無所謂」的數據。4. Flume 收集數據導入 HBase從 Kafka Source 到 HBase Sink 的完整配置4.1 為什麼選 Flume 做管道它和 Kafka Connect 的差別鏈路裡選 Flume 是有道理的。Flume 本來就是為日誌而生它的 agent 可以靈活定義 source、channel、sink尤其適合「消費 Kafka 再寫 HBase」這種場景。你可能會問直接用 Kafka Connect 的 HBase Sink 不就行了Kafka Connect 也是一條路但項目裡既然用 Flume說明作者要的是輕量、可控、易改。Flume 配置是文本文件改完重啟 agent 就生效不需要寫 Java 代碼而且 Flume 對後面接 HBase 的 serializer 擴展更直接寫自定義類也方便。Flume 的三段式結構要理解source 從哪拿數據channel 暫存數據sink 把數據送出去。這個項目裡 source 是 KafkaSourcechannel 可以選 memory 或 filesink 是 HBaseSink。用 memory channel 吞吐高但 agent 進程掛了 channel 裡的數據會丟用 file channel 會持久化但性能稍差。天氣數據重啟丟幾條影響不大不過為了穩定我一般用 file channel。下面配置會把兩者都寫出來你自己取捨。4.2 flume.conf 實操配置Kafka Source File Channel HBaseSinkFlume agent 的配置最怕就是屬性名拼錯。KafkaSource 的參數前綴和普通 Kafka 消費者略有不同很多人在kafka.bootstrap.servers和bootstrap.servers之間反覆試錯。下面是一份可以直接套用的配置把關鍵參數標出來了。agent.sources kafka-source agent.channels file-channel agent.sinks hbase-sink # Kafka Source消費 weather-topic agent.sources.kafka-source.type org.apache.flume.source.kafka.KafkaSource agent.sources.kafka-source.kafka.bootstrap.servers node01:9092,node02:9092 agent.sources.kafka-source.kafka.topics weather-topic agent.sources.kafka-source.kafka.consumer.group.id flume-weather-group agent.sources.kafka-source.kafka.consumer.auto.offset.reset earliest agent.sources.kafka-source.batchSize 100 agent.sources.kafka-source.batchDurationMillis 2000 # File Channel防止 agent 重啟丟數據 agent.channels.file-channel.type file agent.channels.file-channel.checkpointDir /data/flume/checkpoint agent.channels.file-channel.dataDirs /data/flume/data agent.channels.file-channel.capacity 100000 agent.channels.file-channel.transactionCapacity 1000 # HBase Sink寫入 weather 表 agent.sinks.hbase-sink.type org.apache.flume.sink.hbase.HBaseSink agent.sinks.hbase-sink.table weather agent.sinks.hbase-sink.columnFamily cf agent.sinks.hbase-sink.serializer org.apache.flume.sink.hbase.RegexHbaseEventSerializer agent.sinks.hbase-sink.serializer.regex (.*) agent.sinks.hbase-sink.batchSize 200 # 串起來 agent.sources.kafka-source.channels file-channel agent.sinks.hbase-sink.channel file-channel這份配置裡有幾個值得注意的地方。KafkaSource 的kafka.consumer.group.id決定了消費位置一定要和前面第 3 章說的保持一致不要讓 Flume 每次重啟都換 group。batchSize是 Flume 從 Kafka 一次拉多少條batchDurationMillis是等待多長時間兩個配合起來控制消費節奏。如果設太大KafkaSource 一次 poll 會阻塞比較久HBase 端如果寫得慢數據會堆在 channel 裡。file channel 的capacity是能緩存多少事件我設 100000因為 Kafka 裡如果突然湧入大量積壓數據channel 太小會把 agent 直接卡死。HBase sink 的serializer控制每條事件怎麼變成 HBase 的 putRegexHbaseEventSerializer是按正則解析 event如果你在生產者端發的是普通 JSON 字符串最好在項目裡自定義一個 serializer把 JSON 字段拆成 HBase 的列。資源包的src/main/java裡寫的 serializer 類就是幹這個用的。配置寫完後啟動命令是flume-ng agent -n agent -c conf -f flume.conf。這裡的-n agent必須和配置文件裡agent.sources的前綴一致否則 Flume 會報「component not defined」。這是一個非常容易踩的初級坑。啟動後重點看日誌裡有沒有出現KafkaSource成功訂閱、HBaseSink成功創建 table 這類信息沒有的話就要用下一節的 HBase 命令驗證數據到底寫沒寫進去。4.3 HBase 表設計與預分區RowKey 決定寫入速度HBase 表設計決定了這條鏈路能不能穩定寫入。天氣表不需要複雜的列族一個cf就夠。但 RowKey 一定要設計否則所有數據都往一個 region 寫會出現熱點寫入延遲就會飆升。常見做法是用「城市 時間戳」或「城市反轉 時間戳」。如果用city timestamp做 rowkey同一城市的數據連續寫到一個 region查詢方便但寫入往單點壓。用加鹽或哈希前綴可以把寫入分散到多個 region。下面是一個帶預分區的建表語句。hbase shell EOF create weather, {NAME cf, VERSIONS 1, TTL 86400}, {SPLITS [a, g, m, s]} EOFSPLITS按 rowkey 首字母分了 5 個區間寫入時只要 rowkey 首字母分佈均勻請求就會分散到不同 region。如果數據量不大可以不預分區但一旦 Flume 批量寫入單 region 的壓力會立刻暴露。TTL 設 86400 秒也就是一天天氣歷史數據如果沒有分析價值沒必要長期留著。VERSIONS 1表示每個 cell 只保留一個版本避免 HBase 默認存多版本浪費空間。建完表後可以用count weather或scan weather, {LIMIT 10}驗證數據。scan能看到 rowkey 和列值如果 scan 出來是空的但 Kafka 裡有數據說明 Flume 的 serializer 沒匹配上或者 HBase 表名、列族寫錯了。這類問題的排查方法下一章會單獨展開。總之Flume 到 HBase 這一段能不能通最終看scan結果不要只看 Flume 進程還在跑就以為成功。5. 排查與避坑Hive 映射 HBase 與 Superset 展示的常見問題5.1 Hive 建 HBase 映射表的正確姿勢數據落到 HBase 之後直接查 HBase Shell 不方便做聚合分析所以要用 Hive 把 HBase 表映射成一張外部表。Hive 的HBaseStorageHandler能讓 Hive 查 HBase 裡的數據但映射關係必須完全對準。下面是一條實際可執行的建表語句。CREATE EXTERNAL TABLE weather_hbase ( city STRING, temperature DOUBLE, humidity INT, weather STRING, wind STRING ) STORED BY org.apache.hadoop.hive.hbase.HBaseStorageHandler WITH SERDEPROPERTIES ( hbase.columns.mapping :key, cf:temperature, cf:humidity, cf:weather, cf:wind ) TBLPROPERTIES ( hbase.table.name weather );hbase.columns.mapping必須從:key開始表示 HBase 的 rowkey 對應 Hive 的第一個字段city。後面每一列都要寫完整格式列族:列名順序要和 Hive 字段順序一致。這是我見過最容易翻車的地方HBase 表的列名是cf:temperatureHive 映射裡寫成cf:temp查詢結果就是一片 null。另外Hive 的temperature字段最好和 HBase 裡存的字符串能互轉如果生產者發的是MapString, Stringtemperature是字符串 23.5Hive 映射 DOUBLE 時會做轉換但遇到空字符串就會轉失敗返回 null。建議在爬蟲端就把空值統一處理掉不要讓髒數據流到 HBase。建表成功後執行一個最簡單的查詢驗證映射是不是通的比如SELECT city, temperature FROM weather_hbase LIMIT 5。如果結果有值再往下做 Superset 接入如果查出來全 null先回 HBase Shell 用scan確認數據真實存在再用DESCRIBE weather_hbase檢查字段順序。很多問題其實不是 Hive 語法錯而是 HBase 裡的列名和映射對不上。5.2 Superset 接入 Hive 數據源驅動和授權是兩個坑Superset 分析展示是最後一環也是人們最期待的環節但接入 Hive 數據源有兩個常見攔路虎驅動沒裝好用戶認證失敗。Superset 通過 PyHive 連接 HiveServer2SQLAlchemy 連接串一般寫成hive://superset_user:passwordhive-server:10000/default。如果 Superset 所在機器沒有安裝 PyHive 依賴點「Test Connection」會直接報No module named pyhive。這不是 Superset 配置問題而是 Python 環境缺包需要在啟動 Superset 的虛擬環境裡手動裝pyhive[sasl]和thrift。我在部署時就在這個坑上卡了半天一度懷疑是 Hive 端口沒開。認證方面Hive 如果開了 LDAP 或 KerberosSuperset 連接串就得帶上對應認證方式。簡單環境下直接寫hive://userhost:10000/default?authNOSASL比較省事。這裡要提醒不要把 Hive 用戶密碼寫在 Dashboard 的公開配置裡至少用環境變量引用。Superset 接入後數據集列表裡能看到weather_hbase這張表说明 Hive 層面已經通。後續做圖表時如果某個維度顯示不出來先回到 SQL Lab 跑一句SELECT *看原始數據長什麼樣再判斷是 Superset 的問題還是數據本身的問題。5.3 五個高頻踩坑現象、原因、解決這條鏈路從爬蟲到展示至少有五個坑是反覆出現的。每一條我都按「現象 → 原因 → 解決」拆開寫復現的時候直接對號入座。坑一Flume 寫 HBase 很慢agent 日誌出現大量重試現象是 Flume 日誌裡有Blocking waiting for space in channel或 HBase 寫入超時數據一直往 Kafka 消費但 HBase 表行數不漲。原因通常是 HBase 表沒有預分區所有寫請求打到同一個 region同時 Flume sink 的batchSize設置太大單次提交的 put 太多觸發服務端限流。解決方法按 rowkey 做預分區把batchSize調到 100~200並檢查 HBase 的 region server 日誌裡是否有 GC 停頓。如果還慢看是不是 rowkey 前綴太集中加鹽後重新寫。坑二Hive 映射表查出來全是 null現象是SELECT * FROM weather_hbase返回很多行但每一列都是 NULL。原因很直接hbase.columns.mapping裡的列名、順序和 HBase 表實際列不一致。解決方法先用scan weather, {LIMIT 1}看真實列名再和建表語句逐字段比對。特別注意 HBase 裡的列名是cf:temperature不是temperatureHive 映射裡漏掉cf:前綴也會靜默失敗。坑三Superset 看板日期維度亂碼或查不到現象是時間字段顯示成「2025-01-01T00:00:00.000Z」按小時篩選不對。原因是爬蟲端把時間存成了 ISO 字符串Superset 默認按 UTC 解析而 Hive 裡沒有明確轉時區。解決方法在 Hive 表裡把時間字段轉成TIMESTAMP或者 Superset 的數據集字段設置裡把格式化選成對應時區。這不是 bug是時間類型沒統一最好在爬蟲端就存成yyyy-MM-dd HH:mm:ss。坑四Kafka 重啟後 Flume 重複消費大量數據現象是 Kafka 重啟後HBase 裡出現一批重複 rowkey。原因通常是 Flume 的kafka.consumer.group.id沒固定重啟後新 group 從earliest開始重新消費。解決方法固定 group.id建議設置agent.sources.kafka-source.kafka.session.timeout.ms大一點讓 Flume 在 Kafka 重啟期間不被踢出 group。如果想要丟數據而不是重複數據可以把auto.offset.reset改成latest但這要看業務接受哪種不一致。坑五爬蟲接口偶爾返回 403 或空數據現象是 Kafka 裡某一時段數據明顯變少HBase 的 count 也對應下降。原因是被採集端對高頻請求觸發了風控或者接口超時後拋異常導致整個城市循環中斷。解決方法在爬蟲端加User-Agent輪換、請求間隔控制同時把單個城市的抓取異常 catch 住不能讓一個城市失敗影響後面城市。超時時間也要分開設置SocketTimeout 和 ConnectTimeout 都設成 5 秒左右最穩。這個坑看起來不起眼但會直接造成鏈路「看似正常數據缺一段」的假象。6. 全鏈路驗證用一分鐘冒煙測試檢查採集到展示是否真通6.1 冒煙測試腳本每次改完配置我都會強制自己走一遍冒煙測試而不是直接看 Superset 圖表。先確認 Kafka topic 裡有數據再確認 HBase 能查到數據最後在 Superset SQL Lab 執行查詢。三步都過才說明鏈路通了。# 第一步從 Kafka 消費最近 3 條數據 kafka-console-consumer.sh --bootstrap-server node01:9092 \ --topic weather-topic --max-messages 3 # 第二步查 HBase 最新寫入的行 hbase shell -e scan weather, {LIMIT 3, REVERSED true} # 第三步Hive 裡驗證映射表 hive -e SELECT city, temperature, humidity FROM weather_hbase LIMIT 3;6.2 看板製作技巧與日常檢查Superset 裡新建圖表時我習慣先用「表」這種最樸素的圖形驗證數據確認行數、字段、時間範圍都對再做柱狀圖或折線圖。天氣數據最適合做的圖是「各城市實時溫度柱狀圖」和「單個城市 24 小時溫度趨勢線」。SQL Lab 裡可以直接寫SELECT city, MAX(temperature) FROM weather_hbase GROUP BY city這類語句Superset 支持結果可視化調試效率很高。日常檢查則直接看 Kafka consumer group 的 lag如果 lag 長時間居高不下多半是 Flume 或 HBase 寫入出問題而不是爬蟲停了。6.3 一套可複用的檢查習慣這套項目我前後跑過好幾遍最深刻的教訓是不要在 Flume 日誌還沒看明白的時候就去調 Superset。有一次我以為 Hive 映射表建錯了反覆改STORED BY最後發現是 Flume 的 serializer 把 JSON 當成了整個字符串存到一列裡Hive 映射的字段根本對不上。從那以後我每次接到類似鏈路都會先按「Kafka 有數 → HBase 有數 → Hive 有數 → Superset 有數」的順序強制走一遍冒煙測試哪一步斷了就在哪一步查日誌不跨層猜。希望幫到你。本文还有配套的精品资源点击获取