ARTICLE DETAIL

资讯详情

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

从Hadoop到Spark:大数据离线分析到实时流处理完整实战路径

从Hadoop到Spark:大数据离线分析到实时流处理完整实战路径 简介一份面向大数据初学者的一站式学习与实践项目合集覆盖Hadoop生态核心组件与完整学习路径从集群搭建到电商日志分析、Spark实时流处理及数据可视化均有落地代码与案例。资源包共224个文件约5.23MB以Java与Scala源码为主Hadoop/Spark开发辅以Python脚本、XML配置、Properties配置文件、HTML/JS可视化页面、SQL及CSV/data数据集等基本覆盖大数据项目开发各环节。已有79人浏览学习适合希望快速建立大数据知识体系并动手实践的新手参考。资源内包含多个可直接运行的实验项目如电商日志分析、Spark实时处理演示、数据可视化大屏等并配有集群搭建教程与多种数据集可帮助读者理解HDFS存储、MapReduce计算、Spark流处理等关键概念。1. 从一份“大数据学习包”开始先搞懂Hadoop和Spark分别解决什么问题很多新手会去下载一个名为“大数据技术学习与实践项目合集.zip”的资料包解压之后看到 Hadoop、Spark、Kafka、可视化大屏堆了一屏幕却不知道从哪个文件夹开始。我带教研组同学时也常遇到这种状况先讲 Python 基础再讲 Spark SQL等到真正跑离线日志才发现连 NameNode 和 ResourceManager 的区别都没搞清。其实大数据入门最难的地方不是代码量而是脑子里没有一条完整的数据链路。HDFS 负责把数据“放住”MapReduce 或 Spark 负责把数据算清Kafka 把实时流转起来ECharts 再让结果可见。下面就用 Hadoop 电商日志分析、Spark 实时流处理、集群搭建教程和数据可视化案例这条线把从离线到实时、从存储到展示的最小闭环走通。你照着跑一遍会比刷十遍入门视频有效得多。2. 集群搭建与 HDFS先把数据存储链路跑通2.1 单机、伪分布式、完全分布式怎么选在 Hadoop 安装与配置之前先确定部署形态。常见的大数据集群部署策略有三种单机模式不启动任何守护进程直接跑本地文件系统伪分布式在一台机器上同时启动 NameNode 和 DataNode适合入门完全分布式至少三台机器NameNode 和 DataNode 分离是生产环境的基础形态。新手前期不需要急着租三台服务器先用伪分布式把 HDFS 和 MapReduce 跑熟再扩展成 Spark 集群搭建复杂度会平滑很多。我一般建议的安装顺序是准备一台 CentOS 7 或 Ubuntu Server配好 JDK 8、SSH 免密登录如果服务器不能访问外网先配置本地 yum 源否则装 OpenJDK 时会被依赖包卡很久。随后设置HADOOP_HOME环境变量再修改core-site.xml和hdfs-site.xml。下面是伪分布式的核心配置!-- core-site.xml -- configuration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property /configuration!-- hdfs-site.xml -- configuration property namedfs.replication/name value1/value /property /configurationfs.defaultFS指定了 HDFS 的入口地址所有客户端都通过这个地址访问 Namenodedfs.replication是文件块副本数。伪分布式只有一台 DataNode副本数必须设为 1否则hdfs dfs -put上传后会出现 pending_replication 警告因为数据永远等不到第二个副本。配置完成后依次执行hdfs namenode -format start-dfs.sh这里有一个高频率踩坑点namenode -format只能执行一次。如果第二次格式化前没有删除/tmp/hadoop-*和dfs.name.dir指向的目录启动后 datanode 的 ClusterID 会和 namenode 不一致导致 DataNode 不断重试、不出现在 WebUI 上。遇到这种情况删掉数据目录重新格式化即可学习阶段没有存量数据不用心疼。部署模式进程是否分机副本数适用场景单机模式不启动 HDFS不涉及本地快速验证 MapReduce 逻辑伪分布式单台机器启动全进程1Hadoop入门、学习 HDFS 命令完全分布式至少 3 台机器3集群实战、Spark 集群搭建基础2.2 HDFS 写流程从客户端写入到副本确认hdfs 读写流程是大数据面试题里的高频题也是实际排错的基础。客户端要写一个 128MB 的文件时先向 NameNode 发起请求NameNode 检查权限和路径后返回一批可用的 DataNode。客户端把文件按 128MB 切块第一个块传给第一个 DataNode再由这个 DataNode 流水线复制给第二个、第三个。每写一个 chunk 都会计算校验和最后一个 DataNode 通过 RPC 返回 ack客户端收到确认后继续写下一块。这个流程里最容易出问题的是“租约lease”机制。如果某个客户端写了一半宕机HDFS 不会立刻把文件交给别人读而是等 lease 过期。你会在日志里看到Got exception while serving rpc或者previous writer likely failed to write。解决方法是hdfs debug recoverLease -path /user/edu/logs/access.log -retries 3这条命令会强制恢复文件租约但不保证数据完整。恢复后先hdfs fsck检查块健康度确认 0 个坏块再继续处理。2.3 HDFS 常用命令上传、查看与块分布做完集群搭建后日常打交道最多的就是 hdfs dfs 系列命令。下面这组命令覆盖了学习阶段最常用的场景hdfs dfs -mkdir -p /user/edu/logs hdfs dfs -put access.log /user/edu/logs/ hdfs dfs -ls -R /user/edu/logs hdfs fsck /user/edu/logs/access.log -files -blocks -locations-mkdir -p和 Linux 的mkdir -p一致递归创建目录。-put从本地上传文件到 HDFS生产环境也常用其等价命令-copyFromLocal。-ls -R递归查看目录和文件大小。最后一条fsck特别有用它不会读取数据内容只读 NameNode 元数据可以定位某个文件的每个块存在哪些 DataNode 上。如果发现某个块显示Missing replica说明这台 DataNode 的磁盘或者块复制出了问题。相比本地文件系统HDFS 的原子重命名和流式读取决定了它不适合存大量小文件。这也是后来很多人纠结 minio vs hdfs 的原因MinIO 这类对象存储在存图片、JSON 之类的非结构化数据上更灵活HDFS 则强在离线计算的本地性和生态兼容。2.4 集群搭建时的 DataNode 故障与 ZooKeeper 整合写代码报错不可怕集群进程起不来才让人头疼。常见的 DataNode 启动失败有两个原因磁盘权限不对或者 DataFrame 目录里的 VERSION 文件与 NameNode 不一致。第一类用ls -l看dfs.datanode.data.dir的属主第二类直接删除 data 目录后重新格式化最省事。生产环境里还要面对 NameNode 单点故障这就离不开 hadoop 和 zookeeper 整合实战。引入 ZooKeeper 后部署 JournalNode 和 ZKFC两个 NameNode 通过 ZK 竞争 active 状态active 节点挂了standby 在几十秒内接管服务。学习阶段不需要真的搭一套 HA但要知道这条命令的含义hdfs haadmin -transitionToActive nn1提示伪分布式阶段不要为了“像生产”就强行配 HA先搞懂单 NameNode 的日志、进程、端口关系。HA 只是解决高可用不解决业务逻辑。3. Hadoop 电商日志分析离线统计 PV/UV 的最小 MapReduce 任务3.1 电商日志字段与统计口径学习阶段用的 access.log 不需要太复杂但至少要有能够计算 PV/UV 的字段。我习惯用逗号分隔字段分别为 IP、访问时间、用户 ID、请求 URL、商品 ID、下单金额。下面是字段约定字段位置字段名示例用途1ip192.168.1.10按 IP 维度辅助取数2server_time2025-01-10 10:30:00事件时间、日期切片3user_idU12345UV 去重核心字段4request_url/item/1002PV 计数5product_idP1002商品维度统计6pay_amount199.00GMV 计算需求很简单每天有多少访问量PV多少独立用户UV以及每个商品的成交量。这些指标在离线阶段用一行 SQL 也能算但放在 MapReduce 里做是为了理解数据切分为 key-value 的过程。3.2 用 Hadoop Streaming Python 写第一个 MR 任务很多教材默认用 Java 写 Mapper 和 Reducer新手一上来就被 Maven 依赖绊住。我更推荐先用 Hadoop Streaming 跑 Python这样数据清洗逻辑可视、调试成本低后续接数据分析库也更顺。下面是统计 UV 的 mapper 和 reducer。mapper.py#!/usr/bin/env python3 import sys for line in sys.stdin: fields line.strip().split(,) if len(fields) 3: continue ip, server_time, user_id fields[:3] # 以日期作为 keyuser_id 作为 value交给 reducer 去重 date server_time[:10] print(f{date}\t{user_id})reducer.py#!/usr/bin/env python3 import sys cur_date None user_set set() for line in sys.stdin: date, user_id line.strip().split(\t) if date ! cur_date and cur_date is not None: print(f{cur_date}\t{len(user_set)}) user_set set() cur_date date elif cur_date is None: cur_date date user_set.add(user_id) if cur_date is not None: print(f{cur_date}\t{len(user_set)})运行命令hadoop jar $HADOOP_HOME/share/hadoop/tools/lib/hadoop-streaming-*.jar \ -D mapreduce.job.reduces2 \ -input /user/edu/logs/access.log \ -output /user/edu/output/pvuv \ -mapper python3 mapper.py \ -reducer python3 reducer.py \ -files mapper.py,reducer.py参数说明-D mapreduce.job.reduces2控制 Reducer 数量学习阶段设成 1 或 2 都行设太大会把小文件切得更碎。-files会把 Python 脚本分发到所有机器的工作目录不写这个参数会导致No such file or directory。-output一定不能是已存在的目录MapReduce 不会自动覆盖这一步报错最多。上述代码的 reducer 把同一个 day 下的 user_id 全部放入内存 set适合小数据量演示。真实生产里如果一天上亿用户内存会撑爆常见做法是用 HyperLogLog 或者布隆过滤器替换 set。3.3 数据清洗与分区日志里经常有爬虫、空值和异常请求。清洗逻辑一般放在 mapper 开头过滤条件根据业务定。比如只保留 URL 中包含/item/的请求同时丢弃非 2/3/4 开头的状态码valid fields[3] if len(fields) 3 else if valid.startswith(/item/): print(f{date}\t{user_id})清洗完的数据如果分区合理后面写 Spark 的时候会省很多事。按天分区是最基础的策略按小时分区适合活动时段。HDFS 目录结构可以设计成/user/edu/data/date20250110/配合 Hive 或 Spark 的分区裁剪扫描量能差一个数量级。3.4 数据倾斜与小文件治理hadoop 面试题里必考的倾斜问题在电商日志场景里非常容易出现某个头部商品访问量极高对应的 reducer 处理时间远长于其他节点。解决思路是先加随机前缀打散再做二次聚合。另一个问题是小文件电商日志按小时落盘会产生大量 KB 级文件NameNode 内存全被元数据吃掉。处理小文件优先从输入端下手比如用 CombineTextInputFormat 把小文件合并成一个逻辑分片-D mapreduce.input.fileinputformat.input.dir.recursivetrue也可以在输出端提前做一次汇总把按小时日志合并成按天日志。记住一个原则能不改代码就不改代码先通过输入分片和输出合并解决。现象可能原因解决方式个别 reducer 卡到超时key 分布不均随机前缀 两阶段聚合文件数上万、NameNode 内存上升按小时/分钟落盘CombineTextInputFormatreducer 输出文件数量过多并行度设置不合理调低mapreduce.job.reduces磁盘写满中间结果未清理及时清理 /tmp 下中间目录4. Spark 实时流处理Kafka 进门Structured Streaming 接单4.1 为什么选 Spark 做实时流处理离线日志算完之后业务又提出要看“今天截至现在的 GMV”这时候就需要 Spark 实时流处理。Spark 生态里有两套流方案Spark Streaming 基于 RDD 微批Structured Streaming 基于 DataFrame API从 Spark 2.0 开始成为主流。为什么推荐后者因为它的流计算结果和离线 DataFrame 共用同一套 API对流式聚合、事件时间窗口的支持更完整逻辑可以直接从离线批处理搬过来。和 Flink 相比Spark Streaming 的吞吐量在同规模集群下表现不错但延迟是秒级。对于数据大屏这类对延迟不敏感、但需要高吞吐的场景Structured Streaming 很合适。如果未来要并行处理复杂事件规则再考虑 Flink。4.2 最小可运行的 Structured Streaming 代码假设 Kafka 里已经有一个order_topic生产环境里通常配三台 Kafka brokertopic 分 6 个分区。下面这段 PySpark 代码完成“5 秒窗口销售额统计”from pyspark.sql import SparkSession from pyspark.sql.functions import from_json, col, sum, window spark SparkSession.builder \ .appName(ECommerceRealtime) \ .getOrCreate() df spark.readStream \ .format(kafka) \ .option(kafka.bootstrap.servers, node01:9092,node02:9092,node03:9092) \ .option(subscribe, order_topic) \ .option(startingOffsets, earliest) \ .load() # Kafka 消息的 value 是二进制先转字符串再解析 JSON schema user_id STRING, product_id STRING, amount DOUBLE, ts TIMESTAMP orders df.selectExpr(CAST(value AS STRING) AS json) \ .select(from_json(json, schema).alias(data)) \ .select(data.user_id, data.amount, data.ts) # 基于事件时间 ts 开窗watermark 允许迟到 10 秒 result orders \ .withWatermark(ts, 10 seconds) \ .groupBy(window(ts, 5 seconds)) \ .agg(sum(amount).alias(sales)) query result.writeStream \ .outputMode(update) \ .format(console) \ .trigger(processingTime2 seconds) \ .start() query.awaitTermination()几个关键参数startingOffsetsearliest表示从头消费生产环境要设为latest否则每次重启都会重放全部数据withWatermark(ts, 10 seconds)表示允许事件时间比数据到达时间晚 10 秒超过这个阈值就只能等刷新增量trigger(processingTime2 seconds)控制微批触发周期对延迟越敏感设得越短但太短会增加调度开销。代码写完后用 spark-submit 提交spark-submit --master yarn \ --deploy-mode cluster \ --driver-memory 1g \ --executor-memory 2g \ --num-executors 3 \ --executor-cores 2 \ streaming_job.py4.3 Spark 集群搭建与资源参数怎么调如果已经搭建过 Hadoop 完全分布式集群Spark 集群搭建会很简单在所有节点解压 Spark 安装包配置spark-env.sh里的SPARK_MASTER_HOST和SPARK_WORKER_CORES在slaves文件里写上 Worker 节点主机名启动后访问 8080 端口看到 Master 页面即为成功。资源参数里最容易出错的是执行器内存。--executor-memory 2g如果设得太高比如 32gYARN 把一个节点的大部分内存都分给一个执行器其他执行器启动不了任务表现为频繁 OOM。常见经验值是单 executor 1~4GB每个节点挂 2~4 个 executor。--num-executors也不是越多越好要结合总 vCore 数和队列配额算。4.4 实时去重、状态清理与落库流处理里一个让人头疼的问题是去重。比如用户下单事件可能被重发或者 Kafka 消费者重平衡导致重复消费。简单场景下可以给消息加唯一订单号写入外部存储时使用 HBase 的put if absent语义如果想要 Spark 内存内去重可以配合mapGroupsWithState维护一个 state但要注意状态会无限增长必须配合timeoutDuration清理。还有一个容易忽略的坑是水位线。如果上游 Kafka 的 topic 只保留了最近 7 天数据而你设置的水位线超过 7 天Spark 会把大量早期状态堆在内存里。按业务实际设定水位线比如订单场景 10~30 分钟足够。5. 数据可视化案例与上手指南用大屏验证整个链路5.1 把 HDFS 分析结果渲染成可用的大屏前四章跑出来的pvuv结果还是文本接入数据可视化时我会直接用 pyecharts 生成独立 HTML不依赖后端框架。先把 HDFS 上结果拉到本地然后交给 Python 画图。from pyecharts.charts import Bar, Line from pyecharts import options as opts import pandas as pd df pd.read_csv(pvuv_result.txt, sep\t, headerNone, names[date, pv, uv]) bar Bar() bar.add_xaxis(df[date].astype(str).tolist()) bar.add_yaxis(PV, df[pv].tolist()) bar.add_yaxis(UV, df[uv].tolist()) bar.set_global_opts( title_optsopts.TitleOpts(title电商 PV/UV 日报大屏), datazoom_opts[opts.DataZoomOpts()], yaxis_optsopts.AxisOpts(name访问量), ) bar.render(dashboard.html)datazoom_opts会在图表下方生成一个缩放条数据多时不用改代码就能看单日细节。如果你要拼接多个图表直接用Grid或Tab把它们组织到同一张 HTML 里比照着一个大屏模板改 CSS 方便得多。5.2 离线和实时结果做交叉验证大屏好看不代表算得对。我的验证技巧是找一台测试机同时跑离线 MapReduce 和 Spark Streaming 任务用同一份 Kafka topic 从头开始消费比对每天的离线 PV/UV 与实时窗口累计值之间的差值。正常情况下因为没有故障重放差值应当始终在 1% 以内。如果差值较大优先检查 Kafka 分区消费者是否发生了 rebalance再检查 Spark 任务是否在 batch 失败时自动重试导致同一批数据被计算两次。写大屏时不要在页面端做太多聚合计算把聚合下推到 Spark 或 Hive前端只负责渲染。这样即使同一时间段有三个人同时打开页面后端也只计算一次。5.3 按一周时间推进的学习路径与最后一步如果你完全从零开始一周时间可以用这套节奏前三天搭好 Hadoop 伪分布式学会 hdfs 常用命令第四天跑通一个 MapReduce Python 任务第五天引入 Kafka 和 Spark Streaming第六天做可视化第七天把离线和实时结果对齐。做完这条路径后你去刷 hadoop 面试题时会发现很多概念已经落到具体命令上复习效率会明显不同。最后送你一个能立刻用上的小技巧把那张dashboard.html生成命令写进 crontab每天 9 点自动从 HDFS 拉取前一日结果并重新渲染这样每天早上打开浏览器就能看到前一天的数据不用手动开终端。这也是“大数据入门到实战完整学习路径”里最能见到回报的一步。本文还有配套的精品资源点击获取
返回列表