
简介这份资源面向高校大数据、人工智能相关专业的课程设计学习者提供一套基于Hadoop与Spark的中文手写数字实时识别系统完整实现适合作为期末大作业或课程设计参考新手也能借助注释快速理解。压缩包共8个文件约9.06MB包含6个Python脚本、1份PDF实验方案和1段mp4演示视频脚本覆盖HOG特征提取、RDD与DataFrame两种逻辑回归实现、t-SNE可视化及基于Sklearn的模型对比PDF则给出实验方案与文档说明视频用于直观展示系统运行效果。目前已有500人学习下载。读者可从中获得从特征工程、分布式训练到实时识别展示的完整链路理解Hadoop与Spark在图像识别任务中的分工并参考实验报告结构完成自己的课程设计文档具备较高的实际应用与借鉴价值。1. 从一份课程设计说起Hadoop 和 Spark 怎么把中文手写数字识别跑成实时系统很多人第一次看到「基于 Hadoop 和 Spark 的中文手写数字实时识别系统」这个题目第一反应是手写数字识别不是 MNIST 加个 CNN 就完事了吗为什么非要套上 Hadoop 和 Spark我当初也这么想直到自己动手把单机脚本改成集群任务才发现坑全在「实时」和「中文」这两个词上。中文手写数字的书写习惯和 MNIST 里的西文样本差别很大笔画粗细、连笔、倾斜角度都更野单机模型准确率能到 98%一上真实采集的样本就掉到 90% 以下。而「实时」意味着你不能等一批数据攒够了再离线跑得让图片从上传到返回结果控制在秒级。这套课程设计的价值就在于用 Hadoop 存样本和模型、用 Spark 做分布式推理和增量训练把「采集—存储—识别—反馈」串成一条能跑起来的链路。它适合正在做大数据课程设计的学生也适合想搞清楚 Spark 流式推理怎么落地的一线工程师。下面我按自己复现过的路径把选型、搭建、代码和踩坑一次讲透。2. 中文手写数字识别的数据链路HDFS 存什么、Spark 算什么2.1 为什么样本和模型都要进 HDFS单机做手写数字识别图片放本地目录、模型存成.h5或.pt文件就够了。但一旦要做「实时」和「多人并发」本地文件系统立刻成为瓶颈多个推理节点要读同一份模型采集端还在不断写入新样本文件锁和路径冲突会让你怀疑人生。HDFS 的好处是它天生为「一次写入、多次读取」设计模型文件上传一次所有 Spark Executor 都能通过统一 URI 拉取新采集的样本按日期分区写入既不干扰正在读的旧数据也方便后续做增量训练。我一般会这样规划 HDFS 目录# 在 HDFS 上建立项目根目录 hdfs dfs -mkdir -p /handwriting/raw/2025-01-01 # 原始采集图片按天分区 hdfs dfs -mkdir -p /handwriting/processed/train # 预处理后的训练集 hdfs dfs -mkdir -p /handwriting/model # 训练好的模型文件 hdfs dfs -mkdir -p /handwriting/tmp # Spark 临时输出 # 上传一张本地图片做测试 hdfs dfs -put ./sample_0.png /handwriting/raw/2025-01-01/ hdfs dfs -ls /handwriting/raw/2025-01-01/这段命令的逻辑很直白raw放原始数据processed放归一化后的数据model放模型tmp给 Spark 写中间结果。参数上唯一要注意的是分区粒度——按天分区在课程设计的数据量下足够如果采集频率高可以改成按小时。hdfs dfs -ls用来确认文件真的写进去了很多人第一次跑 Spark 读不到数据就是上传路径写错或者权限不对。提示HDFS 默认块大小 128MB手写数字图片单张只有几 KB会产生大量小文件。课程设计阶段数据量不大可以忍但如果要模拟真实场景建议先用 Spark 做一次合并把同一天的图片打包成 SequenceFile 或 Parquet。2.2 Spark 在这条链路里到底承担什么角色Spark 在这套系统里干两件事离线批量训练和近实时流式推理。训练部分用Spark MLlib或者Spark TensorFlow/PyTorch的分布式封装把预处理后的样本分片到多个 Executor 上并行计算梯度推理部分用Spark Structured Streaming监听一个消息队列或文件目录一旦有新图片写入就触发识别把结果写回 HDFS 或数据库。为什么不用 Hadoop MapReduce 做推理因为 MapReduce 的启动开销太大一个作业从提交到出结果动辄几十秒根本谈不上「实时」。Spark 的 DAG 调度和内存计算能把单次推理压到毫秒级配合 Structured Streaming 的微批模式端到端延迟可以控制在 12 秒。这是选型的核心理由也是课程设计里最容易被忽略的对比点。下面是一个用 Structured Streaming 读取新图片并调用模型的骨架代码from pyspark.sql import SparkSession from pyspark.sql.functions import udf, col from pyspark.sql.types import StringType import numpy as np from PIL import Image import io spark SparkSession.builder \ .appName(HandwritingRealtime) \ .config(spark.executor.memory, 2g) \ .config(spark.sql.streaming.checkpointLocation, /handwriting/tmp/checkpoint) \ .getOrCreate() # 加载训练好的模型到广播变量避免每个 task 重复加载 model_bytes open(/handwriting/model/cnn_model.pkl, rb).read() broadcast_model spark.sparkContext.broadcast(model_bytes) def predict(image_bytes): import pickle model pickle.loads(broadcast_model.value) img Image.open(io.BytesIO(image_bytes)).convert(L).resize((28, 28)) arr np.array(img).reshape(1, 28, 28, 1) / 255.0 pred model.predict(arr) return str(int(np.argmax(pred))) predict_udf udf(predict, StringType()) # 监听 HDFS 目录下的新文件 stream_df spark.readStream \ .format(binaryFile) \ .option(path, /handwriting/raw/2025-01-01/) \ .load() result stream_df.withColumn(digit, predict_udf(col(content))) query result.select(path, digit) \ .writeStream \ .outputMode(append) \ .format(parquet) \ .option(path, /handwriting/processed/result) \ .option(checkpointLocation, /handwriting/tmp/checkpoint) \ .start() query.awaitTermination()逻辑说明binaryFile格式让 Spark 把每个文件当成一行二进制数据读入content列就是图片字节流。UDF 里先反序列化模型再预处理图片、跑推理、返回数字字符串。参数上executor.memory给 2g 是课程设计规模下的保守值如果模型是 ResNet 级别要往上加checkpointLocation必须设置否则流式查询重启后会重复处理数据。广播变量是这里的关键优化——如果不广播每个 task 都会重新加载一次模型内存直接爆掉。2.3 从图片到特征预处理为什么不能省中文手写数字的原始图片往往带背景噪点、边框、不同分辨率。直接丢给模型准确率会掉得很惨。我一般会在 Spark 里做三步预处理灰度化、二值化、居中裁剪。灰度化去掉颜色干扰二值化把笔画和背景分开居中裁剪保证数字落在 28×28 画布的中心——这一步对 MNIST 系模型尤其重要因为训练数据本身就是居中的。def preprocess(img_bytes): img Image.open(io.BytesIO(img_bytes)).convert(L) # 二值化阈值 128低于阈值的像素变黑 img img.point(lambda x: 0 if x 128 else 255) # 找到非零区域的边界框 bbox img.getbbox() if bbox: img img.crop(bbox) # 缩放到 20x20 后粘贴到 28x28 中心 img img.resize((20, 20)) canvas Image.new(L, (28, 28), 255) canvas.paste(img, ((28 - 20) // 2, (28 - 20) // 2)) return np.array(canvas) / 255.0参数说明阈值 128 是经验值如果采集的图片偏暗可以调到 100偏亮调到 150。resize((20, 20))留出 4 像素边距是为了模仿 MNIST 的预处理方式。这一步做完模型在真实样本上的准确率通常能提升 35 个百分点。别小看这几步很多人模型训练 loss 降不下去问题就出在预处理和训练数据不一致。3. 集群搭建与 Spark 提交从伪分布式到能跑通的最小配置3.1 Hadoop 伪分布式搭建的四个关键配置课程设计阶段不需要真集群伪分布式足够跑通全流程。核心是改四个文件core-site.xml、hdfs-site.xml、mapred-site.xml、yarn-site.xml。我见过太多人卡在datanode起不来最后发现是hdfs-site.xml里dfs.replication设成了 3伪分布式只有一个节点副本数必须改成 1。!-- 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 property namedfs.namenode.name.dir/name value/usr/local/hadoop/data/name/value /property property namedfs.datanode.data.dir/name value/usr/local/hadoop/data/data/value /property /configuration配置完执行格式化hdfs namenode -format然后start-dfs.sh。用jps检查应该看到NameNode、DataNode、SecondaryNameNode三个进程。如果DataNode没起来九成是data目录权限不对或者之前格式化残留了旧数据删掉data目录重新格式化即可。注意格式化只能执行一次重复格式化会导致clusterID不一致DataNode直接拒绝启动。这是血泪经验我当初重装了三次才记住。3.2 Spark 提交作业时的内存与并行度参数Spark 提交作业时--executor-memory、--num-executors、--executor-cores这三个参数决定了作业能不能跑完。课程设计的机器通常内存有限我一般这样配spark-submit \ --master yarn \ --deploy-mode client \ --num-executors 2 \ --executor-cores 2 \ --executor-memory 2g \ --driver-memory 1g \ --conf spark.sql.shuffle.partitions8 \ handwriting_stream.py参数含义num-executors 2表示两个执行器executor-cores 2表示每个执行器两个核总并行度是 4。shuffle.partitions默认是 200小数据量下会产生大量空 task拖慢速度改成 8 和并行度匹配。如果作业报Container killed by YARN for exceeding memory limits优先加executor-memory其次检查是否有数据倾斜。3.3 用 YARN 还是 Standalone课程设计的取舍YARN 的好处是资源调度成熟能和 Hadoop 共用集群坏处是配置复杂yarn-site.xml里yarn.nodemanager.resource.memory-mb设小了容器起不来设大了宿主机卡死。Standalone 模式配置简单start-all.sh一把梭但资源隔离差。我的建议是如果课程设计只要求跑通用 Standalone 省时间如果要写实验报告体现「大数据生态整合」用 YARN 更有说服力。两者在代码层面没有区别只是spark-submit的--master参数不同。4. 避坑与排查实时识别系统最容易翻车的五个地方4.1 流式查询重复消费导致结果翻倍现象Structured Streaming 作业重启后输出目录里的结果数量是实际图片的两倍。原因没有设置checkpointLocation或者 checkpoint 目录被手动删除Spark 无法记录消费偏移量重启后从头消费。解决始终设置独立的 checkpoint 目录且不要和输出目录混用。如果已经产生重复数据用dropDuplicates按文件路径去重。4.2 模型广播后 Executor 内存溢出现象作业跑几分钟后报java.lang.OutOfMemoryErrorExecutor 被 YARN 杀掉。原因模型文件太大广播变量在每个 Executor 上复制一份加上推理时的中间张量内存不够。解决换更小的模型或者用spark.executor.memoryOverhead增加堆外内存。如果模型超过 200MB考虑用torch.distributed或 TensorFlow Serving 做独立推理服务Spark 只负责调度。4.3 中文手写样本的编码问题现象读取图片时抛UnicodeDecodeError或者文件名乱码。原因采集端保存的文件名包含中文字符HDFS 和 Spark 默认按 UTF-8 处理但某些 Windows 采集工具用 GBK 编码。解决采集端统一用 UUID 命名文件中文标签单独存一张映射表或者在 Spark 里用spark.hadoop.fs.defaultFS配合-Dfile.encodingUTF-8启动。4.4 小文件过多拖垮 NameNode现象训练作业启动极慢hdfs dfs -ls列出几万个文件。原因每张图片单独存一个文件HDFS 元数据压力大。解决用 Spark 的coalesce或repartition把小文件合并成 Parquet每个文件 128MB 左右。课程设计阶段可以在预处理时直接写 Parquet而不是存原始图片。4.5 实时延迟忽高忽低现象大部分请求 1 秒内返回偶尔飙到 10 秒以上。原因Spark 微批的触发间隔和 HDFS 的写入延迟叠加或者某个 Executor 在做 GC。解决把trigger设为processingTime1 second限制微批频率同时监控 Executor 的 GC 日志如果 Full GC 频繁减小executor-memory或改用 G1 垃圾回收器。5. 让识别率再上一个台阶增量训练与模型热更新的具体做法课程设计做完基础版之后如果想拿高分或者真正上线用增量训练是绕不开的。真实场景里用户的手写风格会漂移今天写的「7」和上个月写的「7」可能完全不一样。我的做法是每天把新采集的样本用 Spark 跑一遍预处理和旧样本按 1:3 的比例混合用Spark MLlib的LogisticRegression或MultilayerPerceptronClassifier做一次增量 fit然后把新模型写到 HDFS 的model目录同时保留旧版本。from pyspark.ml.classification import MultilayerPerceptronClassifier from pyspark.ml.linalg import Vectors from pyspark.sql import Row # 读取新旧样本合并后训练 old_df spark.read.parquet(/handwriting/processed/train) new_df spark.read.parquet(/handwriting/processed/new) # 按 3:1 采样避免新样本过拟合 old_sample old_df.sample(False, 0.75) combined old_sample.union(new_df) # 定义网络输入 784隐藏层 128输出 10 layers [784, 128, 10] trainer MultilayerPerceptronClassifier( layerslayers, maxIter100, blockSize128, seed42 ) model trainer.fit(combined) # 保存模型带日期后缀 model.write().overwrite().save(/handwriting/model/mlp_20250101)参数说明layers里 784 是 28×28 展平后的维度128 是隐藏层节点数10 是数字类别。maxIter设 100 在课程设计数据量下足够收敛如果 loss 还在降可以加到 200。blockSize影响矩阵分块128 是内存和速度的平衡点。保存模型时带日期后缀方便回滚——如果新模型在验证集上准确率下降直接切回旧版本。模型热更新用 Spark 的广播变量配合一个定时任务实现每训练完一个新模型写一个latest.txt到 HDFS流式作业每隔 5 分钟读一次这个文件如果版本号变了就重新广播模型。这样不用重启作业就能切换模型对「实时」场景很关键。验证方法上我习惯留 10% 的样本做 hold-out每次增量训练后跑一遍混淆矩阵重点看「1」和「7」、「3」和「8」这些容易混的类别。如果某一类召回率掉得厉害说明新样本里这类数据分布有问题需要人工检查采集端。最后说个我自己的习惯每次改完代码先在本地用 100 张图片跑一遍完整链路确认预处理、推理、输出都正常再提交到集群。集群上调试的成本太高一个参数写错可能要等十分钟才能看到报错。这套系统我前后搭了四遍前三次都栽在「本地能跑、集群报错」上后来养成先本地验证的习惯省了至少一半时间。希望帮到你。本文还有配套的精品资源点击获取