基于Spark的电子书平台数据分析实战:聚合、窗口函数与数据质量 简介面向大数据初学者与毕业设计者的Spark实战资料包围绕电子书平台数据分析场景完整实现下载量Top10统计、用户评分实时计算、分时段在线人数聚合等典型需求可作为课程设计或工程参考。压缩包共51个文件核心为36个Scala源码文件覆盖数据清洗、标签统计与输出逻辑另有9个XML工程配置、2份DOCX项目文档及JSON样例数据整体1.67MB结构紧凑便于快速理解。项目采用常见Spark算子与流式处理思路文档中阐述了三个功能模块的需求拆分和实现过程适合希望掌握Spark批处理与实时统计做法的开发者对照学习。目前已有117人学习下载对于想以真实业务串联Spark API、完善毕业设计代码与文档的读者来说具备较高参考价值。1. 基于Spark的电子书平台数据分析不止是写几个聚合算子的事把电子书平台的下载量统计、评分聚合、阅读时段分析做成Spark作业我第一次拆这套源码时最深的感触是真正难的并不是RDD算子怎么写而是四个业务需求各自的数据口径怎么统一。这份基于Spark的电子书平台数据分析项目正好把「标签下载量、标签内Top10、书籍评分、四时段在线人数」四类指标揉进同一套Maven工程里用u.data和click_data.json两个输入文件喂给Spark作业。它适合刚学完Spark Core、想找一份完整数据分析项目实践来练手的人也适合被公司安排做电子书或内容平台报表、需要参考真实工程结构的朋友。下面我按源码实际的目录和需求拆开讲每一步都会给到能直接跑的代码和参数说明。2. 工程结构与数据源接入先把u.data和click_data.json吃透再说聚合2.1 解压后先看目录Maven工程里哪些文件是真正要关注的拿到E-Book-master.zip后解压出来的目录结构和常见的Spark课程设计项目不太一样它里面除了pom.xml和src目录还有一个电子书平台管理数据分析设计与实现.docx文档和output输出目录。我第一次看到~$开头那个文件就知道是Word临时锁文件可以直接忽略。真正要关心的是下面几个东西路径作用备注pom.xmlMaven依赖与打包配置确认Spark版本和Scala版本src/main/input/u.data评分数据文件格式和MovieLens的u.data一致src/main/input/click_data.json点击下载数据JSON格式包含标签字段output/作业输出目录跑完看结果电子书平台管理数据分析设计与实现.docx项目说明文档含需求描述和模块划分pom.xml里如果用的是Spark 2.x那么Scala版本大概率是2.11如果用Spark 3.x则对应Scala 2.12。这个细节决定了你在IDEA里写代码时选的Scala编译器版本也决定了你打包后能不能在集群上正常提交。我一般先查pom里spark-core的artifactId再定Scala版本两个不匹配时本地能跑、提交到集群就会报NoSuchMethodError的玄学错误。2.2 数据格式与Spark读取参数JSON读取的坑从这一步开始u.data这个文件名的数据格式是用户ID | 书籍ID | 评分 | 时间戳字段之间用竖线分隔和MovieLens的评分数据格式完全一致。第一行数据长这样196 242 3 881250949前面两个数字是用户和书籍的内部ID第三个是1到5的评分第四个是unix时间戳。注意这里的时间戳单位是秒如果后续要用from_unixtime做时段统计直接传这个值就行但如果click_data.json里的时间字段是毫秒13位数字就必须先除以1000这个差别我在避坑章节会专门讲。click_data.json则是另一套结构它承载了标签维度的点击下载记录。读取JSON时有个非常常见的翻车点Spark默认认为JSON文件是一个对象占一行如果文件里每个JSON对象被格式化成了多行就必须显式开启multiLine选项。我第一次读这份数据时没加这个参数结果全部字段解析出来都是null排查了半天才意识到是JSON文件里每条记录被折行了。正确读取方式如下val spark SparkSession.builder() .appName(EBookAnalysis) .master(local[*]) .getOrCreate() val clickDF spark.read .option(multiLine, true) // 每条JSON记录占多行时必开 .json(src/main/input/click_data.json) val ratingDF spark.read .option(delimiter, |) // u.data是竖线分隔 .csv(src/main/input/u.data) .toDF(userId, bookId, rating, ts)这里multiLine选项决定了Spark能不能正确解析多行JSONdelimiter设成竖线后csv读取器才不会把整行当成一个字段。.toDF重命名列名是为了后面写SQL和聚合时不至于用_c0这种无意义列名。跑通这两行代码四个需求的数据源就全部就绪了。2.3 把两个数据源的公共字段对齐click_data.json里存的是点击下载记录字段通常是userId、bookId、tag、clickTime这类u.data里存的是评分记录。要同时做下载量统计和评分统计需要确认两边的书籍ID能否对上。实际项目中如果bookId一致就可以用join把标签信息带到评分数据上如果不一致就只能在各自数据集内部做聚合。这份源码的做法是分开统计下载量模块只依赖click_data.json评分模块只依赖u.data阅读时段模块用u.data的时间戳。这样设计的好处是模块间不耦合坏处是你没法在SQL里直接join出「某标签书籍的评分Top10」所以后面做Top10时只能在标签内部按下载量排。3. 下载量统计模块从groupBy到窗口函数实现标签Top103.1 先明确需求边界总下载量和Top10是两套口径摘要里写的第一个需求是「实时统计各标签电子书的总点击下载量」第二个是「实时统计各标签中电子书下载量Top10」。这两句话看起来都是下载量但前者聚合维度是tag后者聚合维度是tag加bookId。很多新人第一次写会搞混用groupBy(tag)出了总量就以为能直接拿到每本书的排行其实排行要先按tag, bookId分组得到每个书在标签内的下载量再做窗口排序。这里我习惯先定义一个中间视图把点击数据按标签和书聚合好clickDF.createOrReplaceTempView(clicks) val downloadsPerBook spark.sql( |SELECT tag, bookId, COUNT(*) AS downloads |FROM clicks |GROUP BY tag, bookId |.stripMargin)这一段SQL把每条点击记录折叠成「某个标签下某本书的累计下载量」COUNT(*)统计的是点击次数。需要注意的是如果click_data.json里存在同一用户重复点击同一本书的情况这个口径会把重复点击也算进去这一点产品上是否接受要提前确认。我处理时一般会在SQL里先DISTINCT userId, bookId, tag去重再计数除非需求明确要统计总点击次数。3.2 标签总下载量用Spark SQL还是DataFrame API总下载量的实现非常简单SELECT tag, COUNT(*) FROM clicks GROUP BY tag即可。但如果想控制shuffle粒度建议用DataFrame API配合groupBy和agg后面可以接repartition调整分区数。代码里通常长这样val totalByTag clickDF .groupBy(tag) .agg(count(*).alias(total_downloads)) .orderBy($total_downloads.desc)count(*)与count(1)在Spark里的执行计划基本相同但count(userId)会跳过null所以统计记录条数时一律用count(*)。orderBy在数据量大时会产生全局排序如果只是给前端展示Top20名可以改成sortWithinPartitions减少排序开销。这一阶段最需要确认的是tag字段里有没有空字符串或null否则聚合结果里会出现一个空标签行前端展示时会多出一个「无标签」分类。我经常在跑批前先执行SELECT tag, COUNT(*) FROM clicks GROUP BY tag输出到控制台人工扫一眼数据量大的时候就过滤掉null标签再进下一步。3.3 标签内Top10窗口函数比groupBy加take更稳实现每个标签下的下载量Top10最直接的做法是先按tag, bookId聚合再调用ORDER BY downloads DESC取前10条。但要注意如果你把所有数据先按downloads全局排序再take(10)你得到的是全平台Top10而不是每个标签的Top10。正确做法是用窗口函数在tag分区内编号import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions.{row_number, desc} val rankedDF downloadsPerBook .withColumn(rk, row_number().over( Window.partitionBy(tag).orderBy(desc(downloads)) )) .filter($rk 10)Window.partitionBy(tag)把数据按标签切开在每个分区内按下载量降序编号row_number()生成的编号从1开始且不会重复。这里如果把row_number()换成dense_rank()遇到相同下载量时两个书会并列同一个名次但总数可能超过10换成rank()则名次有跳跃。做榜单时我用row_number()最多因为它能严格保证每本书只有一个名次。这个作业跑完后可以把结果write.mode(overwrite).parquet(output/top10)落到output目录。这里有一个边界需要提醒如果某个标签下电子书数量本来就少于10本rk 10会直接输出全部前端要做空值保护。3.4 写文件时的分区与覆盖策略输出结果时经常有人踩mode(overwrite)和mode(append)的坑。离线跑批推荐overwrite因为每次全量重算覆盖不会产生重复数据如果你做的是增量任务才需要append和partitionBy(dt)。我在这份源码里看到的输出是直接写到了output目录没有做分区。生产环境我一般会按日期分区因为后续排查数据时能快速定位某一天的结果否则整个输出目录越来越大重跑一次全量任务会对下游表造成压力。4. 评分与阅读时段模块平均值聚合和时间段分桶的实现细节4.1 评分模块按bookId聚合取均值之前先做数据质量检查评分需求是「实时更新每本电子书的评分基于读者阅读后所给平均值」。数据在u.data里字段顺序是userId、bookId、rating、ts。实现前一定要先过滤掉评分异常值比如rating为0的记录很多平台用0表示用户跳过评分。过滤后按bookId分组求平均值得到的结果再cast成Decimal类型控制精度val ratingDF spark.read .option(delimiter, |) .csv(src/main/input/u.data) .toDF(userId, bookId, rating, ts) val avgRating ratingDF .filter($rating.cast(double) 0) .groupBy(bookId) .agg( avg($rating.cast(double)).cast(decimal(4,2)).alias(avg_rating), count(*).alias(rating_cnt) )filter($rating.cast(double) 0)这一步很关键因为CSV读取所有字段默认是string类型不cast直接求均值的话Spark会把它当字符串处理或直接报类型不匹配错误。cast(decimal(4,2))把评分限制在小数点后两位比如4.67。rating_cnt用于排序和后续过滤样本量过小的书。这里还要说明一个编程习惯UDF函数在非必要时尽量别用。很多初学者拿到这个需求会立刻写一个UDF去算平均值但avg聚合函数是Spark内置的优化算子能走内置聚合器比UDF快一个量级。只有需要自定义加权、去掉最高最低分这类逻辑时才考虑UDF。评分平均值在平台推送场景里样本量少于一定阈值的评分根本不可信源码文档里如果没有明确阈值我会在聚合后加一个filter($rating_cnt 10)防止个别用户刷高分影响推送策略。4.2 阅读时段模块时间戳转小时再映射到四个时段分桶阅读时段需求的输出是四个时间段的人数0:00-6:00、6:00-12:00、12:00-18:00、18:00-24:00。数据源只能用u.data的ts字段因为click_data.json里如果只有点击时间而缺少阅读时长无法准确刻画在线状态。这里我按「该时间点有阅读行为」代表在线人数来处理这是离线场景最合理的近似口径。实现分桶时我见过两种常见写法一种是先用UDF把小时映射成时段字符串另一种是纯SQL的CASE WHEN。UDF不是不能用但纯表达式更有利于Spark做谓词下推和列剪枝我推荐用when函数链来分组import org.apache.spark.sql.functions.{hour, from_unixtime, when, count} val hourDF ratingDF .withColumn(hour, hour(from_unixtime($ts.cast(long)))) .withColumn(time_bucket, when($hour 0 $hour 6, 0-6) .when($hour 6 $hour 12, 6-12) .when($hour 12 $hour 18, 12-18) .otherwise(18-24)) .filter($hour.isNotNull) val timeStat hourDF .groupBy(time_bucket) .agg(countDistinct(userId).alias(online_users))from_unixtime($ts.cast(long))先把时间戳转成日期时间hour函数提取小时数。四个when分支把小时映射到需求定义的区间。countDistinct(userId)统计的是独立用户数不是点击次数这两个数字对外展示的语义完全不同。如果想统计该时段总阅读次数就用count(*)。这里注意如果u.data的时间戳字段有异常值比如0from_unixtime(0)会得到1970-01-01 08:00:00hour结果是8会被归到6-12这个桶里所以前面的filter($hour.isNotNull)只能过滤掉null过滤不了时间戳为0导致的错误归桶。如果出现大范围时间异常我会做一次filter($ts 1000000000L)把明显非法的时间戳排除掉。4.3 评分模块要不要做成准实时摘要里反复出现「实时统计」这个词但如果拿到的数据是静态的u.data所谓实时只能靠Spark Streaming或Structured Streaming模拟。实际上这份源码跑的是离线批处理所谓实时指的是「流式采集后按批次刷新」。如果你真的需要分钟级更新可以在读取JSON时用readStream配合append模式输出但如果只是课程设计或个人练习批量作业每天跑一次完全够用。我一般会建议先把离线版跑通确实有实时要求再升级成Structured Streaming因为实时的窗口和迟到数据水印会引入大量额外复杂度。5. 避坑指南Spark电子书数据分析常见的五个翻车现场5.1 JSON全字段解析为null现象spark.read.json(src/main/input/click_data.json)之后打印schema全是null类型select出来每行都是null。原因click_data.json文件里每个JSON对象被格式化成多行展示例如{\n userId: 1,\n bookId: 2\n}这种带换行缩进的格式。Spark默认把「一行」当作一条完整JSON对象遇到多行格式时解析失败直接把字段推断为null。解决读取时加上option(multiLine, true)这样Spark读到}结束符才算一条记录结束。我在处理这份源码里的JSON时.option(multiLine, true)就是最关键的一行不加这行整个作业等于白跑。5.2 时间戳被当成字符串或毫秒导致时段统计全乱现象from_unixtime($ts)之后得到的时间是1970年或者时段统计结果全部集中在6-12区间。原因第一种情况是u.data的ts字段以字符串形式被csv读取而from_unixtime接受的参数要求是long类型字符串类型不兼容时Spark返回null第二种情况是ts实际是毫秒时间戳13位直接除以1024会让时间落在1970年附近。解决先$ts.cast(long)强制转类型再判断位数如果ts大于100000000000L13位先除以1000L再传给from_unixtime。我在做的过程中会用SELECT min(ts), max(ts)先探查数据范围确认是秒级还是毫秒级再决定函数链。5.3 标签Top10结果与预期不符现象用clickDF.orderBy($downloads.desc).limit(10)得到的结果全是同一个标签的书其他标签一本都没上榜。原因这是经典的理解偏差你先按全量数据排序再limit(10)得到的是全局Top10而需求是每个标签内部的Top10。两个逻辑在数据量小的时候碰巧看起来差不多数据一多就完全对不上。解决用Window.partitionBy(tag).orderBy(desc(downloads))配合row_number()先分区再排序。我每次写完窗口函数后都会刻意打印每个标签下排名前3的书核对一下确认分区键正确。5.4 Spark作业本地能跑、提交到集群就OOM现象本地IDEA里跑local[*]模式没问题打包提交到YARN集群后执行聚合时报ExecutorLostFailure或OOM。原因本地的executor没有shuffle分区大量数据落盘的瓶颈而集群模式下默认spark.sql.shuffle.partitions是200如果你的数据分布极度倾斜比如「计算机」标签的点击量是「哲学」标签的上百倍某个executor会同时处理大量key导致内存溢出。解决提交时显式设置--conf spark.sql.shuffle.partitions500或更大对倾斜严重的tag做加盐操作把单个大key打散成多个随机后缀key聚合后再合并。我一般先用SELECT tag, COUNT(*) FROM clicks GROUP BY tag看分布如果最大标签占比超过30%就上加盐方案这也是资料里Spark内存调优最常见的一个切入点。5.5 输出文件数量爆炸现象写output目录时生成了几百个小文件每次下游读取都慢得要死而且总在报「too many files」。原因Spark任务里每个partition默认写一个文件前面聚合shuffle把分区数调大刀200之后结果就直接以200个文件写到output目录。解决写结果前用repartition(1)或coalesce(1)合并成单文件如果结果数据集较大可以按tag分区写目录写成output/tagxxx/part-00001.parquet这种形式。这样既保证文件数量可控又方便按分区读取。6. 验证与进阶把四个指标串成一条可监控的Spark作业流当四个模块都单独跑通后我建议把整个分析串成一条脚本化作业流而不是在IDEA里手动点四次运行。最简单的做法是把每个聚合结果都写入output下的独立子目录再用一个shell脚本按顺序执行spark-submit提交。因为四个模块共用同一个SparkSession很多同学图省事把所有逻辑写在同一个main函数里导致中间一步失败就要重跑全部这是非常典型的误用。我会把下载量统计、评分统计、时段统计拆成三个job类按照「数据完整性检查 → job1 → job2 → job3 → 结果校验」的方式编排。验证的环节里最实用的检查是拿两份独立数据源做交叉验证。比如click_data.json里的书籍ID集合和u.data里的书籍ID集合交集数量应该与预期一致每个标签的总下载量之和应该等于click_data.json总记录数。我通常会写一个简单的校验类val expected clickDF.count() val actual totalByTag.agg(sum(total_downloads)).collect()(0).getLong(0) assert(expected actual, stotal mismatch: $expected vs $actual)这里expected是原始点击记录总数actual是分标签汇总后的总数如果两者不等说明分组聚合时出现了null标签或重复记录被遗漏。用断言把校验固话到作业里以后每次跑批都会自动检查而不是等下游报表出来才发现数据对不上。进阶的做法是把这套离线批处理迁到Structured Streaming。因为click_data.json改成点击流数据源后标签下载量和Top10都可以用带水印的窗口聚合实现评分模块则用append模式输出到外部存储。但迁移之前一定要明确批处理里算的是全量累计值流处理里默认算的是窗口内增量值两者的数据口径完全不同不能直接拿流结果去套批处理的报表模型。从那以后我每次跑这类榜单需求前都会强制先写一段数据完整性断言再进聚合逻辑哪怕只是课程设计也照样执行。数据的正确性比算得快重要得多希望帮到你。本文还有配套的精品资源点击获取