
简介这份资源包聚焦 Hadoop 与 Spark 两大框架的数据算法实现主要面向大数据开发入门者、需要实战源码参考的工程师以及正在准备大数据相关岗位面试的学习者也适合用于课程设计与毕业设计参考。内容涵盖 MapReduce 的 Map/Reduce 函数实例、RDD 与 DataFrame API 应用以及 Spark SQL、Streaming、MLlib 等模块的算法代码并附有可运行的数据集可直接验证词频统计、分类、回归、聚类等典型场景帮助读者从数据预处理走到模型训练与结果分析。压缩包共有 876 个文件体积约 204MB以 Java 源码360 个、JAR 依赖包242 个、Scala 脚本34 个、Shell 脚本31 个和 Markdown 文档63 个为主同时包含若干 CSV、TSV、PDF 等辅助资料目录结构清晰完整便于按模块检索。已有 686 人学习下载对想系统掌握 Hadoop/Spark 编程范式、提升大数据处理实战能力的学习者而言是一份兼具代码参考与实验素材的实用资源。1. 数据算法 Hadoop/Spark 大数据处理技巧 源代码到底该从哪一层开始啃搜索「数据算法 Hadoop/Spark 大数据处理技巧 源代码」的人多半不是想系统学框架而是手里压着具体需求洗数据、求 Top-N、做关联、去重、组内排序。这类源码包能不能用取决于你从哪一层开始啃。直接通读全部源码很容易被工程胶水代码淹没按高频算法骨架去拆才是最短路径。我按自己的落地顺序把这篇笔记分成四段先讲 MapReduce 侧的数据流向控制再讲 Spark 侧的同款实现接着拆高频算法骨架最后铺集群上才遇得到的坑和验证手段。适合写过 WordCount、想往算法层再走一步的工程师也适合准备把开源算法包改造成内部公共库的团队。2. MapReduce 算法源码从零搭Top-N 如何用几十行代码控住数据流向2.1 为什么说数据算法本质是控制数据流向而不是调 APIMapReduce 的算法能力并不在 API 本身真正决定效率的是数据在 map、shuffle、reduce 之间的流向。每一个 shuffle 边界都伴随一次全量网络传输所谓算法技巧本质上是「在 shuffle 之前把数据量压下去」。Top-N 是最典型也最好用的例子。天真做法是 mapper 把每条记录都发出去reducer 拿全量数据排一次序取前 N 个。代价是 O(全部记录) 的 shuffle数据量一上来就非常被动。换一个写法每个 mapper 在 map 阶段维护一个容量为 N 的局部堆只输出 N 条记录reducer 收到的总量从「全量」降为「mapper 数 × N」。假设 100 个 mapper、N 取 10shuffle 量从几亿条降到一千条这是数量级的差距。这套「map 端剪枝 reduce 端汇总」的模式不只是 Top-N 专用。去重、频次统计、布隆过滤、采样分桶底层都是同一个思路能早收窄就早收窄不要把所有原始数据搬到下游再做决定。理解这一点之后再去看任何一份 Hadoop 算法源码你会先找它的 shuffle 边界在哪、每个 map 任务的输出量是多少而不是先看它调用了哪个类。2.2 手写 Mapper 端局部 Top-N 与 Reducer 端全局 Top-N一份最小可跑的 Top-N 源码核心就是两个类。Mapper 端维护局部堆Reducer 端做全局合并。下面这份代码是常见的实现方式我在注释里标了每个关键步骤的理由。import org.apache.hadoop.io.LongWritable; import org.apache.hadoop.io.NullWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Mapper; import java.io.IOException; import java.util.TreeMap; public class TopNMapper extends MapperLongWritable, Text, NullWritable, Text { // TreeMap 按键升序排列容量限制为 N等价于一个最小堆 private final TreeMapInteger, Text localTop new TreeMap(); private int n 10; Override protected void setup(Context context) { // 从 job 配置里读 top.n避免改业务值时重新编译 n context.getConfiguration().getInt(top.n, 10); } Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { // 假设输入格式itemId,timestamp,score String[] fields value.toString().split(,); if (fields.length 3) { return; // 脏数据直接跳过否则下一行 parseInt 会抛异常 } int score; try { score Integer.parseInt(fields[2].trim()); } catch (NumberFormatException e) { return; // 分数列不是数字同样丢弃 } localTop.put(score, new Text(value)); if (localTop.size() n) { localTop.remove(localTop.firstKey()); // 移除当前最小键堆容量恒为 N } } Override protected void cleanup(Context context) throws IOException, InterruptedException { // 每个 mapper 只输出自己分片内的 Top-N for (Text v : localTop.values()) { context.write(NullWritable.get(), v); } } }Reducers 的写法几乎一样区别在于它收到的是所有 mapper 剪枝后的结果相当于做最后一轮全局合并。import org.apache.hadoop.io.NullWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Reducer; import java.io.IOException; import java.util.TreeMap; public class TopNReducer extends ReducerNullWritable, Text, NullWritable, Text { private final TreeMapInteger, Text globalTop new TreeMap(); private int n 10; Override protected void setup(Context context) { n context.getConfiguration().getInt(top.n, 10); } Override protected void reduce(NullWritable key, IterableText values, Context context) throws IOException, InterruptedException { for (Text v : values) { String[] fields v.toString().split(,); if (fields.length 3) { continue; } int score Integer.parseInt(fields[2].trim()); globalTop.put(score, new Text(v)); if (globalTop.size() n) { globalTop.remove(globalTop.firstKey()); } } } Override protected void cleanup(Context context) throws IOException, InterruptedException { // descendingMap 让分数从高到低输出 for (Text v : globalTop.descendingMap().values()) { context.write(NullWritable.get(), v); } } }这段代码有几个值得注意的细节。第一map 方法里用new Text(value)而不是直接存字符串是为了防止后续修改 value 对象影响堆里已存的数据Hadoop 的 Text 是可变的直接存引用很容易踩到缓存复用导致的脏数据。第二TreeMap 的 key 是 score如果两条记录 score 相同后一条会覆盖前一条这是后面避坑章节要展开的重点。第三所有 key 都写 NullWritable意味着所有数据被哈希到同一个分区全局 Top-N 的正确性靠这个保证代价是只有一个 reducer。2.3 三个必调参数Combiner、Partitioner 与 Reduce 并行度源码写好之后驱动类的参数设置直接决定这个算法在集群上会不会翻车。我一般会把这个 Job 的配置写在独立类里方便套不同的输入输出路径。import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.NullWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; public class TopNJob { public static void main(String[] args) throws Exception { Configuration conf new Configuration(); conf.setInt(top.n, 20); // 业务侧只需改这一行 Job job Job.getInstance(conf, top-n-v1); job.setJarByClass(TopNJob.class); job.setMapperClass(TopNMapper.class); // Reducer 类直接复用为 Combiner因为“取Top-N”是幂等操作 job.setCombinerClass(TopNReducer.class); job.setReducerClass(TopNReducer.class); job.setNumReduceTasks(1); // 全局 Top-N 必须单 reducer job.setOutputKeyClass(NullWritable.class); job.setOutputValueClass(Text.class); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); System.exit(job.waitForCompletion(true) ? 0 : 1); } }Combiner 复用 TopNReducer 在前文提过这里再展开一次Combiner 在 map 端可能执行 0 次、1 次或多次只有算法操作幂等时才能放心复用。求和、取最大最小、取 Top-N 都是幂等的反复执行结果一致但像「求平均值」这种依赖中间态的就不行需要先拆成 sum 和 count 两列的复合结构。遇到从源码包里抄 Combiner 时先确认这一条否则结果会莫名其妙地偏差。Partitioner 在这里用的是默认实现NullWritable 的哈希值恒定所以天然全进一个 reducer。但要注意一旦你想把 Top-N 改成「每个分类维护一个 Top-N」就不能再依赖默认分区器了必须自定义 Partitioner 按分类字段分流这就是下一章二次排序里会讲的三角关系。Reduce 并行度设 1 是正确性优先的取舍数据量破亿后单 reducer 会明显变慢届时就要改成「组内 Top-N 外层再合并」的两阶段方案不要硬扛。3. 用 Spark 重写同一套算法RDD 算子的执行语义比 API 名称更值得吃透3.1 Spark 与 Hadoop 算法表达的差异阶段切分与数据血缘同一个 Top-N 算法在 Spark 里写起来更短但短不等于没讲究。MapReduce 是固定的「map 完必须 reduce」的阶段链Spark 则把作业拆成 DAG宽依赖算子groupByKey、sortByKey、repartition、distinct会触发 shuffle 并切开一个新的 stage。理解这一点你才能解释为什么同样一段代码把某个算子换掉之后作业变快了还是变慢了。Spark 的另一个优势是数据血缘和缓存。一份 RDD 即使被多个动作复用只要你不显式 cache它就会在每次行动时从头重算这是个很容易被忽略的内存坑。算法代码里如果有「先算出一个中间集合、后面循环十次都要用」的场景记得在中间集合上加.cache()否则你以为的复用其实每次都在全量重算数据量大时比 MapReduce 还慢就一点不奇怪。3.2 一份与 MapReduce 对照的 Spark Top-N 代码下面是常见做法里性能较好的一版核心思路和 Hadoop 版一致先在各分区内剪枝再做全局取前 N。from pyspark.rdd import RDD import heapq def top_n(rdd: RDD, n: int 10): 返回全局 Top-N结果为 [(score, itemId)]按 score 降序。 输入 RDD 中每个元素形如 (itemId, timestamp, score)。 def partition_top(iterator): # 每个分区内部先做一次局部剪枝shuffle 量从全量降为 n * 分区数 local heapq.nlargest(n, iterator, keylambda row: row[2]) return local local_top rdd.mapPartitions(partition_top) # takeOrdered 只取前 n 个做的是部分排序不会落一个全量有序集合 result (local_top .map(lambda row: (row[2], row[0])) # (score, itemId) .takeOrdered(n, keylambda pair: -pair[0])) return resultmapPartitions接收的是一个迭代器函数它对整个分区批量处理比map逐条处理更适合做「开堆、灌数据、取结果」这类有状态逻辑。heapq.nlargest内部就是最小堆和 TreeMap 剪枝是同一个套路。takeOrdered(n, key...)是这里最值得记的一个算子它只保证返回前 n 个内部实现是有界堆不产生全量排序结果。keylambda pair: -pair[0]表示按 score 降序注意如果不取负号默认就是升序取最小。对照反面写法很多人图省事直接写成# 反面示例全量排序再 takeN 很小的时候代价不成比例 result (rdd.map(lambda row: (row[2], row[0])) .sortByKey(ascendingFalse) .take(n))sortByKey是一次全量 shuffle 加全量排序即使最终只取 10 条也会把所有数据排好序落盘。数据量一上来这个写法的 shuffle 量比mapPartitions takeOrdered高出好几个数量级。选型时可以抓一条原则只需要前 N 条时永远别做全量排序只有需要看到「完整有序列表」时才用 sort。3.3 数据倾斜在源码层面长什么样从作业日志到加盐代码Spark 数据倾斜最常见的现场是同一个 stage 里大部分 task 几十秒跑完一两个 task 卡了半个多小时去看 Spark UI 的 stage 页会发现某个 task 的 shuffle read 量是平均值的几十倍。源码层面解决倾斜绕不开加盐salting和重新分区这两板斧。import random def add_salt(rdd, hot_key, salt_parts16): 把热点 key 拆成 salt_parts 个带随机后缀的伪 key。 适用于聚合类算法如果算法要求全局有序不能直接加盐。 def transform(row): key, value row # 假设 RDD 元素是 (key, value) if key hot_key: return (f{key}#{random.randint(0, salt_parts - 1)}, value) # 非热点 key 原样保留避免增加无谓的 shuffle return row return rdd.map(transform)加盐之后原本压在一个 reducer 上的热点 key 被分散到 k 个任务上并行度立刻提上来。但要注意适用边界加盐只适用于「先拆后合」的算法——比如计数、求和、去重计数、求各自 Top-N 再归并如果你的算法是全局 Top-N 或者要求全局唯一顺序加盐会把有序性破坏掉盐后缀会混进排序键里。我一般会在加盐前先问自己一句这个 job 的分组语义是什么如果分组的边界都不能变就不加盐改从输入源头做预聚合。4. 高频算法源码拆解二次排序、Join 与去重的骨架与改造路径4.1 一个算法工程的目录结构哪些是骨架哪些是胶水拿到一份结构良好的大数据算法源码包先看目录而不是先看代码。常见的包结构是这样src/main/java/com/example/algo/ ├── topn/ Top-N 算法 │ ├── TopNMapper.java │ ├── TopNReducer.java │ └── TopNJob.java ├── secondarysort/ 二次排序 │ ├── CompositeKey.java │ ├── CategoryPartitioner.java │ ├── CategoryGroupingComparator.java │ └── SecondarySortJob.java ├── join/ 关联算法 │ ├── ReduceSideJoinMapper.java │ ├── ReduceSideJoinReducer.java │ └── MapSideJoinDriver.java └── dedup/ 去重算法 ├── DedupMapper.java └── DedupReducer.java判断一个类是不是骨架看它是否直接参与 shuffle 语义的决定。Mapper、Reducer、Partitioner、Comparator 这四类是骨架改一个字符都可能改变结果Job 驱动类里设置参数的部分是胶水换路径、调并行度都动这里。实战里我建议先从 secondarysort 目录读起因为它同时涉及四个骨架类读完它其他算法的源码基本就能顺着同样的思路扫了。4.2 二次排序复合键、分区器、分组比较器的三角关系二次排序的典型诉求是「按分类分组组内按分数排序」。MapReduce 默认的排序是在整个 key 上做的要实现分组和组内排序必须自己把「分组的维度」和「排序的维度」装进同一个复合键里再配合自定义分区器和分组比较器。三个类缺一个都会出错。import org.apache.hadoop.io.WritableComparable; import org.apache.hadoop.io.WritableUtils; import java.io.DataInput; import java.io.DataOutput; import java.io.IOException; // 复合键category 负责分组score 负责组内排序 public class CompositeKey implements WritableComparableCompositeKey { private String category; private int score; public CompositeKey() {} public CompositeKey(String category, int score) { this.category category; this.score score; } Override public int compareTo(CompositeKey o) { int cmp this.category.compareTo(o.category); if (cmp ! 0) { return cmp; // 先按分类排 } return Integer.compare(this.score, o.score); // 组内按分数升序 } // 序列化与反序列化顺序必须和字段声明一致 Override public void write(DataOutput out) throws IOException { WritableUtils.writeString(out, category); out.writeInt(score); } Override public void readFields(DataInput in) throws IOException { this.category WritableUtils.readString(in); this.score in.readInt(); } public String getCategory() { return category; } public int getScore() { return score; } }复合键写完之后还要加一个按 category 分区、一个按 category 分组。分区器决定「哪些 key 进哪个 reducer」分组比较器决定「进同一个 reducer 后哪些 key 被分到同一组调用一次 reduce」。import org.apache.hadoop.mapreduce.Partitioner; // 分区逻辑只按 category 哈希保证同一分类的数据落在同一个 reducer public class CategoryPartitioner extends PartitionerCompositeKey, Object { Override public int getPartition(CompositeKey key, Object value, int numPartitions) { return (key.getCategory().hashCode() Integer.MAX_VALUE) % numPartitions; } }import org.apache.hadoop.io.WritableComparable; import org.apache.hadoop.io.WritableComparator; // 分组逻辑只比较 category忽略 score public class CategoryGroupingComparator extends WritableComparator { protected CategoryGroupingComparator() { super(CompositeKey.class, true); } Override public int compare(WritableComparable a, WritableComparable b) { CompositeKey ka (CompositeKey) a; CompositeKey kb (CompositeKey) b; return ka.getCategory().compareTo(kb.getCategory()); } }这套三角关系最常见的坑是只写了复合键忘了在 Job 里注册分区器和分组比较器。不注册的后果是分区器用默认的完整 key 哈希同一个 category 的 key 被拆到不同 reducer分组比较器用完整 key 比较每组只会分到一条数据reduce 里看到的「组」就碎了。在驱动类里必须显式调用job.setPartitionerClass和job.setGroupingComparatorClass少一行都是错。4.3 Reduce 端 Join 与 Map 端 Join源码包里的两种现成写法Join 是数据算法里最容易出现倾斜的环节。Reduce 端 Join 的写法是给数据打标签让两种来源在 reducer 里碰头。Map 端 Join 则是把小表塞进缓存在 map 阶段直接查字典跳过 shuffle。// Reduce 端 Join 的 mapper 骨架打上数据源标签 public class ReduceSideJoinMapper extends MapperLongWritable, Text, Text, Text { private String sourceTag; // 来自 job 配置例如 A 或 B Override protected void setup(Context context) { sourceTag context.getConfiguration().get(join.source, A); } Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String[] fields value.toString().split(,); if (fields.length 2) { return; } String joinKey fields[0]; // 关联键 String payload fields[1]; // 业务字段 // 值里带标签reduce 阶段才能区分来源 context.write(new Text(joinKey), new Text(sourceTag : payload)); } }这个写法的代价是两种来源的 key 都要发生 shuffle。如果一张表很小百万行以下更常见做法是把它做成分布式缓存或 Spark broadcast 变量map 端直接查内存字典彻底免掉 Join 的 shuffle。选哪个方案判据只有一个小表能不能装进单机内存。能装就选 Map 端 Join不能装就老实走 Reduce 端 Join或者考虑用大表按 key 预先分桶。去重算法相对简单mapper 把去重字段拼成 keycombiner 阶段先做一次组内去重reducer 再透传一次效果是两层压缩。源码里如果看见一个只有 map 没有 reduce 的 dedup job那是用了 Hive 的 DISTINCT 语义或者把去重前置到了输入清洗层本质上都是同一招在数据扩散出去之前先把重复项灭掉。4.4 把骨架改造成内部库接口抽象与本地回归从源码包抄到业务里最怕的是直接把 Mapper 里的解析逻辑揉进算法逻辑。我习惯做一层极薄的接口抽象把「记录的解析」和「算法的骨架」分开// 调用方只需实现解析器算法流程不感知具体业务字段 public interface RecordParserT { Comparable extractSortKey(String row); T parse(String row); }Top-N 的 Mapper 持有这个接口解析、比较都走它新增一个业务场景时只写一个新的 Parser 类算法主流程一行不动。这么做的好处是回归成本低本地起一个 JUnit 测试喂几十行 fixture 数据断言 Top-N 结果顺序跑完就知道接口改动有没有破坏算法语义。很多团队把源码包直接复制到业务代码里改字段名改到第三处就开始失控根源就是缺少这一层隔离。5. 避坑实录数据算法源码常见的 6 个翻车现场5.1 现象本地跑得好好的上集群一跑就内存溢出本地模式数据量小TreeMap 和堆都毫无压力上了集群每个 mapper 处理的输入分片可能上 GB局部堆里的Text对象数量不变但每个分片的字符串长度和总量都爆了。翻车点往往在new Text(value)频繁创建对象GC 压力飙升。原因本地没有分布式内存限制集群的 map 任务的堆内存是固定值。解决在 Job 里显式设mapreduce.map.memory.mb和mapreduce.reduce.memory.mb同时在 mapper 里控制单条记录的解析缓冲比如只保留参与排序的字段而不是整行原样存进堆。把整行丢进 TreeMap 是最顺手但也最容易爆内存的写法。5.2 现象Top-N 结果多出 N-1 条而且每次跑还不一样这个现象出现时先查是不是 TreeMap 的 key 设置有问题。如果 key 里只放 score分数相同的记录会发生覆盖最终输出不足 N 条如果反过来用整行做 key两条内容不同的记录永远不会视为重复TreeMap 容量失效最终输出超过 N 条。原因没有给 TreeMap 设定「什么算同一条记录」的规则。解决key 用「score 记录唯一 id」的复合结构或者允许分数的重复并让 TreeMap 里存一个 List。每次都跑都是这个结果说明逻辑稳定但结果本身就错了这种问题在测试阶段很难发现因为 Top-N 的边界刚好是那个容易被忽略的相等区间。5.3 现象二次排序的组内顺序完全不对分组是对的但组内顺序乱最常见原因是只写了compareTo没在 Job 里注册setSortComparatorClass。复合键的compareTo默认会被用作排序器但如果代码里重写了排序器类或者用了老版本的 Hadoop API排序可能退化为主键排序。原因排序、分组、分区三个环节各自都有比较器任何一个没有和复合键保持一致结果就会「组对了、序乱了」。解决在驱动类里把三个类全部显式注册用「小数据集 断言顺序」的测试钉死这个行为不要靠肉眼检查输出。5.4 现象Spark 作业用了 cache() 还越跑越慢cache 不是万能的。如果中间集合太大比如超过执行器内存的一半Spark 会把它从内存溢出到磁盘每次重算重新读盘比不 cache 还慢。原因cache 策略默认是MEMORY_ONLY换MEMORY_AND_DISK或者先看看 storage 页的缓存命中率。解决先用rdd.count()估算集合规模再决定用 cache 还是 checkpoint。对超长血缘的链我一般直接上 checkpoint斩断血缘比缓存更治本。5.5 现象改了一个 top.n 配置结果反而变成了取最大数配置从 10 改成 20结果范围扩大表面上正常但某次改成负数之后TreeMap 的 remove 逻辑直接出错输出为空或者全量。原因源码里没有对配置做边界校验负数容量直接进入堆裁剪逻辑。解决在 setup 方法里加一段防御式校验n 小于等于 0 时抛出IllegalArgumentException或强制回退到默认值。这类问题在脚本式的大数据作业里很常见一个配置项就可能让整个算法进入黑匣子状态前 15 分钟排查全耗在核对「到底有没有生效」上。5.6 现象Reducer 收的数据分布极不均匀任务 99% 但卡死默认哈希分区器在少数热点 key 面前毫无还手之力几十个 reducer 里一个扛了 90% 的数据。原因数据本身的 key 分布倾斜算法代码又没做任何预处理。解决在 map 端先做一轮预聚合比如先按 key 的加盐版本做局部计数再走正式 reducer或者按 4.3 节说的把参与 Join 的小表广播出去直接绕开 shuffle。判断是不是倾斜去集群的 job 页面看每个 reducer 的输入字节数如果最大和最小差了一个数量级以上基本就是它。6. 从源码到上线本地验证三板斧与结果核对习惯源码能在本地跑通只算完成了一半。我自己的验证习惯是三板斧本地小样本断言、中间结果落盘对比、独立实现交叉核对。第一板斧是本地 fixture。每个算法骨架配套一个几十行的测试数据明确写出期望结果。Top-N 就断言输出顺序二次排序就断言每个分组内部单调。这一步能把 5.2 和 5.3 这类边界错误在最早期拦下来。第二板斧是中间结果落盘。MapReduce 作业跑完把 reducer 的输入单独导一份出来对比最终输出确认算法语义没被框架的排序规则偷换。第三板斧是独立交叉核对也是我最想强调的。拿 Top-N 来说在本地用 SQLORDER BY score DESC LIMIT 10或者 Python 里直接sorted()跑一遍同样的输入对比两边结果。这里有个技巧不要只比对 Top-N 的名单要连顺序一起比因为某些业务场景里相同分数的次序也有语义。我踩过一次很深的坑两套实现输出的 Top-10 内容完全一致但相同分数记录的先后顺序不同下游消费方把顺序当成了业务含义结果两张报表对不上查了一整天才定位到。养成这个习惯之后我接手任何数据算法源码包的第一件事永远是先写交叉验证脚本再看源码。代码里的逻辑再自信也不如用一套独立实现把结果钉死来得踏实。这个顺序反过来很容易被源码带着走把错误当预期。希望帮到你。本文还有配套的精品资源点击获取