
简介本资源是面向大数据开发工程师与进阶学习者的Hadoop/Spark数据算法实践套件聚焦海量数据场景下的分布式计算原理与工程实现解决算法设计、框架调优及真实任务落地等核心问题。压缩包共876个文件涵盖360个Java核心实现含MapReduce作业逻辑与Spark RDD/DataFrame算子封装、242个依赖JAR包、34个Scala示例代码、63个Markdown技术说明文档以及CSV/TSV/Parquet等多格式测试数据集整体204.27MB结构清晰支持从本地单机调试到集群部署的完整验证链路。已有686人学习下载资源中包含词频统计、ETL流水线、MLlib分类回归模型等典型任务的可运行源码并附带transform.awk等数据预处理脚本及_Success标记文件机制说明便于理解Hadoop输出规范与Spark作业状态管理。1. 这不是又一个“HadoopSpark入门教程”它是一套能直接跑通网约车清洗、交通OD分析、实时日志聚合的生产级源码包专治本地调试失败、集群提交报错、数据倾斜卡死这三类高频翻车现场你是不是也经历过在本地 IDE 里spark-submit跑得好好的一丢到 YARN 上就ClassNotFoundException写了个reduceByKey处理订单流水数据量从 10 万涨到 1000 万就卡在 Stage 3 不动或者hdfs dfs -ls /user/hive/warehouse明明有表spark.sql(select * from traffic_od)却报Table or view not found这不是你代码写得差而是缺一套带真实业务上下文、含完整环境适配逻辑、每行注释都指向具体报错场景的源代码。这份《数据算法 Hadoop/Spark 大数据处理技巧 源代码》不是教学演示玩具——它来自某省会城市网约车监管平台二期的真实工程切片包含 7 个可独立运行的模块伪分布式 HDFS 日志归集脚本、基于 Spark SQL 的交通 OD 矩阵生成器、用 Broadcast Join 优化的司机-车辆维表关联器、带动态分区裁剪的 Hive 分桶写入工具、Spark Streaming 接入 Kafka 的心跳包去重处理器、YARN 资源预估与 Executor 内存反推计算器以及最硬核的——一份把spark-defaults.conf里 23 个关键参数和实际 GC 日志、Shuffle Write 量、Executor Lost 次数做映射的对照表。它适合两类人刚从 Python 数据分析转岗大数据开发需要“抄作业”快速上线的工程师或是已有集群但总在调参、排障、数据倾斜上反复踩坑的熟手。别再看那些“Hello World”式 Demo 了——这次我们直接拆解真实业务里的黑匣子。2. 从单机伪分布到 YARN 集群Hadoop 环境适配不是配置文件搬运而是理解 NameNode 与 DataNode 的心跳契约Hadoop 环境搭建常被简化为“改 core-site.xml、hdfs-site.xml、yarn-site.xml”但真正决定你后续能否跑通 Spark 的是三个隐藏契约NameNode 对 DataNode 心跳超时的容忍阈值、SecondaryNameNode 合并 edits 文件的触发条件、以及 YARN ResourceManager 对 NodeManager 报告状态的采样频率。这份源码包的hadoop-env.sh和hdfs-site.xml并非通用模板而是针对不同硬件做了四档适配开发机4C8G、测试集群16C64G×3、准生产32C128G×5、生产64C256G×10。下面以最常用的开发机档位为例说明关键参数如何联动生效。2.1 伪分布式模式下必须重写的 4 个核心配置项提示所有配置均位于conf/hadoop-dev/目录下start-dfs.sh启动前需执行source conf/hadoop-env.sh加载环境变量否则 JAVA_HOME 不生效导致hadoop version报错。# conf/hadoop-env.sh 关键片段已去除非必要注释 export JAVA_HOME/usr/lib/jvm/java-11-openjdk-amd64 export HADOOP_HEAPSIZE_MAX4096 # 强制 NameNode JVM 堆上限为 4G避免 OOM 后无法响应 DataNode 心跳 export HADOOP_NAMENODE_OPTS-XX:UseG1GC -XX:MaxGCPauseMillis200 export HADOOP_DATANODE_OPTS-XX:UseG1GC -XX:MaxGCPauseMillis100 -Xms2048m -Xmx2048m这段配置的玄机在于DataNode 的-Xms和-Xmx必须严格相等2048m否则在伪分布式下DataNode 启动后会因内存抖动被 NameNode 主动踢出节点列表表现为jps看不到 DataNode 进程但hdfs dfsadmin -report显示Live datanodes: 0。这是很多新手查遍日志都找不到原因的血泪经验。!-- conf/hdfs-site.xml 核心片段 -- property namedfs.namenode.handler.count/name value10/value descriptionNameNode 处理客户端 RPC 请求的线程数。开发机设为 10生产环境需按 CPU 核数×2 设置/description /property property namedfs.datanode.max.transfer.threads/name value4096/value descriptionDataNode 同时处理 Block 传输的线程上限。若 Spark 任务并发读写 HDFS 时出现 Too many open files优先调大此值/description /property property namedfs.client.use.datanode.hostname/name valuefalse/value description强制客户端通过 IP 访问 DataNode。若设为 true在 Docker 或虚拟机中易因 hostname 解析失败导致 Block 读取超时/description /property property namedfs.namenode.avoid.stale.datanode/name valuetrue/value description启用过期 DataNode 避让机制。当 DataNode 心跳中断超过 dfs.namenode.stale.datanode.interval.ms默认30sNameNode 将不再向其分配新 Block/description /propertydfs.client.use.datanode.hostnamefalse是关键中的关键。很多教程教你在/etc/hosts里绑定localhost到127.0.0.1但 Spark Driver 在提交任务时会从 HDFS Client 获取 DataNode 的 hostname若该 hostname 是ubuntu而非127.0.0.1且未在 hosts 中映射就会触发 DNS 查询超时最终表现为BlockMissingException。这个坑我踩过三次每次重装系统都要记一遍。2.2 为什么start-dfs.sh后hdfs dfs -ls /仍报 Connection refused现象执行./sbin/start-dfs.sh后jps显示 NameNode 和 DataNode 进程存在但hdfs dfs -ls /返回Connection refused。原因NameNode 默认监听0.0.0.0:9000但core-site.xml中fs.defaultFS若配置为hdfs://localhost:9000而localhost在部分 Linux 发行版中解析为::1IPv6 地址导致客户端尝试用 IPv6 连接但 NameNode 未开启 IPv6 支持。解决将core-site.xml中的fs.defaultFS明确改为hdfs://127.0.0.1:9000并确保hadoop-env.sh中HADOOP_OPTS包含-Djava.net.preferIPv4Stacktrue# conf/hadoop-env.sh 追加 export HADOOP_OPTS$HADOOP_OPTS -Djava.net.preferIPv4Stacktrue验证方法netstat -tuln | grep :9000应显示127.0.0.1:9000而非:::9000。2.3 YARN 集群提交失败的底层链路排查法Spark 任务提交到 YARN 失败90% 的情况不是 Spark 代码问题而是 YARN ResourceManager 与 NodeManager 的状态契约断裂。源码包中bin/check-yarn-health.sh提供了一套链路检测脚本#!/bin/bash # bin/check-yarn-health.sh echo Step 1: Check ResourceManager status curl -s http://localhost:8088/ws/v1/cluster/info | jq .clusterInfo.state 2/dev/null || echo RM not responding echo Step 2: Check NodeManager registration yarn node -list 2/dev/null | grep -E (RUNNING|DECOMMISSIONED) | wc -l | xargs -I{} sh -c if [ {} -eq 0 ]; then echo No active NodeManager; else echo Active NodeManagers: {}; fi echo Step 3: Check NM log for heartbeat errors if [ -f $HADOOP_LOG_DIR/yarn-*-nodemanager-*.log ]; then tail -50 $HADOOP_LOG_DIR/yarn-*-nodemanager-*.log | grep -i heartbeat\|registration\|unregister | tail -5 else echo NM log not found, check \$HADOOP_LOG_DIR fi这个脚本直击要害第一步确认 RM Web UI 是否存活很多教程忽略这步直接spark-submit导致连接超时第二步用yarn node -list查看 NM 注册状态比jps更可靠NM 进程存在但未注册成功很常见第三步抓取 NM 日志中与心跳相关的关键词因为 NM 每 10 秒向 RM 发送一次心跳若连续 3 次失败RM 会将其标记为LOST此时yarn node -list仍显示RUNNING但实际已不可用。这是集群环境下最隐蔽的坑。3. Spark 任务从本地模式到集群模式的 5 个必改点序列化、依赖传递、路径协议、资源申请、Shuffle 策略本地spark-shell能跑通的代码提交到集群十有八九失败。这不是 Spark 的 Bug而是本地模式绕过了分布式环境的四大约束JVM 类加载隔离、网络传输序列化、分布式文件系统路径解析、Executor 资源调度。这份源码包的spark-submit脚本不是简单封装而是内置了 5 层校验逻辑确保你在敲下回车前就已经规避了 95% 的常见错误。3.1 Kryo 序列化注册表为什么mapPartitions里 new 个对象就报NotSerializableException现象本地spark-shell中rdd.mapPartitions(iter iter.map(x new TrafficRecord(x)))正常集群提交后报org.apache.spark.SparkException: Task not serializable。原因TrafficRecord类未实现java.io.Serializable且未在 Kryo 注册表中声明。Spark 默认使用 Java 序列化对匿名函数捕获的外部对象要求极其苛刻而 Kryo 虽快但必须显式注册类否则反序列化失败。解决源码包中src/main/scala/com/example/spark/serializer/KryoRegistrator.scala提供了标准注册模板// src/main/scala/com/example/spark/serializer/KryoRegistrator.scala import com.esotericsoftware.kryo.Kryo import org.apache.spark.serializer.KryoRegistrator class TrafficKryoRegistrator extends KryoRegistrator { override def registerClasses(kryo: Kryo): Unit { // 必须注册所有在闭包中使用的自定义类 kryo.register(classOf[TrafficRecord]) kryo.register(classOf[DriverInfo]) kryo.register(classOf[VehicleStatus]) // 针对 Scala 集合注册特定实现类而非 trait kryo.register(classOf[scala.collection.immutable.List$Nil.type]) kryo.register(classOf[scala.collection.immutable.$colon$colon[_]]) // 注册常用第三方类如 Joda-Time kryo.register(classOf[org.joda.time.DateTime]) } }使用时在spark-submit中指定spark-submit \ --conf spark.serializerorg.apache.spark.serializer.KryoSerializer \ --conf spark.kryo.registratorcom.example.spark.serializer.TrafficKryoRegistrator \ --conf spark.kryo.registrationRequiredtrue \ # 强制注册避免漏注册导致运行时报错 --class com.example.spark.od.ODMatrixGenerator \ target/traffic-spark-1.0.jarspark.kryo.registrationRequiredtrue是关键开关。它让 Kryo 在反序列化时校验类是否已注册若未注册则立即抛异常而不是等到任务运行中才失败——这极大缩短了排错周期。3.2 依赖包传递--jars和--driver-class-path的生死时速现象代码中用了com.typesafe.config:config:1.4.2本地没问题集群提交后java.lang.NoClassDefFoundError: com/typesafe/config/Config。原因--jars只将 jar 包分发到 Executor ClassPathDriver 端仍需独立加载而--driver-class-path仅影响 Driver不传递给 Executor。两者必须配合使用。解决源码包中bin/submit-od-matrix.sh给出了工业级写法#!/bin/bash # bin/submit-od-matrix.sh APP_JARtarget/traffic-spark-1.0.jar LIBSlib/config-1.4.2.jar,lib/slf4j-log4j12-1.7.32.jar,lib/joda-time-2.10.13.jar spark-submit \ --master yarn \ --deploy-mode cluster \ --name OD-Matrix-Generator \ --jars $LIBS \ # 分发到所有 Executor --driver-class-path $LIBS \ # 加载到 Driver ClassPath --conf spark.driver.extraClassPath$LIBS \ # 兼容旧版本 Spark 的写法 --conf spark.executor.extraClassPath$LIBS \ # 确保 Executor 也能加载 --class com.example.spark.od.ODMatrixGenerator \ $APP_JAR注意--jars参数值是逗号分隔的路径不能有空格--driver-class-path和--conf spark.driver.extraClassPath功能重复但保留双保险因为不同 Spark 版本对此参数的支持有差异Spark 2.x 之后推荐用--conf。3.3 HDFS 路径协议陷阱file://、hdfs://、/的三重幻觉现象spark.read.parquet(data/od_raw)本地能读集群提交后报java.io.FileNotFoundException: File does not exist: hdfs://namenode:9000/user/spark/data/od_raw。原因Spark 对路径的解析规则是若路径以file://开头走本地文件系统以hdfs://开头走 HDFS若无协议前缀如data/od_raw则根据spark.sql.warehouse.dir配置决定——本地模式默认为file:///tmp/spark-warehouse集群模式默认为hdfs://namenode:9000/user/hive/warehouse。解决源码包中所有路径均采用绝对协议路径并提供PathResolver工具类统一管理// src/main/scala/com/example/spark/util/PathResolver.scala object PathResolver { val HDFS_ROOT hdfs://namenode:9000 val WAREHOUSE_ROOT s$HDFS_ROOT/user/hive/warehouse def getRawDataPath(date: String): String s$HDFS_ROOT/data/raw/traffic/$date def getOdOutputPath(date: String): String s$WAREHOUSE_ROOT/od_matrix/date$date def getDriverDimPath: String s$HDFS_ROOT/data/dim/driver_info }在业务代码中强制使用val rawDF spark.read.parquet(PathResolver.getRawDataPath(20231001)) val dimDF spark.read.parquet(PathResolver.getDriverDimPath)这样彻底规避了相对路径带来的不确定性。另外spark.sql.warehouse.dir必须在spark-defaults.conf中显式设置为hdfs://namenode:9000/user/hive/warehouse否则CREATE TABLE语句会创建在本地/tmp下导致 Hive Metastore 找不到表。4. 数据倾斜的 4 种实战解法不是加盐而是看清 Shuffle 的本质是 Key 分布与 Partitioner 的博弈数据倾斜不是“加个随机前缀就能解决”的玄学而是 Key 的哈希分布与 Spark 默认 HashPartitioner 的分区策略不匹配导致的资源浪费。这份源码包的src/main/scala/com/example/spark/tilt/目录下提供了 4 种经过生产验证的解法每种都附带倾斜 Key 的自动识别脚本和效果对比报告。4.1 倾斜 Key 识别用samplecountByValue替代全量groupByKey现象rdd.groupByKey().mapValues(_.size)卡死因为groupByKey会将所有相同 Key 的 Value 拉到同一个 Partition若某个 Key 出现 100 万次该 Partition 就会 OOM。原因全量统计 Key 频次成本太高且无法提前干预。解决源码包中bin/detect-skew-keys.sh使用采样法快速定位 Top 10 倾斜 Key#!/bin/bash # bin/detect-skew-keys.sh INPUT_PATHhdfs://namenode:9000/data/raw/traffic/20231001 OUTPUT_PATHhdfs://namenode:9000/tmp/skew_keys spark-submit \ --class com.example.spark.tilt.KeySkewDetector \ --master yarn \ --conf spark.sql.adaptive.enabledtrue \ target/traffic-spark-1.0.jar \ $INPUT_PATH $OUTPUT_PATH 0.01 # 采样率 1%对应的 Scala 代码// src/main/scala/com/example/spark/tilt/KeySkewDetector.scala def detectSkewKeys(inputPath: String, outputPath: String, sampleFraction: Double)(implicit spark: SparkSession): Unit { import spark.implicits._ // 1. 采样并提取 Key假设是订单 ID val sampledKeys spark.read.parquet(inputPath) .sample(withReplacement false, fraction sampleFraction) .select(order_id) .as[String] .rdd .map(key (key, 1)) .reduceByKey(_ _) .map { case (key, count) (count, key) } .sortByKey(ascending false) .map { case (count, key) (key, count) } .toDF(key, count) // 2. 保存 Top 100 倾斜 Key 供后续处理 sampledKeys.limit(100).write.mode(overwrite).parquet(outputPath) }采样率0.01是经验值对 10 亿条数据采样 1000 万条足够反映分布趋势且耗时控制在 2 分钟内。全量统计可能要 2 小时以上。4.2 方案一Salting加盐——但盐值必须可控不能真随机现象“加盐”后任务变慢甚至更倾斜。原因随机盐值如new Random().nextInt(100)会导致原本一个 Key 被打散到 100 个 Salted Key但这些 Salted Key 的数据量可能极不均衡比如key_A_1有 50 万key_A_2只有 100反而加剧局部倾斜。解决源码包中SaltingStrategy.scala实现了“权重盐值”——根据 Key 的预估频次动态分配盐值数量// src/main/scala/com/example/spark/tilt/SaltingStrategy.scala case class SkewKeyInfo(key: String, estimatedCount: Long, saltRange: Int) object SaltingStrategy { // 从 detect-skew-keys.sh 输出的 parquet 中加载倾斜 Key 信息 def loadSkewKeys(spark: SparkSession, skewPath: String): Map[String, SkewKeyInfo] { val skewDF spark.read.parquet(skewPath) .filter($count 10000) // 阈值可调 .withColumn(saltRange, when($count 1000000, 100) .when($count 100000, 20) .otherwise(5)) .select(key, count, saltRange) .as[(String, Long, Int)] .collect() .map { case (k, c, r) k - SkewKeyInfo(k, c, r) } .toMap } def saltKey(key: String, skewMap: Map[String, SkewKeyInfo]): String { skewMap.get(key) match { case Some(info) val salt (System.currentTimeMillis() % info.saltRange).toInt s$key#$salt case None key } } }使用时val skewMap SaltingStrategy.loadSkewKeys(spark, /tmp/skew_keys) val saltedRDD rawRDD.map { case (key, value) (SaltingStrategy.saltKey(key, skewMap), value) } val aggregated saltedRDD.reduceByKey(_ _) .map { case (saltedKey, sum) val originalKey saltedKey.split(#)(0) (originalKey, sum) } .reduceByKey(_ _) // 最终合并盐值范围saltRange根据预估频次动态设定确保每个 Salted Key 的数据量落在 1~5 万区间这才是真正的“可控加盐”。4.3 方案二Broadcast Join —— 当维表小于 10MB 时这是唯一正解现象orders.join(drivers, driver_id)任务卡在 Shuffle Read 阶段。原因drivers维表虽小10 万行但若未广播Spark 会将其作为大表参与 Shuffle导致所有 Executor 都要拉取全量维表数据。解决源码包中BroadcastJoinOptimizer.scala自动判断维表大小并广播// src/main/scala/com/example/spark/tilt/BroadcastJoinOptimizer.scala def broadcastJoinIfSmall(left: DataFrame, right: DataFrame, joinCol: String): DataFrame { val rightSize right.count() * right.schema.fields.map(_.dataType.defaultSize).sum if (rightSize 10 * 1024 * 1024) { // 小于 10MB val broadcastDF spark.sparkContext.broadcast(right.collect()) left.mapPartitions { iter val dimMap broadcastDF.value.map(r (r.getAs[String](joinCol), r)).toMap iter.map(row { val key row.getAs[String](joinCol) dimMap.get(key).map(dimRow Row.merge(row, dimRow)).getOrElse(row) }) }.toDF(left.schema right.schema) } else { left.join(right, joinCol) } }注意right.count()是行动操作会触发一次 Job但相比后续的 Shuffle Join这点开销微不足道。且defaultSize是 Spark 内部估算的每行字节数足够用于 10MB 量级的粗略判断。5. 避坑Hadoop/Spark 生产环境 5 个高频翻车点与血泪修复方案注意以下问题均来自真实线上事故非理论推测。每一条都对应源码包中docs/troubleshooting.md的详细复现步骤和修复验证命令。5.1 现象spark-submit提交后YARN Web UI 显示 Application 状态为ACCEPTED但 5 分钟后变为FAILED日志中无有效错误信息原因YARN ResourceManager 的yarn.scheduler.maximum-allocation-mb默认 8192小于 Spark 申请的 Executor 内存如--executor-memory 10g导致资源无法分配Application 卡在ACCEPTED状态直至超时。解决检查yarn-site.xml中yarn.scheduler.maximum-allocation-mb和yarn.scheduler.maximum-allocation-vcores确保其大于 Spark 提交参数# 查看当前 YARN 资源上限 yarn rmadmin -getGroups 2/dev/null | head -1 | xargs -I{} yarn scheduler -status # 或直接查配置 grep maximum-allocation $HADOOP_CONF_DIR/yarn-site.xml若需调整修改yarn-site.xmlproperty nameyarn.scheduler.maximum-allocation-mb/name value16384/value !-- 提升至 16G -- /property property nameyarn.scheduler.maximum-allocation-vcores/name value8/value /property然后重启 ResourceManager$HADOOP_HOME/sbin/yarn-daemon.sh stop resourcemanager $HADOOP_HOME/sbin/yarn-daemon.sh start resourcemanager。5.2 现象Spark Streaming 任务运行 2 小时后StreamingListener报ReceiverTracker: Receiver is stopped且 Kafka 消费位点停滞原因spark.streaming.receiver.writeAheadLog.enabletrue开启了 WAL但spark.streaming.driver.writeAheadLog.closeFileAfterWrite未设置导致 WAL 文件持续增长Driver 磁盘写满后崩溃。解决在spark-defaults.conf中强制关闭 WAL 文件句柄spark.streaming.receiver.writeAheadLog.enabletrue spark.streaming.driver.writeAheadLog.closeFileAfterWritetrue spark.streaming.receiver.writeAheadLog.closeFileAfterWritetrue并设置 WAL 存储路径为独立磁盘spark.streaming.driver.writeAheadLog.dirhdfs://namenode:9000/tmp/spark-wal/driver spark.streaming.receiver.writeAheadLog.dirhdfs://namenode:9000/tmp/spark-wal/receiver5.3 现象Hive 表INSERT OVERWRITE后Spark SQL 查询返回空结果但hive -e select * from table能查到数据原因Spark 使用的是自己的 Hive Metastore Client与 Hive CLI 的缓存机制不同。当 Hive 表结构变更如新增分区后Spark 未刷新元数据缓存。解决在 Spark SQL 中执行强制刷新REFRESH TABLE traffic_od; -- 刷新表级元数据 MSCK REPAIR TABLE traffic_od; -- 修复分区适用于按日期自动发现分区或在代码中调用spark.catalog.refreshTable(traffic_od) spark.sql(MSCK REPAIR TABLE traffic_od)5.4 现象hdfs dfs -du -h /user/hive/warehouse显示某表目录大小为 200GB但spark.sql(select count(*) from table).show()返回 0原因该表是外部表EXTERNAL TABLE且 HDFS 上的数据文件被手动删除但 Hive Metastore 中的分区元数据未同步清理导致 Spark 读取时找不到文件。解决先确认表类型DESCRIBE FORMATTED traffic_od;若Type: EXTERNAL_TABLE则用SHOW PARTITIONS查看分区再用ALTER TABLE ... DROP PARTITION清理无效分区SHOW PARTITIONS traffic_od; ALTER TABLE traffic_od DROP IF EXISTS PARTITION (dt20231001);5.5 现象Spark UI 的 Executors 页面显示Used Memory为 0但Storage页面有大量缓存且任务频繁 GC原因spark.memory.fraction默认 0.6设置过高挤压了spark.memory.storageFraction默认 0.5的空间导致存储内存不足缓存数据被频繁驱逐引发 GC。解决调低spark.memory.fraction增大spark.memory.storageFractionspark.memory.fraction0.5 spark.memory.storageFraction0.6并监控Storage页面的Memory Used和Disk Used比例理想状态是 Memory Used 占比 70%Disk Used 10%。6. 用spark-sql命令行做生产级数据探查不是select * limit 10而是构建可复用的探查流水线很多人把spark-sql当成临时查询工具但它其实是 Spark 最被低估的生产力引擎——只要配上正确的配置、UDF 和探查脚本它就能替代 80% 的临时数据分析需求。这份源码包的sql/目录下藏着一套完整的探查流水线从自动识别字段类型、计算空值率、检测数据倾斜到生成建表 DDL 和分区建议全部用纯 SQL 内置函数实现无需写一行 Scala。6.1 自动化字段探查DESCRIBE DETAILANALYZE TABLE的组合拳传统做法是SELECT COUNT(*), COUNT(col1), COUNT(col2) FROM table手动算空值率效率低且无法覆盖所有字段。源码包中sql/probe-schema.sql利用 Spark 3.0 的DESCRIBE DETAIL和ANALYZE TABLE实现一键探查-- sql/probe-schema.sql -- 第一步获取表基础信息位置、格式、分区 DESCRIBE DETAIL traffic_od; -- 第二步强制收集统计信息需 Spark 3.0 ANALYZE TABLE traffic_od COMPUTE STATISTICS FOR ALL COLUMNS; -- 第三步查询统计信息视图Spark 3.2 支持 SELECT col_name, data_type, min, max, num_nulls, num_distincts, avg_col_len, max_col_len FROM system.table_columns WHERE table_catalog spark_catalog AND table_schema default AND table_name traffic_od ORDER BY num_nulls DESC;执行方式spark-sql \ --master yarn \ --conf spark.sql.adaptive.enabledtrue \ --conf spark.sql.adaptive.coalescePartitions.enabledtrue \ -f sql/probe-schema.sqlsystem.table_columns是 Spark 3.2 引入的系统表无需额外部署 Hive Metastore直接暴露列级统计信息。num_nulls和num_distincts是精确值非采样前提是执行了ANALYZE TABLE。6.2 数据倾斜热力图用percentile_approx定位 Key 分布拐点COUNT(DISTINCT key)只能告诉你有多少唯一值但无法揭示分布形态。源码包中sql/skew-heatmap.sql用percentile_approx生成 Key 频次的分位数热力图-- sql/skew-heatmap.sql WITH key_counts AS ( SELECT order_id, COUNT(*) as cnt FROM traffic_od GROUP BY order_id ), stats AS ( SELECT percentile_approx(cnt, 0.5) as median_cnt, percentile_approx(cnt, 0.9) as p90_cnt, percentile_approx(cnt, 0.95) as p95_cnt, percentile_approx(cnt, 0.99) as p99_cnt, percentile_approx(cnt, 0.999) as p999_cnt FROM key_counts ) SELECT Median as percentile, median_cnt as count FROM stats UNION ALL SELECT P90, p90_cnt FROM stats UNION ALL SELECT P95, p95_cnt FROM stats UNION ALL SELECT P99, p99_cnt FROM stats UNION ALL SELECT P999, p999_cnt FROM stats;输出示例percentile count Median 12 P90 85 P95 192 P99 1247 P999 18532这说明99% 的订单 ID 出现次数 ≤ 1247 次但 0.1% 的订单 ID 出现次数高达 18532 次——这就是典型的长尾分布P999 值是中位数的 1500 倍必须针对性处理。6.3 自动生成建表 DDL从 Parquet 文件反推 Schema 并添加分区字段当你拿到一个原始 Parquet 目录如hdfs://namenode:9000/data/raw/traffic/20231001想快速建 Hive 表传统做法是spark.read.parquet(...).printSchema()然后手写 DDL。源码包中bin/gen-ddl-from-parquet.sh用spark-sql的CREATE TABLE USING语法一键生成#!/bin/bash # bin/gen-ddl-from-parquet.sh INPUT_PATHhdfs://namenode:9000/data/raw/traffic/20231001 TABLE_NAMEtraffic_raw PARTITION_COLSdt STRING spark-sql \ --master yarn \ -e CREATE TABLE IF NOT EXISTS $TABLE_NAME USING PARQUET LOCATION $INPUT_PATH TBLPROPERTIES (parquet.compressionSNAPPY); -- p a hrefhttps://download.csdn.net/download/khxu666/10407207 stylecolor:#ec7500;font-size:14px; 本文还有配套的精品资源点击获取 /a img altmenu-r.4af5f7ec.gif srchttps://csdnimg.cn/release/wenkucmsfe/public/img/menu-r.4af5f7ec.gif stylewidth:16px;margin-left:4px;vertical-align:text-bottom;cursor:text; /p