ARTICLE DETAIL

资讯详情

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

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

基于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.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 ListPut 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.loggerINFO,console这里几个参数值得展开参数建议值说明transactionCapacity1000不能超过 capacity否则运行时报 channel 空间不足batchSize500控制单次批量写 HBase 的条数太大容易把 RegionServer 的 memstore 打满kafka.auto.offset.resetlatest在 Flume 层一般用 latest避免从头回放历史日志serializer.rowKeyGenerator自定义类全名该参数由自定义 serializer 内部读取不要拼错包名3. Spark2.2 实时消费与热点窗口统计3.1 读 HBase 还是直连 KafkaFlume 已经把数据落进 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.commitfalse是关键处理完业务逻辑后再手动 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.enabledtrue和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?frommobile HTTP/1.1 200 1024 http://example.com Mozilla/5.0要解析出新闻 ID、ip、访问时间、响应码、UA直接用空格 split 会踩引号和方括号的坑。更稳妥的方式是用正则切字段把 raw line 映射成一个 case classcase 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 里一般会写版本但我实际调通的一套组合是组件版本备注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但需替换自定义 jarHBase1.3.1asynchbase sink 与 HBase 1.x 兼容性最好JDK1.8Spark2.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()); ListPut 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.timezoneGMT8并且解析日志时使用带时区的OffsetDateTime不要用本地时区的LocalDateTime。如果 HBase 出现写入毛刺优先调大hbase.client.write.buffer默认 2MB 可以调到 8MB而不是盲目增加 Flume 的并发线程。这条经验比替换组件版本更立竿见影。本文还有配套的精品资源点击获取
返回列表