简介:面向毕业设计任务及Spark入门开发者,提供一套基于Spark2.2的新闻网大数据实时分析系统完整源码,覆盖新闻日志采集、消息缓存、分布式存储与流式处理的经典链路,适用于需要快速搭建实时行为分析原型的课程项目或论文验证场景。压缩包仅3.64MB,共43个文件,其中scala/java源码与jar包构成核心实现,xml与mf补充运行配置,js和png支撑前端可视化,README与参考步骤txt则说明架构与启动方式,目录划分清晰便于按模块研读。已有54人学习,适合正选相关课题或希望上手Spark流计算的读者。通过阅读项目,可重点理解KfkAsyncHbaseEventSerializer的事件序列化思路、SimpleRowKeyGenerator的HBase行键设计,以及Flume与HBase对接时自定义sink的注册方法;附带图片和样例日志还能辅助还原页面展示效果与模拟数据,为从头搭建类似系统提供可复用的参考。
1. 为什么把日志管道放在 Spark2.2 之前
很多拿到大数据毕设源码的人,第一反应是去翻 Spark RDD 算子,但真正卡住进度的往往是数据入口。基于 Spark2.2 的新闻网大数据实时分析系统,核心链路并不复杂:Flume 采集 web 服务器上的新闻访问日志,经过自定义 HBase Sink 落一份原始数据到 HBase;同时 Spark Streaming 直连 Kafka 消费同一份日志,做窗口聚合,输出热点新闻和用户行为趋势。这套源码的价值在于把「日志怎么进、RowKey 怎么设计、Spark 怎么消费」这三个问题一次性给了可运行答案,特别适合正在做大数据毕业设计、或者想复现实时数仓最小闭环的开发者。先从数据源头拆起。
2. 定制 Flume HBase Sink:从 Event 到 HBase 的最后一公里
2.1 为什么选 Flume 作为消息入口
新闻网站的多台 Web 服务器会产生大量分散的访问日志,Flume 的优势在于 Source 类型丰富。taildir source 支持断点续传,avro source 可以对接 logstash 等其他采集端,kafka source 可以直接消费消息队列里的日志,这些都是直接写 KafkaProducer 要重复造轮子的部分。更关键的是,Flume 有事务语义,source 到 channel 到 sink 的每一跳都可以配置容量和重试参数,单条日志丢失的几率远小于自己写一个采集线程。
源码里flume_hbase目录直接放了一个编译好的flume-ng-hbase-sink.jar,说明作者没有用 Flume 自带的org.apache.flume.sink.hbase.HBaseSink默认行为,而是对它做了二次封装。默认的SimpleHbaseEventSerializer只会把 event body 作为单列写入,对 JSON 日志完全不够用,所以源码里出现了KfkAsyncHbaseEventSerializer.java和SimpleRowKeyGenerator.java,这两个类才是整个 Flume 层的核心。
2.2 KfkAsyncHbaseEventSerializer 解析逻辑
KfkAsyncHbaseEventSerializer继承了 Flume 的AbstractHbaseEventSerializer,核心是覆写getActions()方法,把 Flume Event 转成 HBase 的Put列表。类名里的KfkAsync暗示它处理的是来自 Kafka 的消息,并且写 HBase 时使用异步客户端,避免每一条访问日志都阻塞在 RPC 上。
实际源码结构不长,核心逻辑可以近似抽象为下面这段 Java 代码:
public class KfkAsyncHbaseEventSerializer extends AbstractHbaseEventSerializer { private RowKeyGenerator rowKeyGenerator; private Event event; @Override public List<Put> getActions() { String body = new String(event.getBody(), Charsets.UTF_8); JSONObject obj = JSON.parseObject(body); String rowKey = rowKeyGenerator.generateRowKey( obj.getString("timestamp"), obj.getString("newsId")); Put put = new Put(Bytes.toBytes(rowKey)); put.addColumn(cf, Bytes.toBytes("newsId"), Bytes.toBytes(obj.getString("newsId"))); put.addColumn(cf, Bytes.toBytes("uid"), Bytes.toBytes(obj.getString("uid"))); put.addColumn(cf, Bytes.toBytes("pageUrl"), Bytes.toBytes(obj.getString("pageUrl"))); return Collections.singletonList(put); } @Override public void setRowKeyGenerator(RowKeyGenerator rowKeyGenerator) { this.rowKeyGenerator = rowKeyGenerator; } }这段代码里最容易踩坑的是JSON.parseObject(body),Flume Event 的 body 是byte[],如果上游日志不是标准 JSON,或者编码不是 UTF-8,这里会直接抛异常。我一般会解析失败时把原始 body 写入一个单独的parse_error表,而不是把 agent 直接跑死。另外put.addColumn用到的cf通常是构造时传入的 columnFamily,对应 hbase site 配置里的columnFamily,两者拼写不一致是最低级的错误,但在毕设源码里出现频率非常高。
2.3 SimpleRowKeyGenerator 的边界条件
SimpleRowKeyGenerator.java的存在说明作者知道裸用时间戳做 RowKey 是个坑。新闻日志的写入模式是持续追加,如果 RowKey 用timestamp + newsId,那么同一秒钟的写入都会落在同一个 RegionServer 上,形成典型的写热点。常见做法是把时间戳反转,再拼接一个从 uid 或 newsId 派生的散列值:
public class SimpleRowKeyGenerator implements RowKeyGenerator { public String generateRowKey(String timestamp, String newsId) { String reversedTs = new StringBuilder(timestamp).reverse().toString(); int salt = Math.abs(newsId.hashCode() % 100); return reversedTs + "_" + salt; } }反转时间戳的目的是让越新的数据在字典序上越靠前,同时让连续写入分散到不同的 Region 范围。这里的salt不需要太大,50 到 200 之间就够了,太大反而会让扫描某一天的数据时需要跨过多得多的 Region。面试里被问到的“HBase RowKey 设计为什么不用原样时间戳”,答案就在这个类里。
2.4 flume-ng-hbase-sink 配置实战
整个 Flume agent 的关键配置如下,注意serializer必须指向自定义类,并且rowKeyGenerator是自定义序列化器内部使用的参数:
agent.sources = kafka-source agent.channels = mem-channel agent.sinks = hbase-sink 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 = newslog agent.sources.kafka-source.kafka.consumer.group.id = flume-hbase-group agent.sources.kafka-source.kafka.auto.offset.reset = latest agent.channels.mem-channel.type = memory agent.channels.mem-channel.capacity = 10000 agent.channels.mem-channel.transactionCapacity = 1000 agent.sinks.hbase-sink.type = asynchbase agent.sinks.hbase-sink.table = news_events agent.sinks.hbase-sink.columnFamily = cf agent.sinks.hbase-sink.serializer = com.example.serializer.KfkAsyncHbaseEventSerializer agent.sinks.hbase-sink.serializer.rowKeyGenerator = com.example.serializer.SimpleRowKeyGenerator agent.sinks.hbase-sink.batchSize = 500 agent.sources.kafka-source.channels = mem-channel agent.sinks.hbase-sink.channel = mem-channel启动命令:
flume-ng agent \ -n agent \ -c conf \ -f flume-hbase.conf \ -Dflume.root.logger=INFO,console这里几个参数值得展开:
| 参数 | 建议值 | 说明 |
|---|---|---|
transactionCapacity | 1000 | 不能超过 capacity,否则运行时报 channel 空间不足 |
batchSize | 500 | 控制单次批量写 HBase 的条数,太大容易把 RegionServer 的 memstore 打满 |
kafka.auto.offset.reset | latest | 在 Flume 层一般用 latest,避免从头回放历史日志 |
serializer.rowKeyGenerator | 自定义类全名 | 该参数由自定义 serializer 内部读取,不要拼错包名 |
3. Spark2.2 实时消费与热点窗口统计
3.1 读 HBase 还是直连 Kafka
Flume 已经把数据落进 HBase 了,Spark 作业是不是直接TableInputFormat读 HBase 就行?可以,但那是离线批处理思路。HBase 扫描的延迟在百毫秒到秒级,且 scan 会加大 RegionServer 压力,处理“最近 5 分钟新闻点击量”这种需求时,Spark 直接消费 Kafka 才是实时分析的正确姿势。
| 方案 | 延迟 | 吞吐 | 代码复杂度 | 适用场景 |
|---|---|---|---|---|
| Spark 读 HBase | 秒级 | 受 RegionServer 扫描性能限制 | 低 | 离线报表、历史数据回填 |
| Spark 直连 Kafka | 毫秒级 | 受 partition 数和消费能力限制 | 中 | 热点新闻、实时用户行为分析 |
这套源码里 Flume 的 Kafka source 和 Spark 的 Kafka consumer 用的是同一个 topic,说明作者有意保留了 HBase 里的明细数据,同时又让 Spark 消费 Kafka 做实时计算,两条链路互不干扰。实际生产里这种架构很常见,HBase 那一侧相当于可回溯的原始日志仓库,Spark Streaming 只负责窗口内的聚合。
3.2 createDirectStream 消费代码
Spark2.2 里推荐用的是org.apache.spark.streaming.kafka010.KafkaUtils,需要引入spark-streaming-kafka-0-10_2.11依赖。消费端代码大致如下:
val kafkaParams = Map[String, Object]( "bootstrap.servers" -> "node01:9092,node02:9092", "key.deserializer" -> classOf[StringDeserializer], "value.deserializer" -> classOf[StringDeserializer], "group.id" -> "spark-news-analysis", "auto.offset.reset" -> "latest", "enable.auto.commit" -> (false: java.lang.Boolean) ) val topics = Array("newslog") val stream = KafkaUtils.createDirectStream[String, String]( ssc, LocationStrategies.PreferConsistent, ConsumerStrategies.Subscribe[String, String](topics, kafkaParams) ) val lines = stream.map(_.value()) lines.foreachRDD { rdd => val offsetRanges = rdd.asInstanceOf[HasOffsetRanges].offsetRanges // 业务处理 processRdd(rdd) stream.asInstanceOf[CanCommitOffsets].commitAsync(offsetRanges) }这里enable.auto.commit=false是关键,处理完业务逻辑后再手动 commit,避免数据未处理就提交 offset 导致丢数据。LocationStrategies.PreferConsistent让 Kafka partition 尽量分布在不同 executor 上,避免某个 executor 负载过重。如果 topic 分区多而 executor 少,PreferConsistent会退化成PreferFixed策略,要注意观察 spark ui 里各 executor 的输入速率是否均匀。
3.3 reduceByKeyAndWindow 做热点计算
热点新闻本质上是“最近一段时间内被点击最多的新闻”。Spark DStream 提供了窗口算子,不需要自己维护中间状态:
val windowedCounts = lines .map(json => (extractNewsId(json), 1L)) .reduceByKeyAndWindow( (a: Long, b: Long) => a + b, Seconds(300), Seconds(60) ) windowedCounts .transform(rdd => rdd.sortBy(_._2, ascending = false)) .foreachRDD { rdd => rdd.take(10).foreach(println) }参数里Seconds(300)是窗口长度,代表统计最近 5 分钟的数据;Seconds(60)是滑动间隔,代表每 60 秒输出一次结果。窗口长度大于滑动间隔,每个批次的数据会被重复计算,这是流式窗口的固有语义。下面这个表是常见新闻场景的参数选择参考:
| 业务场景 | 窗口长度 | 滑动间隔 | 说明 |
|---|---|---|---|
| 新闻秒级热点 | 60s | 10s | 实时性最高,但计算量大 |
| 5 分钟热度排名 | 300s | 60s | 最常见的榜单周期 |
| 小时级趋势 | 3600s | 300s | 适合做舆情趋势,延迟可接受 |
需要注意的是,窗口聚合必须开启 checkpoint,否则 driver 宕机后状态无法恢复。在SparkConf里增加spark.streaming.backpressure.enabled=true和spark.streaming.kafka.maxRatePerPartition可以防止 Kafka 积压时打爆 downstream,这两个参数在 Spark2.2 里是控制实时分析稳定性最有效的手段。
4. weblogs 日志的清洗与结果落地
4.1 weblogs 日志格式与字段提取
源码附带的weblogs目录是典型的 nginx 访问日志,每行长这样:
127.0.0.1 - - [02/Jul/2019:14:22:01 +0800] "GET /news/1024?from=mobile HTTP/1.1" 200 1024 "http://example.com" "Mozilla/5.0"要解析出新闻 ID、ip、访问时间、响应码、UA,直接用空格 split 会踩引号和方括号的坑。更稳妥的方式是用正则切字段,把 raw line 映射成一个 case class:
case class NewsLog( ip: String, timestamp: String, method: String, newsId: String, status: Int, ua: String ) val LogRegex = """^(\S+) \S+ \S+ \[([^\]]+)\] "(\S+) (/news/(\d+)[^ "]*)" (\d{3}) \S+ "([^"]*)".*""".r def parseLog(line: String): Option[NewsLog] = { line match { case LogRegex(ip, time, method, _, newsId, status, ua) => Some(NewsLog(ip, time, method, newsId, status.toInt, ua)) case _ => None } }这里要注意正则里/news/(\d+)是业务约定的 URL 规则,如果新闻频道还有/sports/1025或/tech/1026,需要把正则扩展成/news|sports|tech/,否则这些日志会被全部过滤掉,热点排行就会失真。源码里weblogs目录下的样例日志格式和这个正则是匹配的。
4.2 清洗规则与常见误区
日志清洗不是写完正则就完事,实际运行时会遇到三类问题:静态资源请求、爬虫流量、时间字符串不统一。
def isStaticResource(newsId: String): Boolean = { newsId.endsWith(".jpg") || newsId.endsWith(".png") || newsId.endsWith(".css") } val parsedLogs = lines .flatMap(parseLog _) .filter(log => !isStaticResource(log.newsId)) .filter(log => !log.ua.contains("Baiduspider") && !log.ua.contains("Googlebot"))第一个 filter 很直观,凡是没有 newsId 的请求直接丢弃。第二个 filter 需要注意,爬虫会制造大量无意义的点击,如果毕设答辩时被问“热点新闻里为什么全是垃圾内容”,多半就是没过滤爬虫。更严格的做法是维护一个爬虫 UA 黑名单,放到 Redis 里动态加载,而不是写死在代码里。
时间解析也是高频坑。nginx 默认的[02/Jul/2019:14:22:01 +0800]不是标准 ISO 格式,如果要按小时做趋势分析,必须转成yyyy-MM-dd HH:mm:ss。我一般用DateTimeFormatter.ofPattern("dd/MMM/yyyy:HH:mm:ss Z", Locale.ENGLISH),注意必须带Locale.ENGLISH,否则中文服务器上 December 这类英文月份解析直接失败。这个细节在本地 Windows 上测不出来,要到 Linux 生产环境才会暴露。
4.3 结果输出层设计
窗口计算的结果需要写出去才有人看。毕设源码里最常见的是写入 MySQL 或者 Redis。下面是一个用foreachRDD写 MySQL 的骨架:
windowedCounts.foreachRDD { rdd => rdd.foreachPartition { partition => val conn = JdbcUtil.getConnection() partition.foreach { case (newsId, cnt) => val sql = """ |INSERT INTO hot_news(news_id, cnt, window_time) |VALUES (?, ?, ?) |ON DUPLICATE KEY UPDATE cnt = cnt + ? """.stripMargin val ps = conn.prepareStatement(sql) ps.setString(1, newsId) ps.setLong(2, cnt) ps.setTimestamp(3, currentWindowTime) ps.setLong(4, cnt) ps.executeUpdate() ps.close() } conn.close() } }这里每个 executor 都会打开一个 JDBC 连接,必须用连接池而不是每处理一条就DriverManager.getConnection。ON DUPLICATE KEY UPDATE让同一窗口时间内重复输出时做累加而不是报主键冲突。如果数据量更大,可以改成写 Kafka 下游或直接落到 Redis sorted set,让前端接口实时拉取 Top N。毕设项目写到 MySQL 这个粒度已经足够展示实时分析全链路。
5. 从单机调试到集群运行的注意点
5.1 版本匹配清单
这套源码基于 Spark2.2,配套组件如果版本不对,经常出现ClassNotFoundException或者NoSuchMethodError。参考步骤.txt 里一般会写版本,但我实际调通的一套组合是:
| 组件 | 版本 | 备注 |
|---|---|---|
| Spark | 2.2.0 | spark-streaming-kafka-0-10_2.11 必须用 0-10 的坐标 |
| Kafka | 0.10.2.1 | 与 Spark2.2 的 kafka010 consumer API 匹配 |
| Flume | 1.8.0 | 自带 hbase sink,但需替换自定义 jar |
| HBase | 1.3.1 | asynchbase sink 与 HBase 1.x 兼容性最好 |
| JDK | 1.8 | Spark2.2 不支持更高版本 |
最典型的报错是java.lang.NoSuchMethodError: org.apache.kafka.clients.consumer.KafkaConsumer.subscribe,这通常是因为spark-streaming-kafka-0-10和kafka-clients版本不一致。解决方式是检查flume-ng-hbase-sink.jar内部打包的 kafka-client 版本,和 Spark 任务里provided的版本对齐。
5.2 本地先跑通 HBase 写入
不要一上来就启动整个 Flume agent。先在本地单机 HBase 里验证自定义 serializer 能不能正确生成 Put。建表后可以直接用 JUnit 或者一个 main 方法跑:
hbase shell <<EOF create 'news_events', 'cf' EOFEvent event = new SimpleEvent(); event.setBody(("{\"newsId\":\"1024\",\"timestamp\":\"02/Jul/2019:14:22:01 +0800\"," + "\"uid\":\"u1001\",\"pageUrl\":\"/news/1024\"}") .getBytes(StandardCharsets.UTF_8)); KfkAsyncHbaseEventSerializer serializer = new KfkAsyncHbaseEventSerializer(); serializer.initialize(event, Bytes.toBytes("news_events"), Bytes.toBytes("cf")); serializer.setRowKeyGenerator(new SimpleRowKeyGenerator()); List<Put> puts = serializer.getActions(); try (Connection conn = ConnectionFactory.createConnection(conf); Table table = conn.getTable(TableName.valueOf("news_events"))) { table.put(puts); }这个做法的价值是把 Flume、Kafka 全部遮挡掉,先确认 serializer 逻辑正确。我调试时经常发现 JSON 里字段名多一个空格、时间戳格式不对之类的问题,全在这一步能提前暴露。
5.3 集群部署时的连接数与时区问题
在真正往集群上部署这批作业时,有两个问题比业务逻辑更容易导致事故。第一个是 HBase 连接数,每个 executor 任务如果都ConnectionFactory.createConnection,几十个 executor 会把 RegionServer 的 RPC handler 全占满。正确做法是用一个静态Connection对象,并配置hbase.client.write.buffer和hbase.client.connection.impl,让多个线程共享底层连接池。大数据集群部署策略里对 executor 数量和 HBase listener 数目的比例要提前算好,否则加机器反而会拖垮 HBase。
第二个是时区。weblogs日志的时间戳带+0800,而 Spark driver 和 executor 默认使用系统时区,如果集群机器是 UTC,那么窗口聚合的边界会和北京时间错开 8 小时,热点榜看起来是准的,但时间维度全部错位。排查方法很简单,在测试数据里故意用两个跨整点的时间戳,看窗口切换是否符合预期。解决方式是把 Spark 作业启动参数里加-Duser.timezone=GMT+8,并且解析日志时使用带时区的OffsetDateTime,不要用本地时区的LocalDateTime。
如果 HBase 出现写入毛刺,优先调大hbase.client.write.buffer,默认 2MB 可以调到 8MB,而不是盲目增加 Flume 的并发线程。这条经验比替换组件版本更立竿见影。
本文还有配套的精品资源,点击获取