ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

气象数据实时管道:爬虫+Kafka+Flume+HBase实战链路

气象数据实时管道:爬虫+Kafka+Flume+HBase实战链路 简介本资源是一套面向大数据开发初学者与进阶实践者的实时数据处理系统完整实现聚焦天气数据采集、传输与存储全流程解决从网络爬虫到NoSQL数据库落地的典型工程问题。压缩包共17个文件含11个Java核心代码涵盖爬虫逻辑、Kafka生产者/消费者、Flume拦截器及HBase写入模块、3张架构流程图PNG、1个pom.xml依赖配置、1个README.md项目说明及1个.gitignore整体985KB轻量易部署。内容预览显示系统延伸支持Hive与HBase映射及Superset可视化分析体现端到端数据链路设计。已有152人学习下载读者可直接复用爬虫模板、Kafka主题配置、Flume通道定义及HBase表建模脚本并通过结构化目录快速定位各组件集成要点是理解大数据实时管道构建的高实操性参考样本。1. 天气爬虫采集 Kafka 实时分发 Flume 收集导入 HBase为什么这套链路在气象数据中台里不是“炫技”而是刚需你手头有一批城市级实时天气数据——温度、湿度、风速、PM2.5、紫外线指数每分钟更新一次来源是多个公开气象 API如中国气象数据网、和风天气、OpenWeatherMap但它们格式不一、频率不稳、响应超时频发。你想把这些数据“稳住、存住、查得快”而不是每次写个临时脚本跑完就丢。这时候“天气爬虫采集kafka实时分发flume 收集数据导入到 Hbase.zip”这个标题不是一份打包下载的玩具工程而是一套经过生产验证的数据管道骨架它用爬虫做源头稳压器带重试、限流、字段归一化用 Kafka 做流量缓冲与解耦中枢扛住突发峰值、支持多消费者复用用 Flume 做可靠搬运工自动容错、断点续传、字段映射最终落地 HBase——不是因为“HBase 很酷”而是因为它能以毫秒级响应支撑“查某城市过去72小时每10分钟的温度序列”这类典型 OLAP 查询。这套组合不适用于小样本离线分析但对需要持续摄入低延迟查询高写入吞吐的气象监测、IoT 设备告警、城市运行体征平台是当前最轻量、最可控、最容易横向扩展的落地路径。本文只讲怎么从零搭通这条链路不讲概念对比不画架构图只给你能chmod x就跑、tail -f就看到数据进 HBase 的实操。2. 天气爬虫采集不是 requests.get() 一把梭而是带心跳、字段校验、本地缓存的稳态采集器2.1 为什么不能直接用 cron requests 写死 URL真实气象 API 有三类“反爬”机制① 请求头校验User-Agent、Referer 缺一不可② 频控单 IP 每分钟最多 60 次超限返回 429③ 数据签名部分接口需时间戳密钥拼接 MD5。如果只用requests.get()硬刷10 分钟后就会被封且无法区分“网络超时”和“API 限流”导致数据断层。我们采用Requests Retry Backoff Local Cache三层防御# weather_crawler.py import requests import time import json import sqlite3 from urllib.parse import urlencode from tenacity import retry, stop_after_attempt, wait_exponential, retry_if_exception_type # SQLite 本地缓存避免重复请求同一时间点数据防抖 conn sqlite3.connect(weather_cache.db) conn.execute(CREATE TABLE IF NOT EXISTS cache ( city_code TEXT, timestamp INTEGER PRIMARY KEY, data TEXT NOT NULL, updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP )) retry( stopstop_after_attempt(3), waitwait_exponential(multiplier1, min2, max10), retryretry_if_exception_type((requests.exceptions.Timeout, requests.exceptions.ConnectionError)) ) def fetch_weather(city_code: str) - dict: # 构造带签名的 URL以和风天气为例 params { key: your_api_key, location: city_code, language: zh, unit: c } url fhttps://devapi.qweather.com/v7/weather/now?{urlencode(params)} headers { User-Agent: Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36, Accept: application/json } resp requests.get(url, headersheaders, timeout10) resp.raise_for_status() data resp.json() # 字段归一化统一输出结构屏蔽 API 差异 normalized { city_code: city_code, timestamp: int(time.time()), temperature: float(data.get(now, {}).get(temp, 0)), humidity: int(data.get(now, {}).get(humidity, 0)), wind_speed: float(data.get(now, {}).get(windSpeed, 0)), weather_text: data.get(now, {}).get(textDay, unknown), pm25: int(data.get(now, {}).get(pm25, 0)) } # 写入本地缓存防重复、防断电丢失 conn.execute( INSERT OR REPLACE INTO cache (city_code, timestamp, data) VALUES (?, ?, ?), (city_code, normalized[timestamp], json.dumps(normalized)) ) conn.commit() return normalized if __name__ __main__: cities [101010100, 101020100, 101280101] # 北京、上海、深圳编码 for city in cities: try: result fetch_weather(city) print(f[OK] {city} - {result[temperature]}°C, {result[weather_text]}) # 发送到 Kafka下一章 except Exception as e: print(f[FAIL] {city} - {e})提示tenacity是 Python 最成熟的重试库wait_exponential指数退避比固定间隔更抗突发抖动SQLite 缓存不是可选而是必须——当 Kafka 或 Flume 临时故障时爬虫仍能持续写入本地等下游恢复后批量补发这是整条链路“不丢数据”的第一道保险。2.2 爬虫输出到 Kafka为什么不用文件落地再读取很多教程让爬虫先写 CSV/JSON 文件再用 Flume tail -f 监听。这在测试阶段可行但生产环境会出三个问题① 文件轮转时 Flume 可能漏读最后一行② 多进程爬虫并发写同一文件引发锁冲突③ 文件系统 I/O 成为瓶颈尤其当每秒写入 500 条。正确做法是爬虫直连 Kafka Producer# 续上 weather_crawler.py from kafka import KafkaProducer import json producer KafkaProducer( bootstrap_servers[localhost:9092], value_serializerlambda v: json.dumps(v).encode(utf-8), # 关键参数确保消息不丢失 acksall, # 所有副本确认才返回成功 retries5, # 网络失败自动重试 max_in_flight_requests_per_connection1, # 防止乱序HBase 写入依赖顺序 linger_ms10, # 批量攒 10ms 提升吞吐 buffer_memory33554432 # 32MB 缓冲区防瞬时高峰溢出 ) def send_to_kafka(data: dict): try: producer.send(weather-raw, valuedata, keydata[city_code].encode()) producer.flush() # 强制发送避免缓冲区积压 except Exception as e: print(f[KAFKA SEND FAIL] {e}) # 在主循环中调用 result fetch_weather(city) send_to_kafka(result) # 直接发不落盘参数说明acksall是 Kafka 生产者端数据可靠性基石max_in_flight_requests_per_connection1虽牺牲少量吞吐但保证分区内的消息严格有序——这对后续 Flume 按 key 分组写入 HBase 表至关重要linger_ms10是吞吐与延迟的平衡点实测 10ms 下单 Producer 每秒稳定写入 1200 条远超气象数据需求通常 100 条/秒。3. Kafka 实时分发不是装完就完事而是要验证分区、压缩、监控三件套3.1 创建 topic 的最小必要命令分区数与副本因子怎么定气象数据写入特点是写多读少、按城市维度查询、数据生命周期明确保留 90 天。因此 topic 设计必须匹配# 创建 weather-raw topic假设 3 节点 Kafka 集群 kafka-topics.sh --create \ --bootstrap-server localhost:9092 \ --replication-factor 3 \ --partitions 12 \ --topic weather-raw \ --config retention.ms7776000000 \ # 90 天 90*24*60*60*1000 ms --config segment.bytes1073741824 \ # 1GB 段大小减少小文件 --config compression.typelz4 # LZ4 压缩CPU 开销小压缩率够用为什么是 12 分区分区数 ≥ 消费者实例数Flume agent 数否则有消费者空转分区数 ≤ 单节点磁盘 IOPS 承载能力SSD 一般 10K IOPS12 分区 ≈ 800 IOPS/分区安全按城市哈希 key 分区12 是常见城市编码如 101010100模 12 的结果分布较均匀避免热点分区。3.2 验证 Kafka 是否真在“实时”工作三步快速诊断别信kafka-console-consumer.sh看到消息就认为通了。真实链路中Kafka 是承上启下的黑匣子必须验证三件事Producer 端是否真发出去# 查看 producer 日志中的 send success 计数关键 grep send success /var/log/kafka/server.log | wc -l # 对比爬虫日志里的 send_to_kafka 调用次数二者应基本一致Consumer 端是否真消费# 查看 Flume agent 日志搜索 Event took 字样 grep Event took /var/log/flume/flume.log | tail -20 # 正常应看到类似Event took 12ms to process → 表示 Flume 正在实时拉取Topic 是否有堆积# 获取 lag堆积量 kafka-consumer-groups.sh --bootstrap-server localhost:9092 \ --group flume-hbase-group --describe --all-groups # 输出中 LAG 列应长期为 0若持续增长说明 Flume 处理不过来或 HBase 写入慢注意kafka-console-consumer.sh只能验证“消息存在”不能验证“消息被消费”。真正要看的是 Flume 日志里的处理耗时和 Kafka 消费组 lag这是判断实时性的黄金指标。4. Flume 收集数据导入 HBase不是配置完就完而是要字段映射、RowKey 设计、写入幂等4.1 Flume 配置文件核心段source → channel → sink 的精准对应Flume 的flume.conf不是模板复制粘贴就能用。气象数据链路要求每个城市数据独立写入 HBase 表RowKey 必须含时间戳以支持范围扫描且写入失败不能丢数据。以下是生产级配置基于 Flume 1.11# flume.conf # SOURCE: Kafka a1.sources r1 a1.sources.r1.type org.apache.flume.source.kafka.KafkaSource a1.sources.r1.kafka.bootstrap.servers localhost:9092 a1.sources.r1.kafka.topics weather-raw a1.sources.r1.kafka.consumer.group.id flume-hbase-group a1.sources.r1.batchSize 1000 a1.sources.r1.batchDurationMillis 2000 # CHANNEL: File-backed防 Flume 进程崩溃丢数据 a1.channels c1 a1.channels.c1.type file a1.channels.c1.checkpointDir /var/lib/flume/checkpoint a1.channels.c1.dataDirs /var/lib/flume/data a1.channels.c1.capacity 1000000 a1.channels.c1.transactionCapacity 10000 # SINK: HBase关键自定义 serializer 处理 JSON RowKey 生成 a1.sinks k1 a1.sinks.k1.type hbase a1.sinks.k1.hbase.zookeeper.quorum localhost a1.sinks.k1.hbase.zookeeper.property.clientPort 2181 a1.sinks.k1.hbase.table weather_data a1.sinks.k1.hbase.columnFamily cf a1.sinks.k1.serializer org.apache.flume.sink.hbase.HBaseEventSerializer a1.sinks.k1.serializer.rowKeyColumn timestamp a1.sinks.k1.serializer.rowKeyTimestamp true a1.sinks.k1.serializer.payloadColumn data a1.sinks.k1.serializer.ignoreTimestamp false # BINDING a1.sources.r1.channels c1 a1.sinks.k1.channel c1关键点解析file channel是 Flume 可靠性的命脉checkpointDir和dataDirs必须挂载在 SSD 上且预留足够空间建议 ≥ 50GBserializer.rowKeyColumn timestamp表示从 JSON 中提取timestamp字段作为 RowKeyserializer.rowKeyTimestamp true启用 HBase 自动时间戳避免手动设put.setTimestamp()出错serializer.payloadColumn data表示整条 JSON 存入cf:data列后续用协处理器或 Phoenix 查询时可直接SELECT data.temperature FROM weather_data。4.2 HBase 表设计为什么不用city_code timestamp作 RowKey初学者常把 RowKey 设为BJ_1717023456城市时间戳这会导致严重热点所有北京数据都写入同一 RegionServer。正确方案是加盐salting 时间戳反转# 创建表HBase Shell create weather_data, {NAME cf, TTL 7776000}, {SPLITS [0000,1111,2222,3333,4444,5555,6666,7777,8888,9999]}然后在 Flume 的HBaseEventSerializer基础上自定义一个SaltedRowKeyGeneratorJava 类逻辑如下// SaltedRowKeyGenerator.java编译后放入 flume-ng-hbase-sink.jar public class SaltedRowKeyGenerator implements RowKeyGenerator { private static final int SALT_BUCKETS 10; Override public byte[] generateRowKey(Event event) { String json new String(event.getBody()); JSONObject obj new JSONObject(json); String cityCode obj.optString(city_code, unknown); long ts obj.optLong(timestamp, System.currentTimeMillis()); // 反转时间戳使新数据 RowKey 更大避免 Region Split 后老数据全在首 Region long reversedTs Long.MAX_VALUE - ts; // 加盐city_code % 10 作为前缀 int salt Math.abs(cityCode.hashCode()) % SALT_BUCKETS; String rowKey String.format(%01d_%s_%d, salt, cityCode, reversedTs); return rowKey.getBytes(StandardCharsets.UTF_8); } }效果RowKey 变成3_101010100_922337203685477580710 个 salt 前缀将写入压力均摊到 10 个 RegionreversedTs确保新数据总在表末尾追加避免频繁 Region Split 导致 compaction 压力。5. 避坑Flume HBase 链路中最容易翻车的 4 个血泪现场5.1 现象Flume agent 启动后日志疯狂报Failed to open connection to HBase但hbase shell能连原因Flume 使用的 HBase 客户端版本如 2.4.9与 HBase 服务端版本如 2.6.0不兼容特别是hbase-clientjar 包中ConnectionImplementation类签名变更。解决删除 Flumelib/下所有hbase-*.jar从 HBase 服务端$HBASE_HOME/lib/目录拷贝以下 5 个 jar 到 Flumelib/hbase-client-2.6.0.jar,hbase-common-2.6.0.jar,hbase-protocol-2.6.0.jar,hbase-server-2.6.0.jar,htrace-core4-4.2.0-incubating.jar重启 Flume agent。5.2 现象HBase 表里数据存在但scan weather_data返回空或只返回部分数据原因HBase 默认scan只查最新版本VERSIONS1而 Flume 写入时未显式指定版本HBase 自动分配时间戳可能因服务器时钟不同步导致版本混乱。解决创建表时强制单版本create weather_data, {NAME cf, VERSIONS 1}或在 Flume sink 配置中加a1.sinks.k1.serializer.ignoreTimestamp true让 HBase 用服务器本地时间戳。5.3 现象Flume agent 运行几小时后 OOMjava.lang.OutOfMemoryError: Java heap space原因file channel的checkpointDir和dataDirs未定期清理大量.checkpoint和.data文件堆积尤其当 Kafka 消息堆积时channel 缓存暴涨。解决设置 Flume 启动参数-Xms2g -Xmx4g根据机器内存调整添加定时清理脚本每天凌晨执行# /etc/cron.daily/clean-flume-channel find /var/lib/flume/checkpoint -name *.checkpoint -mtime 7 -delete find /var/lib/flume/data -name *.data -mtime 7 -delete5.4 现象Kafka 消费 lag 持续增长Flume 日志显示Event took 5000ms to process原因HBase 写入慢根源常是 RegionServer GC 频繁或 WAL 写入磁盘慢。排查与解决查hbase-regionserver.log搜索GC pause若单次 GC 1s调 JVM 参数-XX:UseG1GC -XX:MaxGCPauseMillis200 -Xms8g -Xmx8g检查 WAL 目录磁盘df -h /hbase/wal若使用机械盘换 SSD 并配置hbase.wal.dir到 SSD 路径临时降级在 Flume sink 配置中加a1.sinks.k1.hbase.batchSize 100默认 1000减小单次 HBase RPC 压力。6. 验证与进阶用一条 SQL 查清过去 24 小时某城市的温度趋势并给它加上告警阈值6.1 用 Phoenix 快速验证数据质量比 HBase Shell 直观 10 倍Phoenix 是 HBase 的 SQL 层安装后无需改代码直接用标准 SQL 查询。对气象数据这是最高效的验证方式-- 连接 Phoenix假设已配置好 JDBC !connect jdbc:phoenix:localhost:2181 -- 查北京过去 24 小时温度序列RowKey 设计决定此查询高效 SELECT TO_CHAR(TO_DATE(CAST(timestamp AS BIGINT) * 1000), YYYY-MM-DD HH24:MI) AS time_point, temperature, humidity, weather_text FROM weather_data WHERE SUBSTR(rowkey, 1, 1) 3 -- salt prefix AND rowkey LIKE 3_101010100_% -- 北京 city_code AND timestamp (CURRENT_TIME() - INTERVAL 24 HOUR) * 1000 ORDER BY timestamp DESC LIMIT 144; -- 每10分钟1条24小时共144条为什么能快Phoenix 将WHERE rowkey LIKE 3_101010100_%下推为 HBase 的PrefixFilter只扫描匹配前缀的 Region避免全表扫描。这是 RowKey 设计正确的直接收益。6.2 给查询加告警用 Phoenix 视图 UDF 实现“高温预警”Phoenix 支持自定义函数UDF我们可以写一个 Java UDF 判断温度是否超阈值// TempAlertUDF.java public class TempAlertUDF extends BaseUDF { Override public Boolean evaluate(Integer temp) { if (temp null) return false; return temp 35; // 35°C 为高温阈值 } }编译打包为temp-alert-udf.jar上传到 HBase 所有 RegionServer 的hbase/lib/然后在 Phoenix 中注册CREATE FUNCTION temp_alert AS com.example.TempAlertUDF USING JAR /path/to/temp-alert-udf.jar;再执行带告警的查询SELECT time_point, temperature, CASE WHEN temp_alert(temperature) THEN HIGH_TEMP_ALERT ELSE NORMAL END AS alert_level FROM ( SELECT TO_CHAR(TO_DATE(CAST(timestamp AS BIGINT) * 1000), YYYY-MM-DD HH24:MI) AS time_point, CAST(JSON_EXTRACT(data, $.temperature) AS INTEGER) AS temperature FROM weather_data WHERE SUBSTR(rowkey, 1, 1) 3 AND rowkey LIKE 3_101010100_% AND timestamp (CURRENT_TIME() - INTERVAL 24 HOUR) * 1000 ) t ORDER BY time_point DESC;结果示例2024-05-30 14:30 | 36 | HIGH_TEMP_ALERT2024-05-30 14:20 | 34 | NORMAL这就是一条可直接对接 Grafana 或钉钉机器人的告警流水线——没有额外中间件全在 HBase Phoenix 内完成。6.3 最后一条实战技巧用zip封装整个部署包但别让它成为运维黑洞标题里的.zip不是随便加的。我习惯把整套链路打包为weather-pipeline-v1.2.zip结构如下weather-pipeline/ ├── conf/ │ ├── kafka/ # server.properties, zookeeper.properties │ ├── flume/ # flume.conf, log4j.properties │ └── hbase/ # hbase-site.xml, regionservers ├── bin/ │ ├── start-all.sh # 一键启 Kafka/ZK/HBase/Flume │ └── check-health.sh # 检查各组件端口、lag、HBase region 状态 ├── data/ │ └── city-codes.csv # 城市编码映射表供爬虫读取 └── lib/ └── temp-alert-udf.jar # 自定义 UDF关键经验.zip里绝不放二进制如 kafka_2.13-3.6.0.tgz只放配置、脚本、JAR。所有软件用 Ansible 或 Docker 安装.zip只是“配置即代码”的载体。这样升级时只需替换 zip 包并./bin/start-all.sh --force-reload不用碰任何安装目录。我吃过亏曾经把 Kafka 二进制打进去结果某次 unzip 覆盖了旧版集群直接脑裂。现在我的 zip 包解压后ls -l第一眼就能看清全是文本心里才踏实。希望帮到你。本文还有配套的精品资源点击获取
返回列表