Hadoop+Spark实战:中文手写数字实时识别系统课程设计
简介这份资源面向高校大数据、人工智能相关专业的课程设计学习者提供一套基于Hadoop与Spark的中文手写数字实时识别系统完整实现适合作为期末大作业或课程设计参考新手也能借助注释快速理解。压缩包共8个文件约9.06MB包含6个Python脚本、1份PDF实验方案和1段mp4演示视频脚本覆盖HOG特征提取、RDD与DataFrame两种逻辑回归实现、t-SNE可视化及基于Sklearn的模型对比PDF则给出实验方案与文档说明视频用于展示系统运行效果。目前已有500人学习下载。读者可据此掌握从特征工程、分布式训练到实时识别展示的完整流程理解Spark与Hadoop在图像识别任务中的协作方式并直接部署运行、对照实验报告完成自己的课程设计省去从零搭建的时间。1. 从一次课程设计翻车说起HadoopSpark 做中文手写数字实时识别到底难在哪每年到了课程设计季总有一批人栽在“实时”两个字上。我见过太多同学把 MNIST 训练脚本跑通、准确率刷到 99%结果答辩时老师随手在画板上写了个歪歪扭扭的“7”系统识别成“1”场面一度非常尴尬。问题不在模型本身而在于整个链路——从手写笔迹采集、图像预处理、特征提取到 Hadoop 分布式存储、Spark 流式推理再到前端实时回传结果任何一环掉链子都会让“实时识别”变成“实时翻车”。这个标题讲的就是这么一套系统用 Hadoop 的 HDFS 存训练数据和模型文件用 Spark 做分布式训练和流式推理前端采集中文手写数字注意是中文语境下的手写习惯不是标准 MNIST 那种规整数字后端实时返回识别结果。它适合正在做大数据课程设计、想找一个能同时体现 Hadoop 和 Spark 两个技术栈、又不想做烂大街的“词频统计”或“日志分析”的同学。核心难点有三个第一中文手写数字的笔画习惯和 MNIST 差异很大直接拿 MNIST 预训练模型迁移效果会打折扣第二Spark Streaming 的微批处理延迟和“实时”之间的平衡第三Hadoop 和 Spark 的版本兼容性这个坑我后面会专门讲。2. 系统架构拆解HDFS 存什么、Spark 算什么、前端怎么接2.1 为什么选 Hadoop Spark 而不是单机 Flask单机 Flask 加一个 CNN 模型也能做手写数字识别但课程设计要体现“大数据”三个字就得让数据流经过分布式组件。Hadoop 在这里的角色是分布式文件系统 HDFS负责存原始手写图片、预处理后的向量数据、训练好的模型文件。Spark 的角色有两个Spark SQL/DataFrame 做离线批量训练Spark Streaming 做在线流式推理。选这个组合的理由很实际HDFS 的副本机制保证了数据不丢Spark 的内存计算让训练速度比 MapReduce 快一个数量级。而且 Spark 的 MLlib 里自带多层感知机分类器不需要额外引入 TensorFlow 或 PyTorch 就能跑通一个还不错的基线模型。如果你的课程设计允许引入深度学习框架也可以在 Spark 里用 Pandas UDF 调用 Keras 模型但那是进阶玩法后面第 6 章会讲。2.2 数据采集与预处理链路中文手写数字的数据来源一般有两种一是自己用 HTML5 Canvas 写一个采集页面让同学帮忙写几百张二是找公开的中文手写数字数据集比如 CASIA-HWDB 里的数字部分。不管哪种来源预处理步骤是一样的# preprocess.py # 将原始手写图片统一为 28x28 灰度图并转为 LibSVM 格式供 Spark MLlib 使用 import cv2 import numpy as np import os def preprocess_image(img_path, size(28, 28)): # 以灰度模式读取 img cv2.imread(img_path, cv2.IMREAD_GRAYSCALE) if img is None: return None # 自适应二值化应对不同光照和笔迹深浅 img cv2.adaptiveThreshold(img, 255, cv2.ADAPTIVE_THRESH_GAUSSIAN_C, cv2.THRESH_BINARY_INV, 11, 2) # 找到数字的边界框并裁剪去掉多余白边 coords cv2.findNonZero(img) if coords is None: return None x, y, w, h cv2.boundingRect(coords) img img[y:yh, x:xw] # 保持长宽比缩放到 20x20再填充到 28x28 if h w: new_h, new_w 20, int(20 * w / h) else: new_h, new_w int(20 * h / w), 20 img cv2.resize(img, (new_w, new_h), interpolationcv2.INTER_AREA) canvas np.zeros(size, dtypenp.uint8) x_offset (size[1] - new_w) // 2 y_offset (size[0] - new_h) // 2 canvas[y_offset:y_offsetnew_h, x_offset:x_offsetnew_w] img return canvas def to_libsvm(img, label): # 将 28x28 图像展平为 784 维向量输出 LibSVM 格式字符串 pixels img.flatten() features .join([f{i1}:{p/255.0} for i, p in enumerate(pixels) if p 0]) return f{label} {features}这段代码的关键参数有三个adaptiveThreshold的blockSize11和C2决定了二值化效果手写笔迹细的时候要调小 blockSizeresize到 20x20 再填充到 28x28 是 MNIST 的标准做法能保留数字的形态特征to_libsvm里只保留非零像素能大幅压缩数据量Spark 读起来更快。预处理完的数据通过hdfs dfs -put上传到 HDFS 的/user/hadoop/digits/raw/目录。这里有个细节LibSVM 格式的文件不要每个样本一个文件而是合并成几个大文件否则 HDFS 的小文件问题会让 Spark 读取时启动几百个 task速度反而变慢。2.3 Spark 离线训练与模型持久化训练部分用 Spark MLlib 的MultilayerPerceptronClassifier虽然它不如 Keras 灵活但胜在原生支持分布式训练不需要额外配置 Python 环境。# train.py from pyspark.sql import SparkSession from pyspark.ml.classification import MultilayerPerceptronClassifier from pyspark.ml.evaluation import MulticlassClassificationEvaluator from pyspark.ml.linalg import Vectors from pyspark.ml.feature import StringIndexer spark SparkSession.builder \ .appName(HandwrittenDigitTraining) \ .master(yarn) \ .config(spark.executor.memory, 4g) \ .config(spark.executor.cores, 2) \ .getOrCreate() # 读取 LibSVM 格式数据Spark 原生支持 data spark.read.format(libsvm).load(hdfs:///user/hadoop/digits/raw/train.libsvm) # 划分训练集和测试集 train_data, test_data data.randomSplit([0.8, 0.2], seed42) # 定义网络结构输入 784 维两个隐藏层分别 256 和 128 个神经元输出 10 类 layers [784, 256, 128, 10] trainer MultilayerPerceptronClassifier( layerslayers, maxIter100, blockSize128, seed42, featuresColfeatures, labelCollabel ) model trainer.fit(train_data) # 在测试集上评估 predictions model.transform(test_data) evaluator MulticlassClassificationEvaluator( labelCollabel, predictionColprediction, metricNameaccuracy) accuracy evaluator.evaluate(predictions) print(fTest Accuracy: {accuracy:.4f}) # 保存模型到 HDFS model.write().overwrite().save(hdfs:///user/hadoop/digits/model/mlp_model)layers参数是网络结构[784, 256, 128, 10]是我试过在课程设计规模下性价比最高的配置再大容易过拟合再小准确率上不去。maxIter100配合blockSize128在 4 个 executor 的集群上大概跑 8 到 12 分钟。blockSize调大能加快训练但吃内存如果 executor 内存只有 2G建议降到 64。模型保存到 HDFS 后Spark Streaming 任务可以直接加载不需要每次重新训练。3. 实时识别链路Spark Streaming 微批处理与前端联调3.1 Spark Streaming 接收手写图片流实时部分的核心是 Spark Streaming 的Receiver或Direct Kafka模式。课程设计里最省事的做法是用 Socket 流前端 Canvas 把图片 Base64 编码后通过 WebSocket 发给一个中转服务中转服务再把图片路径或字节流推给 Spark Streaming 的 Socket 端口。# streaming_inference.py from pyspark import SparkContext from pyspark.streaming import StreamingContext from pyspark.ml.classification import MultilayerPerceptronClassificationModel from pyspark.ml.linalg import Vectors import base64 import cv2 import numpy as np sc SparkContext(appNameDigitStreamingInference) ssc StreamingContext(sc, 2) # 2 秒一个微批 # 加载离线训练好的模型 model MultilayerPerceptronClassificationModel.load( hdfs:///user/hadoop/digits/model/mlp_model) def predict(rdd): if rdd.isEmpty(): return results [] for img_base64 in rdd.collect(): # 解码 Base64 图片 img_bytes base64.b64decode(img_base64) img_array np.frombuffer(img_bytes, np.uint8) img cv2.imdecode(img_array, cv2.IMREAD_GRAYSCALE) # 复用预处理逻辑 processed preprocess_image_from_array(img) if processed is None: continue features Vectors.dense(processed.flatten() / 255.0) prediction model.predict(features) results.append(int(prediction)) # 将结果写回 Kafka 或 Socket供前端轮询 if results: send_results_to_frontend(results) stream ssc.socketTextStream(localhost, 9999) stream.foreachRDD(predict) ssc.start() ssc.awaitTermination()ssc StreamingContext(sc, 2)里的 2 秒是批处理间隔决定了“实时”的上限。调成 1 秒延迟更低但每个批次的数据量可能不够模型推理的吞吐优势发挥不出来。我一般建议课程设计用 2 秒答辩演示时肉眼几乎感觉不到延迟。foreachRDD里对每个 RDD 做collect()是把数据拉到 Driver 端处理这在生产环境是禁忌但课程设计数据量小这样写最简单。如果想让代码更“分布式”可以用mapPartitions在每个 executor 上加载模型副本但模型文件需要提前分发到各节点。3.2 前端 Canvas 采集与结果回显前端部分用 HTML5 Canvas 采集手写笔迹转成 Base64 后通过 WebSocket 发送。这里的关键是笔迹的粗细和颜色要固定否则预处理阶段的自适应二值化会不稳定。// frontend.js const canvas document.getElementById(digitCanvas); const ctx canvas.getContext(2d); let isDrawing false; // 初始化画布为黑色背景、白色笔迹和 MNIST 风格一致 ctx.fillStyle #000; ctx.fillRect(0, 0, canvas.width, canvas.height); ctx.strokeStyle #fff; ctx.lineWidth 12; ctx.lineCap round; canvas.addEventListener(mousedown, (e) { isDrawing true; ctx.beginPath(); ctx.moveTo(e.offsetX, e.offsetY); }); canvas.addEventListener(mousemove, (e) { if (!isDrawing) return; ctx.lineTo(e.offsetX, e.offsetY); ctx.stroke(); }); canvas.addEventListener(mouseup, () { isDrawing false; // 延迟 500ms 后发送避免笔画未完成就识别 setTimeout(sendImage, 500); }); function sendImage() { const dataURL canvas.toDataURL(image/png); const base64 dataURL.split(,)[1]; // 通过 WebSocket 发送到中转服务 ws.send(JSON.stringify({ type: digit, data: base64 })); } // 接收识别结果 ws.onmessage (event) { const result JSON.parse(event.data); document.getElementById(result).innerText 识别结果${result.digit}; };lineWidth 12是经过多次试验的值太细会导致二值化后笔画断裂太粗会让数字粘连。setTimeout(sendImage, 500)是给用户一个“写完抬手”的缓冲避免写到一半就被识别。3.3 端到端联调时的三个关键参数联调阶段最容易出问题的是三个参数Spark Streaming 的批处理间隔、WebSocket 的发送频率、预处理时的图像尺寸。批处理间隔 2 秒、WebSocket 每 500ms 发送一次、图像统一到 28x28这三个值配合起来能保证识别结果在 1 到 3 秒内返回。如果批处理间隔调到 5 秒用户会明显感觉卡顿如果 WebSocket 发送频率太高Spark 端会积压大量小批次反而拖慢整体速度。4. 避坑与排查版本兼容、小文件、序列化与内存溢出4.1 Hadoop 和 Spark 版本不匹配导致 NoSuchMethodError现象Spark 任务提交到 YARN 后立刻报java.lang.NoSuchMethodError堆栈指向 Hadoop 的Configuration类或FileSystem类。原因Spark 编译时依赖的 Hadoop 版本和集群实际安装的 Hadoop 版本不一致。比如 Spark 3.2 默认编译依赖 Hadoop 3.3.1但集群装的是 Hadoop 2.7.7某些 API 签名变了。解决去 Spark 官网下载页面找 “Pre-built for Apache Hadoop 2.7” 或对应版本的包不要用默认的 “Pre-built for Apache Hadoop 3.3”。如果已经装了不匹配的版本重新下载对应版本解压替换SPARK_HOME即可。这个坑我踩过两次血泪经验是装 Spark 之前先hadoop version确认集群版本。4.2 HDFS 小文件过多导致 Spark 任务启动缓慢现象训练任务提交后Spark UI 上显示有几千个 task但每个 task 处理的数据量只有几 KB整体运行时间比预期长好几倍。原因预处理阶段每个手写图片生成一个 LibSVM 文件HDFS 上堆积了大量小文件。Spark 读取时每个小文件至少对应一个 tasktask 启动和调度的开销远大于计算本身。解决在预处理脚本里把多个样本合并成一个大文件每个文件 128MB 左右。可以用hdfs dfs -getmerge把已有小文件合并或者在生成阶段就按批次写入。合并后 task 数量降到几十个训练速度提升明显。4.3 Spark MLlib 模型保存后加载报序列化错误现象离线训练时model.save()成功但 Streaming 任务里MultilayerPerceptronClassificationModel.load()报InvalidClassException或ClassNotFoundException。原因训练和推理用的 Spark 版本不一致或者模型保存时用了overwrite()但 HDFS 上残留了旧版本的元数据文件。解决确保训练和推理在同一个 Spark 版本下进行。如果 HDFS 上有旧模型先hdfs dfs -rm -r删干净再保存。另外MultilayerPerceptronClassificationModel的保存路径下会有data/和metadata/两个目录加载时路径要指到父目录不是data/子目录。4.4 Executor 内存溢出导致任务被 YARN 杀掉现象训练到第 30 轮左右任务突然失败YARN 日志显示Container killed by YARN for exceeding memory limits。原因spark.executor.memory设得太小或者blockSize设得太大导致每个 task 加载的数据块超出内存。MLP 训练时每轮的梯度更新需要缓存中间结果内存需求比预估的高。解决把spark.executor.memory从 2g 调到 4g同时把blockSize从 128 降到 64。如果集群资源有限可以减少maxIter或者用trainValidationSplit做早停。另外spark.memory.fraction可以适当调高到 0.7给计算留更多内存。4.5 前端 Base64 图片过大导致 Socket 缓冲区溢出现象前端发送图片后Spark Streaming 端偶尔收不到数据或者收到截断的 Base64 字符串解码时报binascii.Error。原因Canvas 默认导出 PNG 格式28x28 的图片 Base64 后只有几百字节但如果 Canvas 尺寸设成 400x400Base64 后可能超过 10KB。Socket 流的默认缓冲区可能不够或者 WebSocket 帧被分片。解决在前端canvas.toDataURL(image/jpeg, 0.8)用 JPEG 压缩或者直接把 Canvas 尺寸设成 28x28 再放大显示。更稳妥的做法是在 Base64 字符串前后加分隔符Spark 端按分隔符切分避免粘包。5. 从课程设计到能写进简历模型调优与工程化收尾5.1 用交叉验证和网格搜索找最优网络结构课程设计答辩时老师常问“为什么用这个网络结构”如果只回答“试出来的”会显得不够专业。可以用 Spark MLlib 的CrossValidator和ParamGridBuilder做一轮系统性的参数搜索。from pyspark.ml.tuning import CrossValidator, ParamGridBuilder # 定义参数网格 paramGrid ParamGridBuilder() \ .addGrid(trainer.layers, [[784, 128, 10], [784, 256, 10], [784, 256, 128, 10]]) \ .addGrid(trainer.maxIter, [50, 100]) \ .addGrid(trainer.blockSize, [64, 128]) \ .build() evaluator MulticlassClassificationEvaluator( labelCollabel, predictionColprediction, metricNameaccuracy) crossval CrossValidator(estimatortrainer, estimatorParamMapsparamGrid, evaluatorevaluator, numFolds3) cv_model crossval.fit(train_data) best_model cv_model.bestModel print(fBest layers: {best_model.getLayers()}) print(fBest accuracy: {evaluator.evaluate(best_model.transform(test_data))})numFolds3是课程设计规模下的合理选择5 折会让训练时间翻倍但准确率提升有限。paramGrid里我列了三组网络结构实际跑下来[784, 256, 128, 10]在中文手写数字上表现最稳因为两个隐藏层能更好地拟合笔画细节。5.2 模型持久化与增量训练课程设计通常只要求跑一次训练但如果你想在简历里写“支持增量训练”可以在 Streaming 任务里定期把新采集的样本追加到 HDFS然后每隔一段时间重新训练模型并热更新。# 将新采集的样本追加到 HDFS 训练集 hdfs dfs -appendToFile /local/new_samples.libsvm /user/hadoop/digits/raw/train.libsvm # 重新训练并覆盖旧模型 spark-submit --master yarn --deploy-mode cluster \ --executor-memory 4g --num-executors 4 \ train.py --input /user/hadoop/digits/raw/train.libsvm \ --output /user/hadoop/digits/model/mlp_model_v2热更新的关键是在 Streaming 任务里用model MultilayerPerceptronClassificationModel.load()重新加载而不是重启整个 Streaming 上下文。这样前端用户几乎无感知。5.3 一个容易被忽略的细节中文手写数字的“4”和“9”中文语境下手写数字有几个高频混淆对“4”和“9”在连笔时容易混“1”和“7”在带钩时容易混“0”和“6”在闭合不严时容易混。我在测试集上单独统计过这三对占了总错误率的 60% 以上。解决办法有两个一是在预处理阶段加一个“笔画端点检测”把端点数量作为额外特征拼接到 784 维向量后面变成 786 维输入二是在训练集里针对这三对数字做数据增强比如随机旋转 ±10 度、随机缩放 0.9 到 1.1 倍。我一般会先做数据增强因为改动最小效果也最直接。5.4 答辩演示时的稳定性技巧答辩现场的网络和集群状态不可控我一般会准备一个“降级方案”如果 YARN 资源不够或者 Spark 任务启动失败直接切换到本地模式master(local[*])跑一个轻量版模型。本地模式下把maxIter降到 30layers降到[784, 128, 10]准确率会掉 2 到 3 个百分点但至少能保证演示不翻车。另外前端 Canvas 的笔迹颜色和粗细在演示前一定要校准不同投影仪的色彩还原差异很大白色笔迹在有些屏幕上会偏灰导致二值化后笔画断裂。我的习惯是提前 30 分钟到答辩教室用实际投影仪跑一遍完整流程确认识别结果稳定后再开始讲。这套方案从 HDFS 存储到 Spark 训练再到 Streaming 推理覆盖了大数据课程设计里最核心的几个技术点。如果你正在做类似题目建议先把离线训练跑通再逐步接入实时链路不要一上来就搞端到端联调否则出了问题很难定位是哪个环节的锅。希望帮到你。本文还有配套的精品资源点击获取