简介:面向大数据入门学习者与初级开发者的完整实践合集,涵盖Hadoop电商日志分析、Spark实时流处理、集群搭建教程与数据可视化案例,按入门到实战路径组织,适合希望系统掌握HDFS、MapReduce、Spark Core/Streaming 等核心组件并快速进入真实项目场景的读者。资源共224个文件,压缩包约5.23MB,以Java、Scala源码为主(合计167个),辅以XML配置、Python脚本、HTML/JavaScript可视化页面、Properties配置及SQL、Proto定义、CSV/Data样例数据等,可支撑从代码阅读、环境配置到结果展示的完整闭环。目前已有79人学习下载。通过该资源可获得电商日志离线分析、实时流处理及集群搭建的整套项目代码与配置模板,结合数据可视化案例和ECharts页面,便于边练边学、快速理解大数据技术栈的实际应用方式,是新手构建系统学习路线的实用参考。
1. 先别急着搭集群:这份大数据项目集合到底该从哪下手
拿到压缩包先别急着双击 start-all.sh。解压之后你会看到 ipDatabase.csv、house.csv、u.data、iris.data、echarts.html,外加一个附赠资源.docx。这套东西不是教学 PPT,而是把数据文件当作入口:Hadoop 电商日志分析、Spark 实时流处理、集群搭建教程、数据可视化案例,全围绕这几个数据集展开。如果你已经装好虚拟机、配好 Java,想用真实数据把 HDFS、MapReduce、Spark、ECharts 串起来,这个包会比较顺手。包里还混着几个 .gitignore,说明同一套数据曾被拆到不同项目里维护,正好用来理解多项目结构。我的建议是先按附赠文档把环境过一遍,再按下面的数据链路推进。
2. Hadoop 电商日志分析:从 ipDatabase.csv 到 HDFS 入库与 MapReduce 统计
2.1 为什么先处理 IP 归属地维度表
ipDatabase.csv 在案例里扮演的是 IP 段归属地维度表。真实电商日志通常只记录访问 IP、访问时间、页面 ID、商品 ID、操作类型,不会自带省市信息。要把访问量拆到省份粒度,就必须拿日志里的 IP 去 ipDatabase.csv 里做区间匹配。这个过程放到 MapReduce 里做,就是一次经典的 Reduce 端连接(Reduce Side Join),也是 Hadoop 电商日志分析最常考的知识点。
我先在本地用 Python 确认 CSV 的分隔符、列名和编码:
import csv with open("ipDatabase.csv", encoding="utf-8", errors="replace") as f: reader = csv.reader(f) for i, row in enumerate(reader): if i < 5: print(row) else: break这段脚本做了三件事:读前五行看结构、确认分隔符是不是逗号、暴露出文件编码问题。Windows 导出的 CSV 很多是 GBK,不转码直接传到 HDFS,后面 MapReduce 输出全是乱码占位符。常见做法是统一转成 UTF-8 再 hdfs dfs -put,或者在 Spark 里用 encoding 参数指定。ipDatabase.csv 这类表一般包含起始 IP、结束 IP、国家、省份、城市、运营商几列,注意 IP 要转成整数才能比较大小,字符串比较会出错。
2.2 HDFS 目录规划与数据导入
给项目建目录时,我按数仓分层来组织,避免后续清理和调度时找不到文件:
hdfs dfs -mkdir -p /user/hadoop/warehouse/ods/ip_database hdfs dfs -mkdir -p /user/hadoop/warehouse/ods/access_log hdfs dfs -mkdir -p /user/hadoop/warehouse/app/ip_analysis hdfs dfs -put ./ipDatabase.csv /user/hadoop/warehouse/ods/ip_database/ hdfs dfs -put ./access.log /user/hadoop/warehouse/ods/access_log/参数说明:mkdir -p 会递归创建完整路径;put 后面第一个参数是本地路径,第二个是 HDFS 路径。目录里的 ods 表示原始数据层,app 表示应用结果层。伪分布式环境里 HDFS 默认副本数是 3,但单节点 DataNode 实际上只有 1 份副本,put 操作会因为复制副本不到位而一直等待,最终抛出写文件超时。所以要么在 hdfs-site.xml 里把 dfs.replication 改成 1,要么先确认 datanode 进程已经正常启动。配置如下:
<property> <name>dfs.replication</name> <value>1</value> </property>改完配置要重启 HDFS,或者执行 hdfs dfsadmin -refreshNodes 让配置生效。新手经常在这里卡住,以为是网络问题,实际只是副本数没有按单机环境调整。
2.3 MapReduce 统计 IP 地域分布
统计各省访问量是整套 Hadoop 案例的骨架。用 Java 写完整代码会比较长,这里用 Hadoop Streaming 加 Python 演示思路更直观。先看 Mapper:
#!/usr/bin/env python import sys for line in sys.stdin: fields = line.strip().split(",") if len(fields) < 2: continue src_ip = fields[0].strip() try: # 将 IPv4 转成整数,用于后续区间判断 parts = src_ip.split(".") ip_num = (int(parts[0]) << 24) + (int(parts[1]) << 16) \ + (int(parts[2]) << 8) + int(parts[3]) except Exception: continue print("ip_num:%d" % ip_num)这段 Mapper 只做清洗和 IP 转整数,真正的区间匹配放在 Reducer 里:Reducer 启动时会把 ipDatabase.csv 加载到内存,将日志 IP 逐条二分查找,命中后输出省份和计数。提交命令如下:
hadoop jar /opt/hadoop/share/hadoop/tools/lib/hadoop-streaming-*.jar \ -files mapper.py,reducer.py,ipDatabase.csv \ -mapper "python3 mapper.py" \ -reducer "python3 reducer.py" \ -input /user/hadoop/warehouse/ods/access_log/access.log \ -output /user/hadoop/warehouse/app/ip_analysis/result参数说明:-files 会把本地脚本和维度表打包到 DistributedCache,每个 container 的工作目录里都能读到这些文件;-mapper 和 -reducer 指定解释器和脚本;-input 和 -output 必须是 HDFS 路径,且 output 目录必须不存在。这里把 ipDatabase.csv 直接交给 Reducer 读,等于把 join 挪到了内存里。维度表小于 200MB 时这种方案简单有效,超过这个量级就要用 MapFile 或 HBase 做维度查询,不要硬塞内存。
2.4 伪分布式下任务失败先看这几个地方
任务结束不代表结果正确。我先看 failed 任务的日志,再检查 part-r-00000 前几行,最后确认没有因为 update 模式覆盖掉旧结果。三个高频问题值得记录:
- 输出目录已存在:任务直接抛 FileAlreadyExistsException,删掉再跑。
- Mapper 字段越界:先对 access.log 执行 head -5,确认列数和你代码里取的下标一致。
- Reducer 端内存溢出:观察日志里的 GC overhead limit exceeded,调高 mapreduce.reduce.memory.mb,同时注意 mapreduce.reduce.java.opts 要一起调。
常见错误对照表:
| 报错关键字 | 原因 | 处理方式 |
|---|---|---|
| FileAlreadyExistsException | 输出目录已存在 | hdfs dfs -rm -r 输出目录 |
| Input path does not exist | 输入路径不存在 | 核对目录与文件上传状态 |
| GC overhead limit exceeded | Reducer 堆内存不足 | 同时提高 memory.mb 和 java.opts |
| Incompatible clusterIDs | NameNode 格式化两次 | 清理 DataNode 数据目录后重新格式化 |
到这里,Hadoop 这条线就能闭环了。下一章把 Spark 接进来,目录正好和 HDFS 共用,不需要二次导入。
3. Spark 实时流处理:用 u.data 复现评分流的批与流
3.1 u.data 不是日志,但很适合演流
u.data 是 MovieLens 经典的“用户 ID-电影 ID-评分-时间戳”四列数据,字段只有 4 个,类型清晰,最适合模拟实时评分场景。你拿到的 u.data 是静态文件,但 Spark 的 Structured Streaming 可以监听目录、读取新文件,把静态数据拆成多个小文件放进去,就能伪造出“用户不断打分”的连续流。这章的思路是先讲清楚流处理和批处理的差异,再落到能改参数的代码上。
3.2 用 Structured Streaming 监听评分目录
先确认 u.data 的分隔符是 \t,再写读取逻辑:
from pyspark.sql import SparkSession from pyspark.sql.types import (StructType, StructField, IntegerType, LongType) schema = StructType([ StructField("userId", IntegerType(), True), StructField("movieId", IntegerType(), True), StructField("rating", IntegerType(), True), StructField("timestamp", LongType(), True) ]) spark = SparkSession.builder \ .appName("rating_stream") \ .master("local[2]") \ .getOrCreate() lines = spark.readStream \ .format("csv") \ .schema(schema) \ .option("sep", "\t") \ .load("/tmp/rating_input")参数说明:master local[2] 至少要给两个线程,一个接收数据、一个处理数据,本地只给 1 个会导致任务不输出;sep 指定 Tab 分隔;schema 里的 timestamp 定义成 LongType,后续才能用 from_unixtime 转换;load 的路径是待监控目录,不是单个文件。这里你不需要事先把文件放到目录里,只要保证目录存在。启动后再复制文件进去就可以看到流式输出。
3.3 用窗口统计给实时热门电影排队
流处理里最常用的需求是滑动窗口内统计。下面这段代码每 5 秒输出一次过去 10 秒内被评分次数最多的电影:
from pyspark.sql.functions import window, count, from_unixtime ratings = lines.withColumn( "ts", from_unixtime("timestamp").cast("timestamp") ) hot = ratings.groupBy( window("ts", "10 seconds", "5 seconds"), "movieId" ).agg(count("rating").alias("cnt")) query = hot.writeStream \ .outputMode("complete") \ .format("console") \ .option("truncate", "false") \ .start() query.awaitTermination()逻辑说明:先把时间戳转成 Timestamp 类型,再交给 window 函数。window 参数第一个是窗口长度 10 秒,第二个是滑动间隔 5 秒,合起来就是“每 5 秒滑动一次,计算最近 10 秒的窗口”。outputMode 用 complete,表示每次都输出全量聚合结果,这样 orderBy 才能作用于全局。如果改成 append 模式,只能输出新增行,groupBy 的聚合值会不完整。这是流处理新手最容易混淆的点。
3.4 资源参数与隐藏的数据质量问题
本地提交流任务时,我习惯显式指定内存和分区数,避免默认值拖慢速度:
spark-submit \ --master local[2] \ --driver-memory 2g \ --executor-memory 2g \ --conf spark.sql.shuffle.partitions=4 \ rating_stream.pyspark.sql.shuffle.partitions 默认是 200,本地小数据集用 200 个 shuffle 分区会产生大量空文件,改成 4 更贴合入门场景。另一个隐藏问题是 CSV 解析时的类型转换:schema 里 rating 是 IntegerType,遇到非数字数据时整行会被置成 null,统计结果在没人察觉的情况下变少。我一般的做法是先把 rating 读成 StringType,过滤掉非数字行后再 cast 成 IntegerType,这样异常数据能被看到,而不是被静默吞掉。u.data 本身很干净,但你换到真实用户行为日志时,这一步就是必踩的坑。
4. 集群搭建教程的关键细节:从 Hadoop 伪分布式到 Spark on YARN
4.1 版本选型要放到环境之后考虑
很多资源包里的教程还在用 Hadoop 2.7.3 配 Spark 2.4.0,但 JDK 版本一变就启动失败。我建议优先采用 Hadoop 3.3.x、Spark 3.x、JDK 8 或 11 的组合。选型依据是附赠资源.docx 里是否写明版本矩阵,没写的话按 CDP 或 HDP 的兼容列表来。注意 Spark 3.2 以上虽然支持 Java 17,但 Hadoop 官方对 Java 17 的支持还比较保守,混用容易在 NameNode 启动时出现 UnsupportedClassVersionError,所以别为了追新而把 JDK 拉太高。
4.2 core-site.xml 和 hdfs-site.xml 的最小配置
不管单机还是三节点,核心配置就几项。给一个可以照抄的最小集合:
<property> <name>fs.defaultFS</name> <value>hdfs://node01:9000</value> </property> <property> <name>hadoop.tmp.dir</name> <value>/data/hadoop/tmp</value> </property>这是 core-site.xml。hdfs-site.xml 需要把 NameNode 和 DataNode 的数据目录分开:
<property> <name>dfs.namenode.name.dir</name> <value>file:///data/hadoop/namenode</value> </property> <property> <name>dfs.datanode.data.dir</name> <value>file:///data/hadoop/datanode</value> </property> <property> <name>dfs.replication</name> <value>2</value> </property>参数说明:fs.defaultFS 决定了客户端访问 HDFS 的入口;hadoop.tmp.dir 如果不显式配置,默认落在 /tmp 目录,系统重启后元数据会丢得干干净净,这是新手遇到“重启后 HDFS 起不来”的根源。dfs.namenode.name.dir 和 dfs.datanode.data.dir 必须指向不同目录,否则格式化时会把元数据和数据块混在一起,DataNode 启动后集群 ID 对不上。dfs.replication 在三节点集群配 2,伪分布式配 1。
4.3 启动顺序、格式化和租约恢复
搭建教程里最坑的是顺序问题。正确流程是先修改配置,再执行一次 hdfs namenode -format,然后 start-dfs.sh,用 jps 检查进程,最后 start-yarn.sh。格式化命令只能成功执行一次,第二次格式化会把 NameNode 的 clusterID 换掉,DataNode 还带着旧的 clusterID,启动日志里就会出现 Incompatible clusterIDs。解决办法是清空 dfs.datanode.data.dir 和 dfs.namenode.name.dir 里的内容,统一再格式化,不要只删一边。
注意:hdfs namenode -format 只能执行一次,重复格式化会引发 Incompatible clusterIDs,需要清理数据目录后重新初始化。
HDFS 写文件失败也是高频问题,尤其是网络上常搜到的 previous writer likely failed to write,这是因为旧写操作的租约没有过期,新的写请求拿不到文件锁。缓解办法是找到对应的文件路径后执行:
hdfs debug recoverLease -path /user/hadoop/warehouse/ods/access_log/access.log -retries 3这个命令会让 NameNode 主动恢复文件租约,把没有完成写入的文件标记为可继续写。注意它只对关闭状态的文件生效,如果文件还在正常写入中,不要用这个命令打断。
4.4 Spark on YARN 提交参数与日志查看方式
Spark 任务要跑在 YARN 上,提交命令里最关键的是 --master yarn 和 --deploy-mode。我给一个常见配置:
spark-submit \ --master yarn \ --deploy-mode cluster \ --driver-memory 2g \ --num-executors 3 \ --executor-cores 2 \ --executor-memory 4g \ --class com.example.RatingAnalyzer \ rating.jar参数说明表:
| 参数 | 示例 | 作用 |
|---|---|---|
| --deploy-mode cluster | driver 在集群内运行 | 提交机退出不影响任务 |
| --deploy-mode client | driver 留在本地 | 适合调试,会占用提交机资源 |
| --num-executors 3 | 启动 3 个执行器 | 根据队列资源调整 |
| --executor-cores 2 | 每个执行器用 2 核 | 避免申请超过 YARN 容器上限 |
| --executor-memory 4g | 每个执行器 4GB | 堆内存,注意要留 off-heap 空间 |
cluster 模式下 driver 日志不在提交机,要去 YARN ResourceManager 页面或执行 yarn logs -applicationId 应用ID 查看。我用 client 模式调试时,driver 日志直接打到终端,但任务挂在 Session 上,终端断开任务就被杀。生产环境跑流任务和长任务建议用 cluster 模式,日志统一由 YARN 收集,排错也简单。
4.5 副本数与数据目录对 HDFS 写入的影响
单机伪分布式最常见的报错是“could only be written to 0 of 1 minReplication nodes”,本质是副本需求大于 DataNode 实际副本数。解决办法一行:
hdfs dfs -setrep -R 1 /user/hadoop/warehouse把目录下所有文件副本数降为 1。注意 setrep 只改变已有文件副本数,新文件仍由 dfs.replication 控制,所以更彻底的方案还是改 hdfs-site.xml。三节点集群则要检查 DataNode 是否都上线,hdfs dfsadmin -report 能列出每个 DataNode 的状态。集群搭建不是跑通 start-all.sh 就结束,数据目录、租约、副本数都是回头要查的点。
5. ECharts 数据可视化:用 house.csv 和 iris.data 做可交互的图表
5.1 从 CSV 到 ECharts 的 JSON 转换
echarts.html 是项目集合的最后环。iris.data 是典型的四维特征数据,适合做散点图;house.csv 包含面积、价格、地段,适合展示房价分布。不要把 CSV 直接塞给 ECharts,先转成 JSON 数组:
import pandas as pd import json df = pd.read_csv("house.csv", encoding="utf-8") data = df[["area", "price", "district"]].dropna().to_dict(orient="records") with open("house.json", "w", encoding="utf-8") as f: json.dump(data, f, ensure_ascii=False)dropna() 会把缺失行整行丢弃,to_dict(orient="records") 把 DataFrame 转成 [{area: 89, price: 420}, ...] 这种结构。ensure_ascii=False 必须保留,否则中文 district 会被转成 \uXXXX,ECharts 显示时还要多一步解码。
5.2 用 dataset 组件把数据与坐标轴解耦
ECharts 5 里最推荐的方式是在 option 里配置 dataset,然后用 encode 映射列到轴:
<div id="chart" style="width: 100%; height: 600px;"></div> <script src="https://cdn.jsdelivr.net/npm/echarts@5/dist/echarts.min.js"></script> <script> var chart = echarts.init(document.getElementById("chart")); chart.setOption({ dataset: { source: [ ["area", "price"], [89, 420], [120, 680], [145, 830] ] }, xAxis: { type: "value", name: "面积(m²)" }, yAxis: { type: "value", name: "价格(万)" }, series: [{ type: "scatter", encode: { x: "area", y: "price" } }] }); </script>dataset.source 是二维数组,第一行是列名,encode 里直接引用列名做映射。这样做的好处是切换显示字段时只改 encode,不碰 series 类型。如果 house.json 在本地,要用 python -m http.server 8888 起静态服务,直接双击 HTML 时 file 协议会拦截本地 Ajax 请求。
5.3 用 setOption 合并模式做动态字段切换
iris.data 有四个特征,可视化时要让用户自己选两维映射到 x 和 y。用 select 控件绑定 onchange 事件,让用户切换 x 轴字段:
document.getElementById("x-select").onchange = function () { chart.setOption({ series: [{ encode: { x: this.value, y: currentY } }] }); };setOption 默认做合并,不是整体替换,所以 x 轴字段可以单独更新,其他配置比如网格、图例都保持不变。如果有多个 series,一定要指定 seriesIndex 或 seriesId,否则合并更新会作用到所有序列上。这个技巧比重新 init 一个 chart 实例轻量得多,也不会丢失缩放状态,实际做可视化大屏时是最高频的操作之一。
本文还有配套的精品资源,点击获取