ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

基于 Spark2.x 的新闻网大数据实时分析全链路部署与 Flume HBase 性能优化实践

基于 Spark2.x 的新闻网大数据实时分析全链路部署与 Flume HBase 性能优化实践 简介这是一份面向毕业设计场景的新闻网大数据实时分析可视化系统完整项目资源基于Spark2.x实现数据采集、预处理、实时分析与图表展示全链路适合大数据专业学生、课程设计者及入门开发者参考学习已有45人学习下载。压缩包共35个文件约3.43MB内部包含10个Jar依赖、7个Scala与6个Java源码、前端JS/XML配置资源、3张系统效果图以及部署文档和说明文本。其中Jar包提供Flume/HBase等外部依赖Scala与Java代码对应Spark分析及可视化后端逻辑MD/README和TXT部署文档则指导环境搭建目录划分清晰。资源完整覆盖从Flume日志接入到HBase存储、再由Spark实时分析并Web可视化的关键环节源码中包含HBase自定义Sink等核心实现部署文档详细记录了软硬件准备、网络配置等步骤。凭借项目源码、部署文档与全部数据资料可快速复现系统并支撑完成课程设计、毕业设计或论文实验。1. 基于 Spark2.x 的新闻网大数据实时分析这套源码我十分钟把链路跑通做新闻网大数据实时分析最怕的不是算法推导而是 Flume、Kafka、Spark、HBase 四个组件各说各话。我拿到这套基于 Spark2.x 的毕业设计源码时扫了一眼目录就发现代码量不大真正值钱的是那个参考步骤.txt和改好的 Flume 序列化器。把 weblogs 日志从采集端一路送到可视化大屏只花了一个晚上中间最大的坑卡在 HBase 的 RowKey 设计上。这套资源适合正在做大数据综合实训、课程设计或者想完整复现一条实时分析业务流的人。它不是教学 demo是一套能真正跑出图表和指标的 News_Spark 项目。2. 项目拆包文件清单与从日志到可视化的整条数据链路2.1 文件清单哪些文件真正影响启动结果我习惯先把压缩包完整列一遍避免后面拿着 README 找不到入口。这套News_Spark-master解压后根目录下杂项不少但真正在链路里起作用的就这几类。文件 / 目录作用我的判断sparkStu/pom.xmlMaven 构建文件定义 Spark、Kafka、HBase 依赖构建入口依赖版本写得很全sparkStu/srcSpark 作业源码包含数据清洗、统计逻辑核心代码建议优先读flume_hbase/KfkAsyncHbaseEventSerializer.java自定义 Flume HBase Sink 的 Event 序列化器关键改造点解决异步写入flume_hbase/SimpleRowKeyGenerator.javaRowKey 生成器核心热点问题就靠它解决flume_hbase/SimpleHbaseEventSerializer.java官方示例序列化器对比参考实际没采用flume_hbase/flume-ng-hbase-sink.jar打包好的 Flume HBase Sink 扩展 jar直接放到 Flume lib 下即可weblogs/采集的新闻网站日志样本直接喂给 Flume 的测试数据参考步骤.txt部署操作步骤全程最有价值的文档z_pic/news1.png等可视化界面截图判断前端预期效果README.md基础说明写得太简略别只靠它拿到源码先别急着 mvn package。我建议按README → 参考步骤.txt → 源码的顺序读因为 README 只给了项目背景真正的启动命令、表结构、路径配置全在参考步骤.txt里。如果你只照着 README 操作大概率会在环境变量上卡住。2.2 数据链路Flume 采集、Kafka 缓冲、Spark 清洗、HBase 存储、可视化读取整套系统的数据流是一条清晰的管道我拆包后第一件事就是先把链路画出来再对照代码逐个确认。常见做法是 Flume 挂一个 source 监听日志目录把每行日志包装成 Event通过 Kafka Channel 或直接写入 Kafka topicSpark Streaming 从 Kafka 拉数据做窗口统计结果落 HBase可视化层定时扫描 HBase 做展示。环节组件数据形态说明采集Flume原始 weblogs 日志监听目录或 tail -F缓冲Kafka序列化后的 Event削峰填谷隔离生产与消费处理Spark StreamingDStream / RDD按窗口做 PV、UV、来源统计存储HBase含 RowKey 的 KV按时间前缀存储方便扫描展示Web 可视化JSON 数据定时查询 HBase刷新图表weblogs里的样本日志字段比较规整基本是访问时间、IP、URL、状态码、来源页这几类。Spark 端做的清洗主要是切分字段、过滤掉 404 和空 referer再按分钟窗口对请求量、独立 IP 数、热门 URL 做聚合最终把结果写入 HBase 的不同列族。可视化端拿到的已经是聚合结果不再做二次计算所以大屏刷新基本在秒级。这套链路里最容易出问题的不是 Spark 作业本身而是 Flume 端写入 HBase 时的序列化方式。如果不改官方自带的SimpleHbaseEventSerializer你会发现日志全被写进同一个单元格根本没法按时间查询这也是我下一章要重点展开的原因。3. 部署环境准备JDK、ZooKeeper、Kafka、HBase、Spark 的接地气配置3.1 版本匹配的三条原则这套项目基于 Spark2.x但名字只到 2.x具体到小版本得看pom.xml。我先说三条通用的版本匹配原则这是部署任何 Spark 项目都能用的经验。第一Spark 2.x 和 Kafka 客户端的兼容性要看 Scala 版本。如果pom.xml里 Spark 是 2.4 系列通常对应 Scala 2.11 或 2.12Kafka 客户端选kafka-clients而不是kafka_2.11避免把整套 Kafka 服务端依赖带进 Spark 作业。第二HBase 版本尽量选 1.x 或 2.x 里与 Spark 集成文档最全的系列因为 Flume 的 HBase sink 对 HBase 客户端的 API 兼容性比较敏感。第三ZooKeeper 版本跟着 Kafka 和 HBase 走避免一个 3.4 一个 3.6 导致会话超时。我在部署时习惯先把pom.xml里的依赖树打出来看一遍。命令是cd sparkStu mvn dependency:tree deps.txt看deps.txt的目的很简单确认有没有多个版本的 HBase-client 或 Kafka-clients 互相冲突。常见问题是 Flume 自带的 lib 里有一套 HBase-clientSpark 作业的 fat jar 里又带一套运行时就会报NoSuchMethodError。这个现象在后面避坑章节我会单独讲。3.2 四个必须改的配置文件整套系统部署时真正需要手改的配置文件就四个。我按修改顺序列出来并标注了关键参数。文件位置关键项说明flume-env.shFlume conf 目录JAVA_HOME、FLUME_CLASSPATH必须把 HBase 和 Kafka 相关 jar 加进来flume.confFlume conf 目录source、channel、sink 的定义指向 weblogs 目录和 Kafka topichbase-site.xmlHBase conf 目录hbase.zookeeper.quorum客户端连接 HBase 的地址spark-env.shSpark conf 目录SPARK_MASTER_HOST、SPARK_DRIVER_MEMORY单机部署时核心配置Flume 的flume.conf里sink 部分要指定改好的序列化器全类名。我一般这么写 HBase sinkagent.sinks.hbaseSink.type com.example.flume.sink.hbase.AsyncHBaseSink agent.sinks.hbaseSink.table news_weblogs agent.sinks.hbaseSink.columnFamily info agent.sinks.hbaseSink.serializer com.example.flume.sink.hbase.KfkAsyncHbaseEventSerializer agent.sinks.hbaseSink.serializer.rowKeyGenerator com.example.flume.sink.hbase.SimpleRowKeyGenerator这里AsyncHBaseSink是异步写入的关键直接决定吞吐量serializer指向项目自带的改造类rowKeyGenerator必须指定否则会用默认的 Event 头生成 RowKey热点问题立刻暴露。参数table和columnFamily要和 HBase 里预创建的表结构完全一致否则写入时报NoSuchColumnFamilyException。3.3 启动顺序与验证命令部署时组件启动顺序有讲究顺序错了会出现反复重连、会话过期的问题。我通常按这个顺序操作# 1. 启动 ZooKeeper zkServer.sh start # 2. 启动 HDFS如果 HBase 依赖 HDFS start-dfs.sh # 3. 启动 HBase start-hbase.sh # 4. 启动 Kafka kafka-server-start.sh -daemon config/server.properties # 5. 创建 Kafka topic kafka-topics.sh --create --topic news-log --partitions 3 --replication-factor 1 --bootstrap-server localhost:9092 # 6. 启动 Spark 作业在 sparkStu 目录 mvn clean package -DskipTests spark-submit --class com.news.spark.StreamingApp \ --master local[4] \ target/news-spark-1.0.jar # 7. 启动 Flume flume-ng agent --conf conf --conf-file conf/flume.conf --name agent -Dflume.root.loggerINFO,console验证链路是否通先别急着看大屏。我先在 Kafka 里确认消息进来了kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic news-log --from-beginning如果 Kafka 能看到 weblogs 的原始日志说明 Flume 采集成功。接着在 HBase shell 里扫描表hbase shell scan news_weblogs, {LIMIT 5}看到带时间戳前缀的 RowKey 和字段列说明 Flume 到 HBase 写通。最后再确认 Spark Streaming 的窗口统计是否落表这样每个环节都能快速定位问题位置。4. 改造 Flume 的 HBase 序列化器异步批量写入与 RowKey 加盐4.1 自带序列化器的两个缺陷官方自带的SimpleHbaseEventSerializer在生产场景基本不能用这是我对照源码得出的结论也是我推荐你一定要用项目里改好的类的原因。第一个缺陷是同步写入。SimpleHbaseEventSerializer在HBaseSink里是逐条 put每条 Event 都等待 HBase 返回确认吞吐量被锁死在网络往返上。新闻网站日志在高峰时段每秒十几万条同步写根本扛不住Kafka 堆积会越来越严重。第二个缺陷是 RowKey 设计粗暴。官方实现默认用Event的 header 里的 timestamp 做 RowKey或者干脆用随机 UUID。前者会产生顺序递增的 RowKey所有写入落在 HBase 的同一个 region 上立刻形成热点 region后者虽然分散了压力但完全没法按时间范围扫描可视化层查数据只能全表扫。项目里改好的KfkAsyncHbaseEventSerializer同时解决了这两个问题。它的核心思路是使用AsyncHBaseSink通过回调批量提交 put同时把 RowKey 交给独立的SimpleRowKeyGenerator生成兼顾散列和时间范围查询。4.2 改写 KfkAsyncHbaseEventSerializer异步回调与列族映射先看关键代码。我摘取了序列化器里最核心的两段逻辑一段是构造 put 请求一段是异步回调处理。public class KfkAsyncHbaseEventSerializer implements AsyncHbaseEventSerializer { private final ListPutRequest putRequests new ArrayList(); private byte[] table; private byte[] columnFamily; private Event event; Override public void initialize(byte[] table, byte[] columnFamily, Event event) { this.table table; this.columnFamily columnFamily; this.event event; } Override public ListPutRequest getActions() { putRequests.clear(); String body new String(event.getBody(), StandardCharsets.UTF_8); String[] fields body.split(\\|); if (fields.length 4) { return putRequests; } String rowKey SimpleRowKeyGenerator.generateRowKey(fields[0], fields[1]); PutRequest put new PutRequest(table, rowKey.getBytes(), columnFamily, url.getBytes(), fields[2].getBytes()); putRequests.add(put); return putRequests; } Override public void onFailure(Exception e) { // 记日志落本地文件防止丢数据 logger.error(HBase put failed, e); } }逻辑说明initialize负责把 Flume 传入的表名、列族和 Event 缓存下来getActions是每个 Event 处理时都会被调用的方法它先切分日志字段再调用SimpleRowKeyGenerator生成 RowKey最后构造PutRequest放进列表。注意我这里返回的是一个列表因为一条日志可能需要写入多列例如 URL、状态码、响应时长每个列对应一个PutRequest。参数说明body.split(\\|)要求 weblogs 里的日志字段以竖线分隔如果原始日志是空格或 Tab 分隔必须在这里改成对应的正则否则fields.length 4会直接丢掉数据。columnFamily定义成info对应 HBase 表创建时的列族名两者不一致会报错。onFailure里我只打印错误日志并忽略如果你想保证数据不丢可以在失败时把 Event 重新写回 Kafka我一般会加一个 retry 队列这套源码里没写属于可扩展点。4.3 SimpleRowKeyGenerator时间戳加随机前缀避开写热点RowKey 是整个 HBase 设计的灵魂。如果直接拿时间戳当 RowKey 前缀HBase 按字典序存储时同一秒的写入会全部打到同一个 region导致单点过热。项目里的SimpleRowKeyGenerator用了加盐方案我看源码的时候觉得这段值得单独拿出来说。public class SimpleRowKeyGenerator { private static final SimpleDateFormat SDF new SimpleDateFormat(yyyyMMddHHmm); public static String generateRowKey(String timestamp, String ip) { String timePart SDF.format(new Date(Long.parseLong(timestamp))); String salt String.valueOf((ip.hashCode() 0x7fffffff) % 100); return salt _ timePart _ ip; } }逻辑说明timePart把毫秒时间戳格式化到分钟精度作为时间范围扫描的基准salt取 IP 哈希值对 100 取模这样可以保证同一个用户在一分钟内始终落同一个 region而不同用户会分散到不同 region。最终 RowKey 的形态是盐值_时间_IP既支持查看“某个时间段的所有新闻访问”又分散了写入压力。参数说明% 100的 100 是散列桶的数量。如果你的 HBase 表预分区为 10 个 region建议改成% 10让盐值范围与 region 数量匹配否则会出现某些 region 空置、某些 region 过载。这个数字不是拍脑袋定的要和建表时的预分区数量对齐。另外SDF不是线程安全的多线程下必须用ThreadLocal或改为传参格式化否则在 Flume 高并发场景下会偶发性生成错误的时间字符串。这段代码是整套系统里我最想强调的部分因为它直接影响 HBase 的写入性能和查询效率也是你在答辩时能讲清楚的设计点。5. 参考步骤.txt 里的五条避坑记录从启动失败到数据不落地我把参考步骤.txt和实际调试过程中遇到最典型的五个问题整理出来每条都是真实翻车记录。按“现象 → 原因 → 解决”的格式写方便你复现时直接对照。5.1 Kafka 消费者持续堆积消费速度追不上生产速度现象运行kafka-consumer-groups.sh查看 lag 时group 的 lag 持续上涨Spark Streaming 处理速度远远低于 Flume 写入速度。原因Kafka topic 的分区数只有一个Spark Streaming 读取时只能开启一个 receiver 或一个 partition 的并行度单线程消费成为瓶颈。另一个原因是 Flume 侧batchSize设置过大Kafka channel 堆积了太多 Event但 Spark 端每次拉取的 maxOffsetsPerTrigger 没有调大。解决优先把 topic 分区数改为 3 或 6再调整 Spark Streaming 的并行度。我实际修改的方式是kafka-topics.sh --alter --topic news-log --partitions 6 --bootstrap-server localhost:9092对应 Spark 端设置val kafkaParams Map[String, Object]( bootstrap.servers - localhost:9092, max.partition.fetch.bytes - 10485760, max.poll.records - 5000 ) val stream KafkaUtils.createDirectStream[String, String]( ssc, LocationStrategies.PreferConsistent, ConsumerStrategies.Subscribe[String, String](topics, kafkaParams) )参数说明max.poll.records控制每次拉取的最大记录数调大后单次处理量提升max.partition.fetch.bytes控制每次拉取的字节上限。两个参数要配合调整如果只调max.poll.records而字节上限不变还是会被截断。5.2 HBase 写入出现 Region 热点部分 RegionServer 负载极高现象HBase 的 Web UI 上某个 region 的 store 文件数和请求数远超其他 region写入延迟从几毫秒涨到几百毫秒。原因RowKey 没有做加盐处理或者盐值范围远大于 region 数量。时间戳顺序递增的 RowKey 会被 HBase 持续写入同一个 region 的尾部形成顺序写热点。解决使用SimpleRowKeyGenerator中的加盐逻辑并让盐值范围和 HBase 表的预分区数量匹配。建表时用如下方式预分区hbase shell create news_weblogs, info, {NUMREGIONS 10, SPLITALGO HexStringSplit}然后用% 10替换% 100让盐值分布在 0-9 之间与 10 个 region 一一对应。改完后观察 region 请求分布热点立即消失。5.3 Spark Streaming 运行一段时间后自动停止日志里出现 checkpoint 异常现象作业在前几分钟正常运行半小时后突然停止日志中提示 checkpoint 目录不可用或序列化失败。原因checkpoint 目录配置到了本地磁盘但多次重启后目录被污染同时 checkpoint 里保存了不能序列化的对象比如连接 Kafka 的非序列化客户端。解决把 checkpoint 指向 HDFS 的独立目录并在每次修改代码后清空该目录重新启动。我实际使用的方式ssc.checkpoint(hdfs://localhost:9000/spark/news-checkpoint)修改代码或依赖版本后先删除 checkpoint 目录再提交作业防止旧状态干扰新逻辑。5.4 Flume 启动报ClassNotFound: AsyncHbaseEventSerializer或NoSuchMethodError现象执行flume-ng agent命令时控制台直接报序列化器类找不到或者报put方法签名不匹配。原因自定义序列化器的 jar 没有放到 Flume 的 lib 目录Flume 启动时加载不到或者flume-ng-hbase-sink.jar和 HBase 客户端版本不一致导致方法签名错误。解决把flume_hbase/flume-ng-hbase-sink.jar软链到 Flume 的lib目录并确认 HBase 客户端的 jar 也在同目录。我一般会写一个启动脚本动态把项目 jar 加入FLUME_CLASSPATHexport FLUME_CLASSPATH/your/path/flume-ng-hbase-sink.jar flume-ng agent --conf conf --conf-file conf/flume.conf --name agent这是最简单也最不容易忘的解决办法省得每次手动复制 jar。5.5 可视化大屏数据不刷新HBase 查询超时现象图表区域长时间空白浏览器控制台显示接口请求报超时HBase 日志里出现 scanner 超时异常。原因可视化层扫描 HBase 时没有设置缓存和限制条件某些扫描请求要获取大量行远超 HBase 默认的 scanner 超时时间。解决优化 scan 的查询范围加限制条件同时设置缓存。我在可视化代码里一般这么调Scan scan new Scan(); scan.setStartRow(Bytes.toBytes(0_202401010000)); scan.setStopRow(Bytes.toBytes(z_202401012359)); scan.setCaching(500); scan.setLimit(1000);参数说明setStartRow和setStopRow利用 RowKey 的时间和盐值范围做索引扫描避开全表扫描setCaching(500)指每次从 region server 取 500 行减少网络往返数量setLimit(1000)兜底限制最大返回行数可视化的展示一般只需要最新一批数据加 limit 防止内存溢出。6. 验证链路与进阶把默认图表改成能说服答辩官的动态展示6.1 验证整条链路的三步走部署完成后我习惯按“造数 → 看 Kafka → 看 HBase → 看大屏”四步验证不是直接刷新页面。第一步手动往日志目录灌几条带明确特征的测试数据。比如修改一条日志的 URL 为/test-pv-check时间戳换成当前分钟。第二步在 Kafka consumer 里确认这条日志出现在 Spark 日志里确认窗口计算触发。第三步在 HBase shell 里 scan 当天的数据确认这条 URL 的 PV 计数已经累加。第四步才是打开可视化页面确认对应的折线或柱状图出现了这个测试访问数。这个过程看起来繁琐但能帮你把“数据没出来”的问题快速定位到具体环节而不是在大屏里瞎猜。6.2 进阶技巧之一把静态轮询改成事件驱动默认可视化端大概率是定时轮询 HBase间隔几秒查一次这样数据有延迟且查询压力不小。我一般会改用 WebSocket 推送Spark Streaming 每批次完成时主动通知前端刷新。const ws new WebSocket(ws://localhost:8088/ws/news); ws.onmessage function(event) { const data JSON.parse(event.data); drawChart(data); };配合后端更新常见架构是在 Spark 每批 result 写完后同时向 Redis 发布一条消息Web 端订阅 Redis channel 再推送前端。6.3 进阶技巧之二预聚合表与明细表分离如果毕业设计想往深了讲可以把 HBase 分成两张表一张存分钟级聚合结果news_agg_minute一张存原始访问明细news_weblogs。聚合表供大屏查询明细表保留数据备查和回溯实验。这样既减少大屏查询压力又能展示你理解了存储层设计思想。我后来每次部署这套项目都会强制走一遍“先压测 Kafka 分区 → 再改 RowKey 盐值 → 最后验证 HBase 扫描范围”的流程少了任何一步都会在下一次启动时以某种方式翻车。把那五条避坑记录贴到自己的部署笔记里基本上能帮你少走掉一整晚的调试弯路。希望帮到你。本文还有配套的精品资源点击获取
返回列表