ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

大数据学习路径:Hadoop离线分析到Spark实时流处理与可视化

大数据学习路径:Hadoop离线分析到Spark实时流处理与可视化 简介一套面向入门学习者的大数据技术学习与实践项目集合涵盖Hadoop电商日志分析、Spark实时流处理、集群搭建与数据可视化案例构成从零到实战的完整学习路径。借助真实项目代码与配置新手可快速理解分布式存储、离线计算与实时计算的核心方法。包内共224个文件以Java和Scala源码为主分别是Hadoop与Spark项目核心实现XML和properties文件用于项目配置与部署Python脚本辅助数据预处理HTML与JavaScript用于可视化界面CSV和data文件为测试数据集SQL负责数据查询proto文件用于数据序列化。压缩包仅5.23MB轻量便捷。目前已有80人学习下载。整体内容按入门到实战组织学习者能理清大数据项目目录结构掌握从数据采集、HDFS存储到离线分析、实时计算、可视化展示的完整开发链路。附带的多个数据集与文档可支撑课程设计、自学实践与面试复习是一份实用的大数据入门宝典。1. 大数据技术学习与实践项目集合一条从离线到实时的完整学习路径“大数据技术学习与实践项目集合”这个标题代表的不只是一份压缩包而是一条链条Hadoop电商日志分析让你理解离线批处理怎么做Spark实时流处理补齐实时链路集群搭建教程负责把环境弄起来数据可视化案例则把结果变成别人能看懂的作品。零基础、或者写过单体应用但没碰过分布式的开发者最容易卡在“书看完了却不知道从哪下手”环境起不来、日志不会洗、任务一跑就OOM。这篇笔记按这条路径把每一步的关键命令、参数和坑拆开照着做可以跑通一个完整的大数据项目闭环。2. Hadoop伪分布式搭建与HDFS常用命令先把读写流程跑通2.1 为什么新手先从伪分布式开始而不是直接搭三节点集群很多教程上来就让你准备三台虚拟机配免密、配NTP、改hosts。这套流程对理解“集群”是有用的但对一个刚接触Hadoop的人来说环境问题会直接把学习热情耗光网络不通、内存不够、各节点进程起不来最后连JobTracker日志都看不懂。伪分布式的意思是所有角色都跑在本机NameNode、DataNode、ResourceManager、NodeManager各占一个进程大家共用同一份本地磁盘。它和真实集群的读写路径完全一致——客户端还是先问NameNode拿元数据再写DataNode只是DataNode只有一份副本数必须设成1。先在这个模式下跑通HDFS读写、跑通一个MapReduce或Spark任务再考虑多节点集群和HA。后面要上高可用、和ZooKeeper整合的实战那是环境已经熟练之后的事。2.2 Hadoop伪分布式搭建四个配置文件和一条格式化命令以Hadoop 3.x为例先装JDK8并配好JAVA_HOME下载解压Hadoop后需要改四个配置文件都在$HADOOP_HOME/etc/hadoop/下。这步是新手第一个翻车点少配一个文件后面就跑不起来。第一个是core-site.xml指定NameNode地址和临时目录configuration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property property namehadoop.tmp.dir/name value/data/hadoop/tmp/value /property /configurationfs.defaultFS是HDFS的统一入口客户端和NameNode通信都靠它hadoop.tmp.dir是元数据和数据块的公共根目录生产上会用多块磁盘分开写学习阶段放在一块盘上没问题。第二个是hdfs-site.xml设置副本数和NameNode/DataNode各自的存储目录configuration property namedfs.replication/name value1/value /property property namedfs.namenode.name.dir/name value/data/hadoop/name/value /property property namedfs.datanode.data.dir/name value/data/hadoop/data/value /property /configuration副本数设1是因为真只有一个DataNode设了3反而会一直安全副本缺失告警。生产环境这两个目录要分盘放系统盘崩了数据还在。剩下的yarn-site.xml和mapred-site.xml分别是资源调度和MapReduce运行框架的配置伪分布式里只需把mapreduce.framework.name设为yarn。配好后执行格式化命令注意这条命令只有第一次需要执行# 本机免密登录为ssh启动脚本做准备 ssh-keygen -t rsa -P -f ~/.ssh/id_rsa cat ~/.ssh/id_rsa.pub ~/.ssh/authorized_keys chmod 600 ~/.ssh/authorized_keys # 格式化NameNode重复执行会导致DataNode起不来 hdfs namenode -format # 启动HDFS和YARN start-dfs.sh start-yarn.sh # 查看进程是否齐全 jps格式化本质是生成空的元数据和clusterID。这里最需要注意的是“只格式一次”如果后面因为配置修改想重新初始化也要先清掉DataNode的current目录再重来否则DataNode会报clusterID不一致这个坑在第5章展开讲。jps能看到NameNode、DataNode、SecondaryNameNode、ResourceManager、NodeManager五个进程才算成功。2.3 HDFS读写流程与新手必须会的6条命令HDFS写文件时客户端调用FileSystem.create联系NameNodeNameNode在元数据里创建文件并返回可写入的DataNode列表客户端按块大小默认128MB把数据切成块以流水线方式依次发给第一个DataNode再逐级复制到副本节点整个块写完DataNode返回确认。读文件时客户端先问NameNode要文件分布在哪些块、在哪些DataNode上NameNode按网络拓扑排序返回客户端选最近的DataNode建立连接直接读同时做checksum校验。这个“元数据走NameNode、数据直连DataNode”的设计是理解后面所有HDFS故障排查的基础。日常项目里最常用的命令就这几条# 建目录并上传日志 hdfs dfs -mkdir -p /data/log hdfs dfs -put access.log /data/log/ # 查看文件列表和块分布 hdfs dfs -ls /data/log hdfs dfs -du -h /data/log # 查看文件内容tail适合看实时写入的日志 hdfs dfs -cat /data/log/access.log | head hdfs dfs -tail /data/log/access.log # 删除目录注意加上递归参数 hdfs dfs -rm -r /data/log另外还有一条hdfs dfsadmin -report能直接看到DataNode的容量和健康状况。刚搭建完一定要跑一次确认DataNode状态不是In Service而是正常可用。上传文件后用-du -h看大小和本地文件比对能验证副本和块的实际存储情况。这些命令在后面的电商日志分析里会反复用到建议单独开一个终端窗口反复敲熟。3. Hadoop电商日志分析从原始日志到PV、UV、TopN指标3.1 离线分析选型MapReduce、Hive还是Spark SQL电商日志分析是这个项目集合的核心场景数据通常是Nginx打出来的access log一行就是一次访问记录。要对这类数据算PV、UV、TopN商品有三种常见做法原生MapReduce、Hive、Spark SQL。方案学习成本执行效率适合场景原生MapReduce高代码模板繁琐慢中间结果落盘入门理解Map-Shuffle-Reduce原理Hive低写SQL即可中底层还是MR或Tez纯离线报表、团队里SQL多Spark SQL中API接近SQL快内存计算复杂ETL、需要快速迭代的分析我的建议是两条腿走先用一个单词计数的MapReduce把“MapShuffleReduce”三阶段跑通理解框架在做什么真正做日志分析直接上Spark SQL因为学习路径最短而且Spark的DataFrame API和后面实时流处理是同一套学一次用两处。头歌这类实训平台上的“实验5 HDFS和MapReduce综合实训”练的就是前半段很多项目集合里的电商日志分析案例其实就是把这段实训换成了Spark实现。3.2 日志清洗用正则和PySpark把一行日志拆成结构化字段原始日志长这样192.168.1.1 - - [10/Oct/2023:13:55:36 0800] GET /product/123 HTTP/1.1 200 4321 https://www.example.com/ Mozilla/5.0统计前必须先清洗过滤空行和爬虫、按空格切分、把时间格式改掉。用PySpark读文本再用正则提取字段from pyspark.sql import SparkSession from pyspark.sql.functions import regexp_extract, col spark SparkSession.builder \ .appName(log-parse) \ .master(local[*]) \ .getOrCreate() # 从HDFS读原始文本一行一条记录 raw spark.read.text(hdfs://localhost:9000/data/log/access.log) # 正则按顺序提取ip、日期、请求方式、URL、状态码 log_pattern r^(\S)\s\S\s\S\s\[([^\]])\]\s(\S)\s(\S)\s\S\s(\d{3}) parsed raw.select( regexp_extract(raw.value, log_pattern, 1).alias(ip), regexp_extract(raw.value, log_pattern, 2).alias(time_str), regexp_extract(raw.value, log_pattern, 3).alias(method), regexp_extract(raw.value, log_pattern, 4).alias(url), regexp_extract(raw.value, log_pattern, 5).cast(int).alias(status) ) # 过滤掉没匹配上的行状态码非200的也算有效但单独标记 clean parsed.filter(col(ip) ! ) clean.show(10, truncateFalse)regexp_extract的第三个参数是取正则里第几个括号的组匹配不上返回空串状态码直接cast成int后面按状态码分组就方便了。清洗这一步最容易忽视的是脏数据量生产日志里大约有5%~10%的行格式不完整比如请求行末尾缺引号、状态码是四位数。如果过滤条件写得太严格会把正常行也删掉写得太松统计结果又明显偏小。建议先count一下总数再和清洗后的行数对比做到心里有数。3.3 指标计算落地PV、UV和TopN的代码与结果输出清洗完就是三个经典指标。PV最简单一行就是一次访问分组计数即可UV要去重思路是精确还是近似TopN则是分组后按次数排序取前10。这里给出一套可直接跑的代码from pyspark.sql.functions import count, approx_count_distinct, desc, to_date # 先把时间字符串转成日期方便按天统计 clean clean.withColumn(date, to_date(col(time_str), dd/MMM/yyyy:HH:mm:ss)) # PV按天统计访问次数 pv clean.groupBy(date).agg(count(*).alias(pv)) # UV按天统计去重IP数 uv clean.groupBy(date) \ .agg(approx_count_distinct(ip).alias(uv)) # TopN统计被访问次数最多的10个URL topn clean.groupBy(url) \ .agg(count(*).alias(pv)) \ .orderBy(desc(pv)) \ .limit(10) # 结果以parquet格式写回HDFS供可视化阶段读取 uv.write.mode(overwrite).parquet(hdfs://localhost:9000/result/uv)approx_count_distinct基于HyperLogLog算法内存占用固定误差在1%~2%左右如果要求精确去重需要repartition后把IP收敛到单个分区再countDistinct数据量大了会明显变慢。实际生产里UV一般都用近似值因为产品看的是量级和趋势不是精确个位数。结果写成parquet而不是CSV是因为列式存储体积小、带schema后续Spark读回来自动识别字段比文本格式高效得多。4. Spark实时流处理用Structured Streaming接Kafka做窗口统计4.1 实时链路选型DStream还是Structured Streaming离线日志分析跑完一批是一次性的实时流处理则是数据一边进一边算。常见的链路是业务日志写进KafkaSpark实时消费聚合结果落地到MySQL或HDFS最后可视化展示。Spark里有两套流API老一点的Spark Streaming基于DStream把流切成一批批RDD再做处理新一点的Structured Streaming基于DataFrameAPI和Spark SQL几乎一样语义上也更强。新项目不要犹豫直接选Structured Streaming。原因有三个第一它处理JSON更方便from_json直接解析成结构化字段第二它原生支持exactly-once语义配合Kafka偏移量自己管理重启不容易重复算第三它的窗口和水位线API更直观和离线代码能共用DataFrame逻辑。在很多“Spark实战”案例和近期热门的“spark中读取json”场景里能找到的代码基本都是这个方向。4.2 读Kafka数据源的最小可运行代码先确保本地有Kafka并能生产数据然后提交这个Spark作业。这段代码从Kafka订阅order-topic把value解析成JSON按商品ID统计订单数from pyspark.sql import SparkSession from pyspark.sql.functions import from_json, col from pyspark.sql.types import StructType, StructField, StringType, LongType spark SparkSession.builder \ .appName(kafka-streaming) \ .master(local[2]) \ .getOrCreate() # 定义JSON的schema字段类型必须和实际数据严格匹配 schema StructType([ StructField(order_id, StringType()), StructField(product_id, StringType()), StructField(amount, LongType()), StructField(ts, LongType()) ]) # readStream从Kafka拉数据这是流处理的入口 stream spark.readStream \ .format(kafka) \ .option(kafka.bootstrap.servers, localhost:9092) \ .option(subscribe, order-topic) \ .option(startingOffsets, earliest) \ .option(failOnDataLoss, false) \ .load() # Kafka的value是字节数组先转string再按schema解析 parsed stream.select( from_json(col(value).cast(string), schema).alias(data) ).select(data.*) # 按商品ID累计订单数 result parsed.groupBy(product_id) \ .agg(count(*).alias(cnt)) # 每10秒把结果输出到控制台 query result.writeStream \ .outputMode(update) \ .format(console) \ .trigger(processingTime10 seconds) \ .option(checkpointLocation, /data/checkpoint/order-topic) \ .start() query.awaitTermination()startingOffsets设为earliest表示从最早的未消费消息开始调试最方便生产环境一般用latest只处理新数据。failOnDataLoss默认true意思是Kafka记录被删或偏移量不可用时任务直接失败开发阶段改成false能避免因为手动删过topic导致起不来。checkpointLocation必须单独用一个目录流重启后靠它恢复偏移量和状态路径不要和别的工作共用否则状态错乱。4.3 窗口计算、水位线与checkpoint参数设置电商场景里“最近5分钟每件商品的订单量”“最近1小时UV”都需要窗口聚合。Structured Streaming的窗口有两种滚动窗口固定长度互不重叠适合统计每分钟订单量滑动窗口固定长度加上滑动步长适合看每5分钟更新一次的最近10分钟趋势。代码里加一个window函数和watermarkfrom pyspark.sql.functions import window # ts是毫秒时间戳转成TimestampType parsed parsed.withColumn(event_time, (col(ts) / 1000).cast(timestamp)) # 10分钟窗口5分钟滑动一次允许事件迟到1分钟 windowed parsed \ .withWatermark(event_time, 1 minute) \ .groupBy( window(col(event_time), 10 minutes, 5 minutes), col(product_id) ) \ .agg(count(*).alias(cnt))withWatermark允许一定延迟的数据晚到窗口关闭后迟到超过水位线的数据会被丢弃这比Spark Streaming的粗粒度窗口灵活得多。两个参数最容易调坏一是窗口长度太短导致结果碎片化二是水位线太小导致晚到数据全丢。调试阶段建议把processingTime设为“5 seconds”窗口设为“10 minutes, 5 minutes”这样几十秒内就能看到多批结果确认逻辑没问题再调回生产参数。另外流任务重启后如果修改了聚合逻辑旧checkpoint里的状态和新代码对不上任务会一直报错这时需要删掉checkpoint目录重跑这点在文档里写得很隐晦属于典型血泪经验。5. 大数据项目避坑指南5个高频故障的现象、原因与解决5.1 现象NameNode格式化了两次DataNode全部起不来跑伪分布式或者多节点集群时进程日志里频繁出现Incompatible clusterIDsDataNode启动后立刻退出。原因在于hdfs namenode -format是允许重复执行的每次格式化都会生成一个新的clusterID而DataNode的存储目录里还是第一次格式化时的ID。两边ID对不上DataNode认为自己在向一个陌生的NameNode上报块直接拒绝注册。解决方法是先看$dfs.namenode.name.dir/current/VERSION里的clusterID再对比DataNode的current/VERSION确认不一致后把DataNode存储目录下的文件全部删掉重新启动DataNode让它重新向NameNode注册。学习环境操作最快但不建议再执行一次-format那会连NameNode的元数据一起重置原本上传的数据全没了。格式化命令加了-force也一样修复的是逻辑不是数据这个操作没有后悔药。5.2 现象Spark任务Container OOM调大内存也没用一个Spark批处理任务跑了几分钟日志里出现Container killed on exceeding memory limits就算把spark.executor.memory从1G调到4G照挂不误。原因往往是两个一是默认Java序列化把对象膨胀了3到5倍大量内存花在序列化副本上二是某个分区数据倾斜单task处理的数据远超平均。解决分两步。先改序列化在spark-submit或Session配置里加spark.serializerorg.apache.spark.serializer.KryoSerializer注意还要注册需要用到的类否则退化成反射序列化没省多少。再看Spark UI里每个task处理的数据量如果有task明显高于其他对key做加盐分散或先repartition再聚合。内存参数也不只是executor.memory一项spark.memory.offHeap.enabled和spark.executor.memoryOverhead在容器环境下同样影响OOM判断。5.3 现象ECharts图表空白后端接口返回的类型不对用Flask写API给前端ECharts用浏览器请求接口返回的JSON看起来完全正常但页面只有坐标轴没有柱子和线条控制台也没有报错。这类问题排查时是典型的“看起来都对了”。原因多半是字段类型没对齐ECharts的xAxis要求字符串数组yAxis的data要求数值数组而PyMySQL查出的pv字段以字符串返回前端拿123这种字符串画柱状图ECharts的value轴会认成非法值直接跳过。还有日期字段不一致、接口返回了null值也会导致图形空白。解决方法是后端在返回前统一int(pv)、str(date)None替换成0前端做一层Number(d.pv)兜底。注意在可视化案例里清洗和统计阶段为了查得动把字段设成string省事这个偷懒会一路传到图表最好在接口层统一消毒。5.4 现象集群节点时间不同步任务随机失败像玄学三个节点的集群任务有时候成功有时候失败日志里或明或暗出现clock skew、lease expired甚至ZooKeeper连接频繁丢失。排查代码、网络、参数都无果最后发现各节点系统时间差了十几分钟。HDFS和ZooKeeper对时间差都有容忍上限节点时间偏差过大NameNode会拒绝处理来自“未来”或“过去”节点的块报告租约也会提前过期表现就是任务随机失败。解决方法是给每个节点配NTP或chrony同步测试环境至少也要手动校准还要注意虚拟机一键克隆出来的多个节点克隆时的快照时间残留会让时间差更明显。这个问题和代码质量无关遇到偶发跨节点故障一定先检查时间别先动参数。5.5 现象distcp迁移大批量数据时报checksum错误用hadoop distcp把一份历史日志从旧集群搬进新集群跑到一半报ChecksumException。原因有两类跨Hadoop版本迁移时块校验和算法不一致或者源文件在拷贝过程中被上游任务修改第二次读取时checksum对不上。解决方案是按优先级加参数-skipcrccheck跳过校验和对比适合迁移纯历史数据、已确认内容完整的场景如果目标目录可能已有部分文件加-update只同步新增或更新的文件只想看差异不实际拷贝用-diff先生成差异列表。distcp本身是MapReduce作业跑大批量时还要关注-m参数控制的map数默认20个map在PB级数据量下容易超时生产环境一般按数据量手动推算map数而不是一路默认跑到底。6. ECharts数据可视化案例用Flask把分析结果做成可演示作品6.1 最小可视化架构Flask提供JSON、前端ECharts渲染常见的数据可视化案例比如网约车订单分析、农产品价格可视化结构其实都一样Spark分析结果落到MySQL或HDFS后端提供一个JSON接口前端用ECharts渲染。最小实现只要两个文件。后端是Flask接口从MySQL查出每天PV返回JSONfrom flask import Flask, jsonify import pymysql app Flask(__name__) app.route(/api/pv) def pv(): conn pymysql.connect(hostlocalhost, userroot, password123456, databasebigdata) cur conn.cursor() cur.execute(select date, pv from daily_pv order by date) return jsonify([{date: str(r[0]), pv: int(r[1])} for r in cur.fetchall()]) if __name__ __main__: app.run(port5000)前端页面用fetch拉接口再交给ECharts画柱状图关键是把接口里的字段对上body div idchart stylewidth:800px;height:400px;/div script srchttps://cdn.jsdelivr.net/npm/echarts5/script script fetch(/api/pv) .then(res res.json()) .then(data { const chart echarts.init(document.getElementById(chart)); chart.setOption({ xAxis: { type: category, data: data.map(d d.date) }, yAxis: { type: value }, series: [{ type: bar, data: data.map(d d.pv) }] }); }); /script /body6.2 验证学习成果的三个自测方法第一个自测是“删掉重来”把checkpoint目录删掉、把HDFS上的结果目录清空从原始日志重新跑一遍全链路。能不看笔记把环境、清洗、统计、接口串起来才算真会了。第二个是“换一份数据”把电商access.log换成自己生成的一批订单JSON或者换一个数据集结构不变但业务字段不同逼着你改schema、改清洗规则比照着教程敲代码收获大得多。第三个是“写实验手记”每次操作把命令、参数、异常日志原样贴在笔记里旁边标注这次为什么改了这个值。这条链路我也反复走过最早的教训是只管跑通不记录参数翻车之后翻文档找原因时间浪费比多刷两个教程都多后来养成边跑边记的习惯效率提升很明显。希望帮到你。本文还有配套的精品资源点击获取
返回列表