news 2026/9/11 10:48:03

基于Spark2.2的新闻日志实时分析:Flume自定义Sink与Streaming实践

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
基于Spark2.2的新闻日志实时分析:Flume自定义Sink与Streaming实践

简介:面向毕业设计任务及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.javaSimpleRowKeyGenerator.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

这里几个参数值得展开:

参数建议值说明
transactionCapacity1000不能超过 capacity,否则运行时报 channel 空间不足
batchSize500控制单次批量写 HBase 的条数,太大容易把 RegionServer 的 memstore 打满
kafka.auto.offset.resetlatest在 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 秒输出一次结果。窗口长度大于滑动间隔,每个批次的数据会被重复计算,这是流式窗口的固有语义。下面这个表是常见新闻场景的参数选择参考:

业务场景窗口长度滑动间隔说明
新闻秒级热点60s10s实时性最高,但计算量大
5 分钟热度排名300s60s最常见的榜单周期
小时级趋势3600s300s适合做舆情趋势,延迟可接受

需要注意的是,窗口聚合必须开启 checkpoint,否则 driver 宕机后状态无法恢复。在SparkConf里增加spark.streaming.backpressure.enabled=truespark.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.getConnectionON DUPLICATE KEY UPDATE让同一窗口时间内重复输出时做累加而不是报主键冲突。如果数据量更大,可以改成写 Kafka 下游或直接落到 Redis sorted set,让前端接口实时拉取 Top N。毕设项目写到 MySQL 这个粒度已经足够展示实时分析全链路。

5. 从单机调试到集群运行的注意点

5.1 版本匹配清单

这套源码基于 Spark2.2,配套组件如果版本不对,经常出现ClassNotFoundException或者NoSuchMethodError。参考步骤.txt 里一般会写版本,但我实际调通的一套组合是:

组件版本备注
Spark2.2.0spark-streaming-kafka-0-10_2.11 必须用 0-10 的坐标
Kafka0.10.2.1与 Spark2.2 的 kafka010 consumer API 匹配
Flume1.8.0自带 hbase sink,但需替换自定义 jar
HBase1.3.1asynchbase sink 与 HBase 1.x 兼容性最好
JDK1.8Spark2.2 不支持更高版本

最典型的报错是java.lang.NoSuchMethodError: org.apache.kafka.clients.consumer.KafkaConsumer.subscribe,这通常是因为spark-streaming-kafka-0-10kafka-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' EOF
Event 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.bufferhbase.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 的并发线程。这条经验比替换组件版本更立竿见影。

本文还有配套的精品资源,点击获取

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/11 10:47:37

deer-flow实战:可视化编排LLM工作流,打造可靠AI业务流

说实话&#xff0c;第一次看到 deer-flow 这个项目名称&#xff0c;我以为是某个做数据管道的新玩具。结果真正上手之后&#xff0c;我发现它解决的是我一直很头疼的问题&#xff1a;怎么把一堆 LLM 调用编排成一条可靠、可观测、能上生产的业务流。过去我们聊智能体&#xff0…

作者头像 李华
网站建设 2026/9/11 10:46:26

iOS 4.3审核被拒怎么办?小蟹iOS混淆4.3实战解析

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/11 10:45:33

从MovieLens实战看协同过滤的落地边界与工程细节

简介&#xff1a;面向计算机相关专业学生、算法初学者及需要推荐系统参考的开发者&#xff0c;这份资源基于 MovieLens 公开数据集&#xff0c;实现了一个完整可运行的协同过滤推荐算法项目。内容涵盖数据预处理、用户/物品相似度计算、评分预测与结果评估等核心环节&#xff0…

作者头像 李华
网站建设 2026/9/11 10:45:29

盗图与图片指纹:原创检测不是吓唬人

盗图与图片指纹&#xff1a;原创检测不是吓唬人 一次盗图投诉的代价清单&#xff1a; 「被同行投诉盗图的时候&#xff0c;我以为顶多删个图。结果&#xff1a;链接下架、扣分、被投诉的那批图全部换掉、申诉期两周。最气的是图的来源——供货商统一发的素材包&#xff0c;全行…

作者头像 李华
网站建设 2026/9/11 10:43:39

基于Java springboot化妆品推荐系统(源码+lw+部署文档+讲解等)

温馨提示&#xff1a;本人主页置顶文章(点我)开头有 CSDN 平台官方提供的学长联系方式的名片&#xff01; 温馨提示&#xff1a;本人主页置顶文章(点我)开头有 CSDN 平台官方提供的学长联系方式的名片&#xff01; 温馨提示&#xff1a;本人主页置顶文章(点我)开头有 CSDN 平台…

作者头像 李华