
简介本资源是一个基于Hadoop生态的实战型大数据分析项目面向大数据初学者与Java开发者聚焦海量酒店数据的分布式处理与统计分析。项目完整实现MapReduce编程模型涵盖数据清洗、HDFS存储、Java编写的Mapper/Reducer逻辑及结果汇总可支撑省市酒店数量分布、平均房价等典型业务指标计算适用于高校课程设计、岗位技能实训与Hadoop集群实操演练。压缩包共79个文件含21个Java源码与21个编译后class文件构成核心MapReduce程序、12个CRC校验文件、7个XML配置含Hadoop集群参数与Maven构建配置、6个part-r-00000输出文件及2个CSV原始数据含hotel.csv整体体积758KB结构清晰便于本地调试与集群部署。已有2096人学习下载提供从数据导入、作业提交到结果解读的全流程配套包含说明.txt操作指南、hadoop-hotel工程目录及可直接运行的脚本与配置显著降低Hadoop环境搭建与调试门槛。1. 为什么用 Hadoop 处理全国酒店数据不是“大炮打蚊子”而是必须的硬需求你手头有一份来自文旅局、OTA 平台或爬虫采集的全国酒店数据——34 个省级行政区、333 个地级市、2843 个县级单位每家酒店含名称、地址、经纬度、星级、价格区间、开业年份、房型数量、实时库存、近30天订单量、用户评分、标签亲子/商务/民宿/电竞等……原始 CSV 总量超 12GB单文件最大 800MB每天新增 50 万条记录要求 T1 完成清洗、去重、地理编码、多维聚合按省/市/星级/价格带/标签组合、热力图生成、异常价格预警。这时候用 pandas 读取卡死、MySQL 导入报错、本地 Python 脚本跑 6 小时还没出结果——这不是你代码写得差是数据规模已经越过了单机处理的物理天花板。Hadoop 在这里不是“学着玩”的分布式玩具而是唯一能扛住「高通量数据处理」压力的工业级底座它把 12GB 数据自动切片InputSplit、分发到多个节点并行解析、容错重试失败任务、统一管理海量中间结果。尤其当你要做「省市维度交叉对比」比如对比长三角 vs 成渝地区四星酒店均价波动、「标签共现分析」“温泉亲子”组合在北方冬季是否显著增长、「时空联合聚合」某省会城市周末夜间 22:00–24:00 的低价房剩余率这些操作天然需要 MapReduce 或 Spark on YARN 的 shuffle 能力——而伪分布式搭建、YARN 资源调度、HDFS 文件块策略正是让这个项目从“跑不起来”到“稳态产出”的关键链路。适合正在做课程设计、企业数仓接入、或真实业务中接手酒店类数据治理的工程师不是纯理论派是明天就要改配置、跑 job、查日志的人。2. 从零构建 Hadoop 伪分布式环境避开 JDK 版本陷阱与 core-site.xml 的隐形坑Hadoop 伪分布式Pseudo-Distributed Mode不是“单机模拟集群”而是真启动 NameNode/DataNode/ResourceManager/NodeManager 四个守护进程共享同一台机器的 CPU 和内存但完全复用生产环境的配置逻辑和通信协议。这对调试数据处理流程、验证 MR/Spark 作业行为、压测 HDFS 写入吞吐量比完全分布式更高效、比本地模式更真实。我们以 Hadoop 3.3.6当前 LTS 最稳定版本为例全程基于 Ubuntu 22.04 OpenJDK 11注意Hadoop 3.x 不支持 JDK 17这是血泪经验。2.1 JDK 与 Hadoop 环境变量两个必须死守的硬约束提示hadoop_home配置错误是 73% 的伪分布式启动失败根源。不要用export HADOOP_HOME/opt/hadoop这种模糊路径必须指向解压后的完整绝对路径且该路径下必须存在sbin/start-dfs.sh和etc/hadoop/core-site.xml。# 下载并解压务必用官方二进制包非源码编译版 wget https://downloads.apache.org/hadoop/common/hadoop-3.3.6/hadoop-3.3.6.tar.gz tar -xzf hadoop-3.3.6.tar.gz -C /opt/ # 设置 JAVA_HOMEOpenJDK 11 是 Hadoop 3.3.6 唯一官方认证版本 export JAVA_HOME/usr/lib/jvm/java-11-openjdk-amd64 export PATH$JAVA_HOME/bin:$PATH # 关键HADOOP_HOME 必须精确到解压目录且不能有软链接 export HADOOP_HOME/opt/hadoop-3.3.6 export PATH$HADOOP_HOME/bin:$HADOOP_HOME/sbin:$PATH # 验证 hadoop version # 应输出 3.3.6且无 Could not find or load main class 错误参数说明JAVA_HOME必须指向 JDK 11 的jre上一级目录即含bin/java的路径不是 JRE 路径HADOOP_HOME若指向/opt/hadoop软链接start-dfs.sh会因pwd解析失败导致hdfs namenode -format报NoClassDefFoundError所有环境变量需写入~/.bashrc并source ~/.bashrc否则sbin/下脚本无法继承。2.2 四大核心配置文件每个property都决定数据能否落地伪分布式本质是“单机上的最小集群”因此core-site.xml、hdfs-site.xml、yarn-site.xml、mapred-site.xml必须全部修改缺一不可。重点不是抄模板而是理解每个参数的物理意义core-site.xml定义 HDFS 访问入口与本地缓存策略configuration property namefs.defaultFS/name valuehdfs://localhost:9000/value !-- 注意不是 file:///这是 HDFS URI -- /property property namehadoop.tmp.dir/name value/opt/hadoop-3.3.6/data/value !-- 必须手动创建且 chmod 755 -- /property /configuration逻辑说明fs.defaultFS是所有 Hadoop 客户端包括hadoop fs -ls的默认文件系统地址。设为hdfs://localhost:9000后hadoop fs -put local.csv /input/实际走的是 HDFS 协议而非本地文件系统。hadoop.tmp.dir是 NameNode 元数据、DataNode 块存储、YARN 临时文件的根目录必须提前mkdir -p /opt/hadoop-3.3.6/data chmod 755 /opt/hadoop-3.3.6/data否则格式化失败。hdfs-site.xml控制数据块副本与存储位置configuration property namedfs.replication/name value1/value !-- 伪分布式只需 1 副本设为 3 会导致 DataNode 启动失败 -- /property property namedfs.namenode.name.dir/name valuefile:/opt/hadoop-3.3.6/data/namenode/value /property property namedfs.datanode.data.dir/name valuefile:/opt/hadoop-3.3.6/data/datanode/value /property /configuration参数说明dfs.replication1是伪分布式铁律。Hadoop 默认要求至少 3 个 DataNode 才允许replication3单节点强行设 3 会导致start-dfs.sh后jps看不到 DataNode 进程。namenode.name.dir和datanode.data.dir必须是绝对路径且目录需手动创建mkdir -p /opt/hadoop-3.3.6/data/{namenode,datanode}。yarn-site.xml激活资源调度器YARN 是 Spark/Hive 的运行底座configuration property nameyarn.nodemanager.aux-services/name valuemapreduce_shuffle/value /property property nameyarn.resourcemanager.hostname/name valuelocalhost/value /property property nameyarn.nodemanager.resource.memory-mb/name value4096/value !-- 根据你的机器内存设建议 ≥3GB -- /property /configuration逻辑说明yarn.nodemanager.aux-servicesmapreduce_shuffle是 MapReduce 作业能运行的关键它启用了 ShuffleHandler 服务负责 Reduce 阶段拉取 Map 输出。yarn.nodemanager.resource.memory-mb设太小如默认 1024MB会导致 Spark 作业因内存不足被 YARN Kill设太大超过物理内存 80%则系统 OOM。mapred-site.xml绑定 MapReduce 框架到 YARNconfiguration property namemapreduce.framework.name/name valueyarn/value !-- 强制使用 YARN不是 local 模式 -- /property /configuration注意此文件默认是mapred-site.xml.template需先cp mapred-site.xml.template mapred-site.xml再修改。mapreduce.framework.nameyarn是伪分布式区别于本地模式的核心开关——没有它hadoop jar ...会降级为单线程执行彻底失去分布式意义。3. 酒店数据 ETL 流程从原始 CSV 到可分析 Hive 表的全链路实操全国酒店数据不是结构规整的数据库导出表而是混合来源的“脏数据沼泽”地址字段含括号、换行符、emoji经纬度有空值、字符串“NULL”、坐标系混用WGS84 vs GCJ02价格区间写成“¥200-¥500”或“200元起”星级标注为“★★★☆”或“准五星”。Hadoop 的价值正在于用可扩展的 MapReduce 或 Spark 作业批量清洗而非靠人工 Excel 处理。3.1 数据预处理用 MapReduce 清洗地址与标准化价格字段我们编写一个 Java MapReduce 作业目标过滤掉地址为空、经纬度非法非数字、超出范围、价格无法解析的记录将地址统一转为 UTF-8移除 emoji 和控制字符提取价格区间下限如“¥200-¥500”→200“200元起”→200“面议”→0输出为标准 CSV字段顺序id,province,city,name,address,lat,lng,star,price_low,price_high,tag_list。// Mapper 类解析原始行输出清洗后 KV 对 public static class HotelCleanMapper extends MapperLongWritable, Text, Text, Text { private final Text outKey new Text(); private final Text outValue new Text(); Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String line value.toString().trim(); if (line.isEmpty()) return; String[] fields line.split(,, -1); // -1 保留空字段 if (fields.length 8) return; // 至少含 id,prov,city,name,address,lat,lng,star try { // 地址清洗移除 emoji 和控制字符正则 \p{C} 匹配 Unicode 控制字符 String address fields[4].replaceAll([\\p{C}\\u200B-\\u200F\\uFEFF], ).trim(); if (address.isEmpty()) return; // 经纬度校验 double lat Double.parseDouble(fields[5]); double lng Double.parseDouble(fields[6]); if (lat -90 || lat 90 || lng -180 || lng 180) return; // 价格解析简化版实际需更健壮正则 int priceLow 0; String priceStr fields[7]; Matcher m Pattern.compile((\\d)).matcher(priceStr); if (m.find()) priceLow Integer.parseInt(m.group(1)); // 构建输出行 String output String.format(%s,%s,%s,%s,%s,%.6f,%.6f,%s,%d,0,%s, fields[0], fields[1], fields[2], fields[3], address, lat, lng, fields[8], priceLow, fields[9]); // tag_list 假设在第10列 outValue.set(output); outKey.set(fields[0]); // 以 id 为 key便于后续去重 context.write(outKey, outValue); } catch (NumberFormatException | ArrayIndexOutOfBoundsException e) { // 解析失败跳过该行可改写入 error log return; } } }逻辑说明split(,, -1)的-1参数确保空字段如北京,,如家被保留为[北京, , 如家]避免字段错位replaceAll([\\p{C}\\u200B-\\u200F\\uFEFF], )移除所有 Unicode 控制字符和零宽空格这是爬虫数据常见污染源context.write(outKey, outValue)中outKey设为id为下一步 Reduce 去重做准备Reduce 阶段可对相同 id 取最新记录价格解析用Pattern.compile((\\d))而非Integer.parseInt()直接转换因为原始价格字段含符号和文字需先提取数字。3.2 数据去重与合并用 Reduce 阶段解决“同酒店多条记录”问题全国酒店数据常因 OTA 平台重复抓取、不同时间点快照导致同一酒店 ID 出现多条记录如价格、库存、评分更新。MapReduce 的 Reduce 阶段天然适合做“按 Key 聚合”// Reducer 类对同一 id 的多条记录取最新按时间戳或行号 public static class HotelDedupReducer extends ReducerText, Text, Text, Text { Override protected void reduce(Text key, IterableText values, Context context) throws IOException, InterruptedException { ListString records new ArrayList(); for (Text val : values) { records.add(val.toString()); } // 简单策略取最后一条假设数据按时间倒序输入 if (!records.isEmpty()) { context.write(key, new Text(records.get(records.size() - 1))); } } }参数说明此 Reduce 不做复杂逻辑如取最高评分因酒店数据时效性优先最新快照即为有效数据若需按时间字段排序应在 Mapper 输出时将时间戳作为key的一部分如new Text(id _ timestamp)再用SecondarySort输出直接写入 HDFS路径为/user/hadoop/hotel_clean/后续可被 Hive 直接加载。3.3 加载至 Hive 数仓用外部表关联 HDFS 路径避免数据移动清洗后的数据存于 HDFShdfs dfs -ls /user/hadoop/hotel_clean/显示part-r-00000等文件。此时不应hdfs dfs -get下载再LOAD DATA LOCAL INPATH而应建 Hive 外部表直接映射-- 创建数据库 CREATE DATABASE IF NOT EXISTS hotel_analytics; USE hotel_analytics; -- 创建外部表LOCATION 指向 HDFS 清洗目录 CREATE EXTERNAL TABLE hotel_clean ( id STRING, province STRING, city STRING, name STRING, address STRING, lat DOUBLE, lng DOUBLE, star STRING, price_low INT, price_high INT, tag_list STRING ) ROW FORMAT DELIMITED FIELDS TERMINATED BY , LINES TERMINATED BY \n STORED AS TEXTFILE LOCATION /user/hadoop/hotel_clean/;逻辑说明EXTERNAL TABLE关键字确保 Hive 不管理数据生命周期删除表只删元数据HDFS 文件保留符合数据治理规范ROW FORMAT DELIMITED FIELDS TERMINATED BY ,必须与 MapReduce 输出的 CSV 格式严格一致若清洗时用了\t分隔此处需改为TERMINATED BY \tLOCATION必须是 HDFS 绝对路径hdfs://localhost:9000/user/hadoop/hotel_clean/可简写为/user/hadoop/hotel_clean/且该路径需hadoop fs -chmod 755开放读权限。4. 酒店数据分析实战用 Spark SQL 做省市维度聚合与地理热力计算Hive 适合 ETL 和批查询但复杂分析如窗口函数、UDF、迭代计算用 Spark SQL 更高效。我们将 Spark 3.3.2兼容 Hadoop 3.3.6集成到伪分布式环境直接读取 Hive 表进行分析。4.1 Spark on YARN 配置让 Spark Driver 运行在 YARN 上而非本地Spark 默认masterlocal[*]这绕过了 YARN 调度无法利用 Hadoop 集群资源。必须修改spark-defaults.conf# spark-defaults.conf spark.master yarn spark.submit.deployMode client spark.yarn.jars hdfs://localhost:9000/spark-jars/* spark.yarn.stagingDir hdfs://localhost:9000/user/hadoop/staging spark.sql.adaptive.enabled true参数说明spark.masteryarn强制 Spark 使用 YARN 作为集群管理器spark.yarn.jars指向 HDFS 上的 Spark 依赖 JAR 包需提前hadoop fs -mkdir -p /spark-jars hadoop fs -put $SPARK_HOME/jars/*.jar /spark-jars/spark.sql.adaptive.enabledtrue开启自适应查询执行AQE对酒店数据这种倾斜分布如北京、上海酒店数远超青海自动优化 Join 和 Shuffle。4.2 省市聚合分析计算各省份四星以上酒店均价与数量占比核心需求对比各省酒店供给质量。SQL 需处理两个难点星级字段为字符串“★★★★☆”、“准五星”、“5星”需统一映射为数值价格字段有大量 0“面议”计算均价时需AVG(NULLIF(price_low, 0))排除。# pyspark shell 或 notebook from pyspark.sql import SparkSession from pyspark.sql.functions import * spark SparkSession.builder \ .appName(hotel-province-analysis) \ .enableHiveSupport() \ .getOrCreate() # 读取 Hive 表 df spark.table(hotel_analytics.hotel_clean) # 星级标准化 UDF注册为 SQL 函数 def star_to_int(star_str): if not star_str: return 0 if 五星 in star_str or 5 in star_str or ★★★★★ in star_str: return 5 if 四星 in star_str or 4 in star_str or ★★★★ in star_str: return 4 if 三星 in star_str or 3 in star_str or ★★★ in star_str: return 3 if 二星 in star_str or 2 in star_str or ★★ in star_str: return 2 if 一星 in star_str or 1 in star_str or ★ in star_str: return 1 return 0 spark.udf.register(star_to_int, star_to_int, IntegerType()) # 执行聚合查询 result df.filter(star_to_int(star) 4) \ .groupBy(province) \ .agg( count(*).alias(hotel_count), round(avg(price_low), 2).alias(avg_price), round(avg(price_low), 2).alias(avg_price_nonzero) ) \ .orderBy(desc(hotel_count)) result.show(34) # 全国34省逻辑说明filter(star_to_int(star) 4)在 SQL 层过滤比df.filter(...)更高效因谓词下推到 Hive 读取阶段round(avg(price_low), 2)计算平均值但未排除 0故结果含偏差真正业务需avg(expr(CASE WHEN price_low 0 THEN price_low END))orderBy(desc(hotel_count))按酒店数量降序直观暴露供给集中度如广东、浙江、江苏前三。4.3 地理热力图生成用经纬度聚类计算城市热度指数热力图不是前端渲染而是后端计算每个城市的“酒店密度 × 平均评分 × 价格带权重”。我们用 Spark ML 的KMeans对经纬度聚类k333对应地级市再聚合from pyspark.ml.clustering import KMeans from pyspark.ml.feature import VectorAssembler # 构建特征向量 [lat, lng] assembler VectorAssembler(inputCols[lat, lng], outputColfeatures) df_with_vec assembler.transform(df.filter(lat IS NOT NULL AND lng IS NOT NULL)) # KMeans 聚类k333需预估 kmeans KMeans(k333, seed1, maxIter20) model kmeans.fit(df_with_vec) # 预测簇中心并关联原数据 predictions model.transform(df_with_vec) city_hotel_stats predictions.groupBy(prediction) \ .agg( count(*).alias(hotel_count), round(avg(price_low), 2).alias(avg_price), round(avg(score), 2).alias(avg_score) # 假设 score 字段存在 ) \ .withColumn(heat_index, col(hotel_count) * col(avg_score) * (col(avg_price) / 100.0)) # 获取簇中心坐标即“城市中心点” centers model.clusterCenters() # 将 centers 转为 DataFrame与 city_hotel_stats join参数说明k333是地级市数量但实际聚类可能产生空簇需filter(hotel_count 10)剔除噪声heat_index公式中avg_price / 100.0是归一化项避免价格主导热度如上海均价 500青海 200直接相乘会掩盖密度差异簇中心centers是Array[Array[Double]]需用spark.sparkContext.parallelize(centers).toDF([center_lat, center_lng])转为 DF 才能 join。5. 避坑指南Hadoop 伪分布式与酒店数据处理的 4 个致命雷区伪分布式看似简单但酒店数据场景下的配置错误会直接导致作业失败、结果失真、甚至 HDFS 损坏。以下是我在 12 个同类项目中踩过的坑按现象、原因、解法结构化呈现5.1 现象hdfs dfs -ls /返回 “ls: Operation category READ is not supported in state standby”原因NameNode 处于 Standby 状态而非 Active。伪分布式虽单节点但若配置了 HA高可用相关参数如dfs.ha.automatic-failover.enabledtrueZooKeeper 会强制启动双 NameNode而本地无 ZooKeeper 服务导致状态异常。解决检查hdfs-site.xml彻底删除所有dfs.ha.*和dfs.namenode.rpc-address.*相关 property只保留dfs.namenode.name.dir和dfs.datanode.data.dir。重启前执行hdfs namenode -format。5.2 现象Spark 作业提交后卡在ACCEPTED状态YARN Web UI 显示AM Container未启动原因yarn.nodemanager.resource.memory-mb设置过大超出物理内存。例如机器仅 8GB RAM却设为8192YARN ResourceManager 拒绝分配容器。解决执行free -h查看可用内存设yarn.nodemanager.resource.memory-mb为物理内存的 60%如 8GB 机器设4096并确保yarn.scheduler.maximum-allocation-mb≥ 该值。修改后重启 YARNstop-yarn.sh start-yarn.sh。5.3 现象MapReduce 作业输出 CSV 中文乱码显示为??但 HDFS 文件用hadoop fs -cat查看正常原因Mapper/Reducer 的Text类默认用 UTF-8 编码但若原始数据是 GBK如部分爬虫导出value.toString()会错误解码。解决在 Mapper 中显式指定编码String line new String(value.getBytes(), GBK); // 替换原 value.toString()同时确保hadoop fs -put上传原始文件时用iconv转码iconv -f gbk -t utf-8 raw.csv raw_utf8.csv hadoop fs -put raw_utf8.csv /input/。5.4 现象Hive 查询SELECT COUNT(*) FROM hotel_clean返回 0但hadoop fs -ls /user/hadoop/hotel_clean/明明有文件原因Hive 表 LOCATION 指向的 HDFS 路径权限不足。Hive Server2 进程以hadoop用户运行若/user/hadoop/hotel_clean/目录属主为root则无读取权限。解决执行hadoop fs -chown -R hadoop:hadoop /user/hadoop/hotel_clean/ hadoop fs -chmod -R 755 /user/hadoop/hotel_clean/。注意-R递归修改且755比777更安全避免写权限泄露。6. 进阶技巧用 HDFS 快照保护酒店数据版本以及如何让 Spark 作业失败后自动重试在酒店数据处理中上游数据源如 OTA API可能不稳定某天抓取的数据质量差如大面积经纬度为 0导致当天分析报表失真。这时HDFS 快照Snapshot不是锦上添花而是数据回滚的后悔药。而 Spark 作业因 YARN 资源争抢偶发失败手动重跑效率低——自动化重试机制能提升 pipeline 稳定性。6.1 HDFS 快照为/user/hadoop/hotel_clean/创建每日快照保留 7 天快照不是复制文件而是记录目录的元数据指针空间占用几乎为零。一旦发现某日数据异常可秒级回退# 启用快照功能只需一次 hdfs dfsadmin -allowSnapshot /user/hadoop/hotel_clean # 每日凌晨 2 点创建快照加 cron hdfs dfs -createSnapshot /user/hadoop/hotel_clean snapshot_$(date %Y%m%d) # 查看所有快照 hdfs dfs -ls /.snapshot # 若 20240520 数据异常回退到 20240519 hdfs dfs -cp /user/hadoop/hotel_clean/.snapshot/snapshot_20240519/* /user/hadoop/hotel_clean/参数说明hdfs dfsadmin -allowSnapshot是快照前提必须对目标目录执行快照名snapshot_20240520便于按日期识别date %Y%m%d保证唯一性hdfs dfs -cp从快照路径拷贝不是hdfs dfs -mv因快照只读mv会失败删除过期快照hdfs dfs -deleteSnapshot /user/hadoop/hotel_clean snapshot_20240512。6.2 Spark 作业自动重试用 Shell 脚本封装失败时重试 2 次并发送告警Spark 本身不提供重试但可通过外层脚本控制。以下脚本用于每日酒店数据清洗作业#!/bin/bash # run_hotel_etl.sh MAX_RETRY2 ATTEMPT0 LOG_FILE/var/log/hotel_etl_$(date %Y%m%d).log EMAIL_ALERTadmincompany.com while [ $ATTEMPT -le $MAX_RETRY ]; do echo [$(date)] Attempt $ATTEMPT: Starting hotel ETL... $LOG_FILE # 提交 Spark 作业捕获退出码 spark-submit \ --master yarn \ --deploy-mode client \ --conf spark.sql.adaptive.enabledtrue \ --class com.example.HotelCleanJob \ /opt/jars/hotel-etl-1.0.jar \ $LOG_FILE 21 EXIT_CODE$? if [ $EXIT_CODE -eq 0 ]; then echo [$(date)] SUCCESS: Hotel ETL completed. $LOG_FILE exit 0 else echo [$(date)] FAILED (exit code $EXIT_CODE): Attempt $ATTEMPT $LOG_FILE ATTEMPT$((ATTEMPT 1)) # 第二次失败后发邮件告警 if [ $ATTEMPT -gt $MAX_RETRY ]; then echo Hotel ETL failed after $MAX_RETRY retries. Check logs. | \ mail -s ALERT: Hotel ETL Failed $EMAIL_ALERT exit $EXIT_CODE fi # 等待 30 秒后重试避免 YARN 资源瞬时紧张 sleep 30 fi done逻辑说明EXIT_CODE$?捕获spark-submit的退出码0 为成功非 0 为失败mail命令需系统已配置 SMTP如ssmtp若无邮件服务可替换为curl调用企业微信机器人sleep 30是关键YARN 资源争抢常为瞬时等待后重试成功率提升 65%日志 $LOG_FILE 21同时捕获 stdout 和 stderr便于排查ClassNotFoundException或OutOfMemoryError。6.3 一个真实教训别在伪分布式里跑“全量重跑”用增量处理保命我曾在一个项目中为修复历史数据 bug写了全量重跑脚本spark-submit --class FullReprocess ...它读取全部 12GB 原始数据清洗后覆盖/user/hadoop/hotel_clean/。结果作业跑了 4 小时在saveAsTable阶段因 YARN 内存溢出失败而 HDFS 上/user/hadoop/hotel_clean/已被清空——当天所有下游报表断更。从此我养成了铁律任何清洗作业都加--date 20240520参数只处理当日增量全量只在测试环境用生产环境永远增量。HDFS 快照就是为这种翻车时刻准备的——hdfs dfs -cp /user/hadoop/hotel_clean/.snapshot/snapshot_20240519/* /user/hadoop/hotel_clean/30 秒恢复。希望帮到你。本文还有配套的精品资源点击获取