ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

从零搭建离线数仓:Hadoop+Spark+Flume完整实践笔记

从零搭建离线数仓:Hadoop+Spark+Flume完整实践笔记 这篇《大数据实践笔记2》拖了挺久总算整理出来了。我平时在一家互联网公司做数据开发业余会接一些大数据相关的毕设辅导和面试指导这篇笔记其实是我自己从零搭一套“可复现的离线数仓小项目”的完整记录。核心链路不搞虚的用3台Linux服务器搭一个真实的多节点Hadoop集群用Python生成模拟电商用户行为日志再用Flume采集到HDFS用Spark SQL做清洗和指标计算最后通过Spring Boot接口把结果喂给一个ReactTypeScriptECharts的可视化大屏。这东西适合谁一类是准备大数据毕业设计的同学这套链路正好覆盖“数据采集-存储-计算-展示”的完整闭环另一类是刚转行大数据、想亲手练集群部署的工程师还有一类是准备大数据面试的人因为项目里遇到的很多坑面试官真的会问到。这篇不是理论教程我尽量把每一步的“为什么这么做”和“我踩过的坑”都写清楚照着做能跑通少走弯路。1. 动手之前先把整套链路想明白1.1 从伪分布式到真集群差的不只是机器数量我第一篇实践笔记写的是在单台上搭伪分布式Hadoop当时很多功能都能跑但说实话那种方式会掩盖掉分布式系统真正的难点。单机上跑MapReduce不用考虑数据本地性不用管网络传输连节点宕机都不用操心。可一旦换成真集群任务调度、资源隔离、节点间通信、数据倾斜这些问题全都会冒出来这才是大数据项目里最有价值的部分。所以我决定在笔记2里直接上“一主三从”的迷你集群。三台Worker节点不算多但对学习和毕设来说完全够用而且能把分布式该有的问题都暴露出来。这个选择背后有个很现实的原因如果你只是跑单机伪分布式简历上写“熟悉Hadoop”是没有说服力的但如果你能在面试里讲清楚“我在3台机器上搭了集群遇到过DataNode起不来、Spark任务OOM、数据倾斜”这些真实问题效果完全不一样。1.2 整体流程离线数仓的简化版这套项目的处理流程其实就是很多公司离线数仓的缩小版环境准备3台Linux服务器划分角色配置网络、JDK、SSH免密、时钟同步。部署Hadoop和Spark提交一个测试任务确认分布式环境可用。用Python脚本生成模拟电商行为日志通过Flume实时监控日志目录并写入HDFS。用Spark SQL读取HDFS上的原始日志做数据清洗、去重、维度补充再按不同业务口径计算指标。把计算结果写入MySQL前端通过Spring Boot接口读取最后用ECharts做成可视化大屏。最后统一整理踩坑记录和调优参数。这套流程选型上故意避开了太重的组件比如没有上Hive、没上Kafka。原因很简单在3台低配机器上硬塞太多组件光是维护就能把人劝退。学习阶段先保证链路完整跑通把HDFS、Yarn、Spark这些核心组件吃透后面再加组件是水到渠成的事。1.3 集群规模怎么定先算一笔资源账好多同学上来就问“我需要几台机器”其实这个答案取决于你的数据量。我先做了个大概估算单条电商行为日志按250字节算一天模拟500万条大约1.25GB。HDFS默认3副本所以日增原始数据约占3.75GB存储。还要留出中间结果、临时文件、Spark Shuffle临时目录的空间一般按原始数据的2到3倍预留。一个月的数据量预留200GB左右比较稳妥。所以我的节点规划如下节点角色配置服务角色磁盘master8核16GNameNode、ResourceManager、SecondaryNameNode200Gworker014核8GDataNode、NodeManager200Gworker024核8GDataNode、NodeManager200Gworker034核8GDataNode、NodeManager200G这不是拍脑袋定的。NameNode和ResourceManager都是内存大户尤其NameNode的元数据会随文件数量增长给16G比较稳。而DataNode主要吃磁盘IO和带宽4核8G足够。如果条件有限可以用虚拟机替代实体机但每台虚拟机的内存最好不低于4G否则跑Spark任务时很容易OOM。另外强烈建议用云服务器因为可以按小时计费开三台临时机器练习练完直接释放成本很低。2. 集群部署实操三台机器从裸机到跑通2.1 基础环境初始化不做这步后面全是坑先把系统的坑填掉不然装到一半再回头排错很窝火。以下操作在所有节点都要执行。第一步配置hosts。每台机器都要配确保能通过主机名互相访问因为Hadoop内部通信强依赖主机名要是靠IP解析节点一换IP就全线崩盘。# /etc/hosts 192.168.1.10 master 192.168.1.11 worker01 192.168.1.12 worker02 192.168.1.13 worker03第二步安装JDK。Hadoop 3.3.x要求JDK8起步我用的是JDK8。安装之后记得配置JAVA_HOME并把它写进/etc/profile。这一步容易被忽略但很多启动失败都跟JAVA_HOME没配好有关。第三步配置SSH免密登录。主节点需要免密登录到所有节点包括自己否则start-dfs.sh启动过程中会反复要求输密码。我在master上执行ssh-keygen -t rsa -P -f ~/.ssh/id_rsa ssh-copy-id master ssh-copy-id worker01 ssh-copy-id worker02 ssh-copy-id worker03第四步同步时间。这一点比很多人想的更重要。Hadoop集群和HDFS通信有超时机制如果节点间时间偏差太大会出现各种让人摸不着头脑的RPC异常。我在这三台机器上都配置了NTP同步或者至少手动校准一次# 手动校准示例 date -s $(curl -s --head http://www.baidu.com | grep -i ^date: | sed s/^[Dd]ate: //g)2.2 Hadoop集群安装与启动下载Hadoop 3.3.4版本比较稳妥太新的版本可能和Spark有兼容性问题。下载解压到/opt目录后重点修改这几个配置文件。core-site.xml配置默认文件系统和临时目录configuration property namefs.defaultFS/name valuehdfs://master:9000/value /property property namehadoop.tmp.dir/name value/data/hadoop/tmp/value /property /configurationhdfs-site.xml配置副本数和NameNode/DataNode数据目录。这里有个关键点生产环境会把NameNode元数据目录单独放一块盘但我们这个项目直接在/data下建目录就够了。property namedfs.replication/name value3/value /property property namedfs.namenode.name.dir/name value/data/hadoop/namenode/value /property property namedfs.datanode.data.dir/name value/data/hadoop/datanode/value /propertyyarn-site.xml要指定ResourceManager跑在master否则默认起在提交任务的节点property nameyarn.resourcemanager.hostname/name valuemaster/value /property property nameyarn.nodemanager.resource.memory-mb/name value6144/value /property property nameyarn.nodemanager.resource.cpu-vcores/name value3/value /property上面这两个资源参数很关键它们告诉Yarn每个NodeManager最多能用多少内存和CPU。请注意worker是8G内存系统本身要占到1.5G左右再扣掉DataNode等服务NodeManager给6144MB是比较稳的给太多会导致机器卡死。workers文件写上所有DataNode主机名每行一个。然后先启动HDFS再启动Yarnhdfs namenode -format start-dfs.sh start-yarn.sh格式化NameNode只有第一次需要以后千万别随便执行否则集群ID变了DataNode全部启动失败。启动后分别在master和worker上执行jps检查进程master上应该有NameNode、SecondaryNameNode、ResourceManager每个worker上应该有DataNode、NodeManager再用hdfs dfsadmin -report确认三台DataNode都在线。只有这一步通过HDFS才真正可用。2.3 Spark部署跑在Yarn上的资源调度Spark我用的是3.3.2这个版本对Hadoop 3.x支持很好。配置方面除了SPARK_HOME环境变量关键改两个文件。spark-env.sh里显式指定Java目录避免Spark找不到JDK同时设置Driver内存export JAVA_HOME/usr/local/java export SPARK_DRIVER_MEMORY2gspark-defaults.conf里把提交模式固定为Yarn并设置Executor资源spark.masteryarn spark.driver.memory2g spark.executor.memory4g spark.executor.cores2 spark.sql.shuffle.partitions100注意这个shuffle.partitions参数很多Spark性能问题都是因为它引发的。默认值是200在小集群上200个分区会导致大量小任务在节点间频繁传输反而更慢我调成了100。如果发现某个Stage有大量数据倾斜再结合具体数据量做调整这个后面细说。部署完成后提交一个简单的spark自带示例验证跑通spark-submit --class org.apache.spark.examples.SparkPi \ --master yarn \ --deploy-mode cluster \ /opt/spark/examples/jars/spark-examples_2.12-3.3.2.jar 10能正常输出圆周率结果说明Spark和Yarn的配合没问题。3. 数据从哪来造日志、采日志、查日志3.1 用Python脚本生成模拟电商行为日志没有真实业务数据的时候我们完全可以自己造一份看起来像样的数据。我用Python写了个小脚本按电商场景生成用户行为日志每条日志包含用户ID、商品ID、类目、行为类型浏览/加购/下单/支付、渠道、省份、时间戳等字段。行为类型按真实业务比例设置权重浏览70%、加购15%、下单10%、支付5%。生成脚本核心逻辑如下import random import time users [fu{i:06d} for i in range(100000)] items [fi{i:06d} for i in range(50000)] categories [fc{i:02d} for i in range(50)] actions [view, cart, order, pay] weights [70, 15, 10, 5] channels [app, h5, pc] provinces [广东, 江苏, 浙江, 山东, 四川, 北京, 上海, 湖北] def generate_log(timestamp): return { user_id: random.choice(users), item_id: random.choice(items), category_id: random.choice(categories), action: random.choices(actions, weightsweights)[0], channel: random.choice(channels), province: random.choice(provinces), timestamp: timestamp } # 每100万条写一个文件文件按小时切割更贴近真实情况 for batch in range(5): with open(f/data/flume/logs/access_{batch}.log, w) as f: for _ in range(1000000): ts int(time.time()) log_json json.dumps(generate_log(ts)) f.write(log_json \\n) time.sleep(1)写文件时我故意按批次写入这样Flume的spooldir源监控目录时有持续的数据流能看到采集过程实时进行。如果一次性生成几十个G的文件Flume跑起来也容易因为处理不及时产生积压。3.2 Flume采集让日志自动进入HDFSFlume是个日志采集工具它的核心就是三部分Source读取数据源、Channel临时缓冲、Sink输出到目的地。这里我选spooldir作为Source监控日志目录HDFS作为SinkChannel用内存通道。配置如下a1.sources r1 a1.sinks k1 a1.channels c1 a1.sources.r1.type spooldir a1.sources.r1.spoolDir /data/flume/logs a1.sources.r1.fileHeader false a1.sinks.k1.type hdfs a1.sinks.k1.hdfs.path hdfs://master:9000/data/ods/access/%Y%m%d a1.sinks.k1.hdfs.fileType DataStream a1.sinks.k1.hdfs.rollInterval 3600 a1.sinks.k1.hdfs.rollSize 134217728 a1.sinks.k1.hdfs.rollCount 0 a1.channels.c1.type memory a1.channels.c1.capacity 10000 a1.channels.c1.transactionCapacity 1000 a1.sources.r1.channels c1 a1.sinks.k1.channel c1这几个参数是调出来的。hdfs.rollInterval控制按时间滚动文件3600秒切一个文件避免单个文件过大hdfs.rollSize设置128MB就滚动一次因为后续Spark按文件块读取时128MB是最舒服的块大小不会因为大量小文件导致任务数爆炸。还有一点HDFS Sink的path里用了%Y%m%d日期变量配合timestamp拦截器可以把数据自动按天分区存储这是很重要的数仓习惯。启动命令flume-ng agent --conf conf --conf-file /opt/flume/conf/access-agent.conf --name a1 -Dflume.root.loggerINFO,console看到控制台不断输出Event数量说明采集进HDFS了整个链路就通了。3.3 采集后的数据质量检查数据进HDFS之后别急着写计算逻辑先花两分钟验证一下数据质量。我在HDFS上先看一眼目录结构和文件大小hdfs dfs -ls hdfs://master:9000/data/ods/access/20250601 hdfs dfs -du -s hdfs://master:9000/data/ods/access/20250601然后抽样看看前几行内容确认格式没问题hdfs dfs -cat hdfs://master:9000/data/ods/access/20250601/*.log | head -5建议再做一个简单的记录条数统计。因为Flume偶尔会丢数据尤其是Source读文件过快、Channel积压时丢一行都说明配置有瓶颈。我习惯用wc -l对比源文件行数和HDFS上的行数能快速判断采集是否完整。数据质量这步可能看起来不刺激但真正跑起分析时你会感谢自己多看了这一步。4. 离线分析用Spark SQL把日志变成指标4.1 指标定义先想清楚要算什么数据有了以后不要急着写代码。我习惯先在纸上把指标列清楚否则很容易在Spark SQL里来回改。这次项目我选了4个比较典型的大数据指标分时PV和UV、商品点击Top10、行为转化漏斗、分省支付金额分布。为什么选这些因为它们分别对应了去重计数、排序统计、多条件聚合、多维分组基本覆盖了面试里最常见的SQL场景。定义好指标后我在Spark里建了一张临时表把HDFS上的原始日志读进来DatasetRow raw spark.read() .json(hdfs://master:9000/data/ods/access/20250601); raw.createOrReplaceTempView(ods_access);如果文件里既有JSON又有脏数据建议还是先转成DataFrame统一schema后面SQL就好写了。4.2 清洗过滤脏数据和异常值我在生成数据的时候还算规整但因为Flume采集或者文件写入时偶尔会出现空行所以清洗这步不能省。清洗要处理的典型问题有关键字段为空、行为类型不在合法枚举内、时间戳异常、重复记录。这步我用一条SQL完成CREATE OR REPLACE TEMP VIEW dwd_access AS SELECT user_id, item_id, category_id, action, channel, province, from_unixtime(ts, yyyy-MM-dd HH:mm:ss) AS event_time, date_format(from_unixtime(ts, yyyy-MM-dd HH:mm:ss), HH) AS hour FROM ods_access WHERE user_id IS NOT NULL AND item_id IS NOT NULL AND action IN (view, cart, order, pay) AND ts 0这里我新增了hour字段因为后面按小时聚合分时PV/UV时不需要每次都用date_format再算一遍。这种“在清洗层尽量把能算的维度先算好”的思路在真实数仓里叫维度退化可以大大简化后续计算。4.3 指标计算Spark SQL的几种典型场景分时PV和UV是我最看重的指标因为它同时考察count和count distinctSELECT hour, count(*) AS pv, count(DISTINCT user_id) AS uv FROM dwd_access GROUP BY hour ORDER BY hour这里有个值得注意的性能知识点count(DISTINCT)在实际大数据量下非常贵因为它需要把所有user_id分布到各个节点去重后再汇总。如果数据量到了亿级建议改用approx_count_distinct做近似去重误差率可以在1%以内速度却能快好几倍。商品点击Top10直接按item_id分组排序SELECT item_id, count(*) AS cnt FROM dwd_access WHERE action view GROUP BY item_id ORDER BY cnt DESC LIMIT 10;转换漏斗则是典型的条件聚合一行SQL算出浏览、加购、下单、支付各自的总次数SELECT sum(CASE WHEN action view THEN 1 ELSE 0 END) AS view_cnt, sum(CASE WHEN action cart THEN 1 ELSE 0 END) AS cart_cnt, sum(CASE WHEN action order THEN 1 ELSE 0 END) AS order_cnt, sum(CASE WHEN action pay THEN 1 ELSE 0 END) AS pay_cnt FROM dwd_access这4个结果我都写回HDFS的hive分区表或者直接写MySQL。生产上一般会先落到Hive分区表做沉淀再同步到MySQL给前端查。这里我选择直接把这个结果用JDBC写入MySQLresult.write() .mode(SaveMode.Overwrite) .jdbc(jdbc:mysql://master:3306/dw, t_hour_pv_uv, props);4.4 任务调度让分析每天自动跑离线分析最怕的事就是“人肉跑数”今天记得跑明天忘了就废了。我在这套项目里用crontab做最简单的调度每天凌晨2点执行一次spark-submit把前一天的数据算一遍0 2 * * * /opt/spark/bin/spark-submit --class com.demo.EtlRunner \ --master yarn \ --deploy-mode cluster \ /opt/app/ecommerce-analysis.jar为什么选凌晨2点因为凌晨业务低峰期集群资源空闲而且上游Flume的数据在凌晨1点左右基本都同步完毕数据完整性有保障。如果以后数据源变多可以考虑上Azkaban或DolphinScheduler这类调度平台但项目初期crontab足够。5. 可视化大屏ECharts把结果端到前端5.1 技术选型ReactTSECharts的组合可视化大屏前端我选了ReactTypeScriptECharts这套组合现在是做数据展示类项目的主流选择。React负责组件化布局TypeScript保证接口数据类型不出错ECharts负责绘图。如果你只是纯粹想快速看到效果用一个HTML文件引入ECharts也完全可行但考虑到毕设或者项目展示需要ReactTS更像一个完整的前端正规军也更好扩展。后端我写了一个非常简单的Spring Boot服务只提供两个接口一个查分时PV/UV一个查商品Top10。Spring Boot直连MySQL返回JSON给前端。为什么要单独写个后端而不是前端直接读MySQL一是安全考虑不能把数据库账号密码暴露在浏览器端二是为了演示完整的前后端分离架构这在简历上是个加分项。5.2 大屏布局与ECharts图表配置大屏的视觉比例一般按1920x1080设计布局思路是中间大、两侧小。我用了CSS Grid布局分成三个区域左侧放商品Top10柱状图中间放核心指标数字和分省地图右侧放分时PV/UV折线图和转化漏斗图。顶部是项目标题和日期。核心的折线图配置大概是这样的const option { tooltip: { trigger: axis }, legend: { data: [PV, UV] }, grid: { left: 60, right: 20, top: 40, bottom: 30 }, xAxis: { type: category, data: hours }, yAxis: { type: value }, series: [ { name: PV, type: line, smooth: true, data: pvData }, { name: UV, type: line, smooth: true, data: uvData } ] };在React里我建议用useRef保存ECharts实例用useEffect请求数据和初始化图表。不要每次请求都重新创建实例那样会有明显卡顿。5.3 大屏适配与性能优化大屏最常见的坑就是浏览器缩放后图表错位。我的做法是基于1920x1080设计稿然后用transform: scale做整体缩放这样在16:9的屏幕上显示效果最理想。这个方案是兼容性成本最低的。还有几个性能细节页面里图表超过5个时要确保图表在组件卸载时执行dispose否则会有内存泄漏ECharts数据更新时用setOption而不是再次init数据请求加loading状态避免白屏。这些看着琐碎但都是实际开发一定会遇到的。6. 我这一路踩过的坑和排查思路6.1 部署和运行阶段的典型故障这套项目我在不同机器上复现过很多次整理了一份常见问题速查表现象可能原因排查和解决办法DataNode进程起不来NameNode和DataNode的clusterID不一致检查/data/hadoop/namenode/current/VERSION和datanode的VERSION统一clusterID任务在Yarn上一直ACCEPTED不执行NodeManager可用内存不足查看yarn-site.xml的nodemanager内存配置预留足够系统内存Spark执行时频繁GC甚至OOMexecutor内存太小或者数据倾斜调大executor-memory倾斜键加随机前缀做两阶段聚合集群数据节点显示宕机节点时间不同步引起RPC超时配置NTP同步保证每台机器时间一致Flume采集日志丢失Channel容量太小Source写入过快调大memory channel的capacity或改用file channel前端大屏加载缓慢一次性查询太多数据MySQL里预聚合统计接口只查当日汇总数据6.2 关于Spark调优的几个实际心得很多人在小集群上喜欢照搬网上的调优参数动辄给executor分配16G内存结果节点直接卡死。我自己调优的原则是让每个Executor占用的总内存不超过节点内存的70%。比如worker节点8GNodeManager分配6GSpark Executor最多给4G剩下留给系统和其他进程。shuffle分区数也是个大坑。有一次我跑分时PV/UV数据只有几百万条但shuffle.partitions还是默认的200导致每个任务只处理很少的数据整个作业卡了十几分钟。后来改成100再配合coalesce把结果文件合并速度快了3倍。小集群上宁可分区少一点让每个任务多吃点数据也别让任务碎片化。还有一个小技巧写结果到MySQL时先repartition(1)把结果文件合并成一个再写入否则Spark会产生几十个小文件写数据库时频繁建连接慢得离谱。6.3 一定不要忽略数据倾斜数据倾斜是分布式计算里最容易遇到又最让人头疼的问题。我这次造的数据故意埋了个坑某个类目的商品被大量用户浏览导致group by item_id时这个类目所在的Task处理的数据量是其他Task的几十倍跑起来像有一个“孤岛”任务迟迟结束不了。解决办法我用了两种第一种是做两阶段聚合先给每个item_id加一个随机前缀分散到多个Task完成第一轮聚合后再去掉前缀做第二轮聚合。第二种是过滤掉高频噪点数据比如只统计前1000个热门商品把极热数据单独处理。面试时如果你能讲清楚这两种思路比背一百个八股文都管用。大数据实践笔记2到这里整个项目就完整跑通了。说实话我在不同机器上搭了三次才敢把配置和排错思路写出来前两次都倒在各种“小问题”上比如hosts没配全、clusterID没统一、Flume的Chanel容量太小。但恰恰是这些坑让我对集群的理解比看十篇文档都深。如果你也在做大数据相关的事别怕慢强烈建议亲手把集群搭一遍。等你把日志采进去、数据算出来、图表亮起来的那一刻那种踏实感绝对不是只看教程能给的。
返回列表