共计 2160 个字符,预计需要花费 6 分钟才能阅读完成。
背景痛点
处理大规模人脸数据集时,单机方案通常会遇到几个典型问题:

- 内存溢出:300w 张人脸图片(假设每张图片 100KB)加载到内存需要约 300GB,远超单机内存容量
- I/ O 等待:传统机械硬盘顺序读取速度约 100MB/s,完整加载数据集需要 50 分钟以上
- 特征提取速度慢:使用 OpenCV 的 DNN 模块处理单张图片约需 100ms,单线程处理全部数据需要 34 小时
存储层面对比:
| 方案 | TPS(事务 / 秒) | QPS(查询 / 秒) | 存储效率 |
|---|---|---|---|
| MySQL 单机 | 500-1000 | 200-500 | 30-40% |
| HDFS 集群(3 节点) | 5000+ | 3000+ | 70-80% |
技术方案
架构设计
采用 HDFS+Spark 的 Lambda 架构,关键设计点:
- 数据分片策略:
- 将原始图片按 128MB/block 大小分片存储
- 每个 HDFS block 对应 1 个 Spark partition
-
目录结构示例:
/dataset/yyyyMMdd/hour=00/block_001.parquet -
核心优化手段:
- 列式存储:使用 Parquet 格式存储特征向量,实测压缩比达 1:4
- 索引构建:采用 FAISS 的 IVF2048 索引类型,召回率 >95%
- 内存优化:Executor 内存分配公式:
(节点内存 - 1GB) * 0.8 / executor 数
代码示例
特征提取 UDF
from pyspark.sql.functions import udf
from pyspark.sql.types import ArrayType, FloatType
import cv2
@udf(ArrayType(FloatType()))
def extract_face_features(img_bytes):
"""
:param img_bytes: 二进制图片数据
:return: 512 维特征向量
"""
nparr = np.frombuffer(img_bytes, np.uint8)
img = cv2.imdecode(nparr, cv2.IMREAD_COLOR)
# 使用 OpenCV FaceRecognizer
recognizer = cv2.dnn.readNetFromTorch('openface.nn4.small2.v1.t7')
blob = cv2.dnn.blobFromImage(img, 1.0/255, (96,96), (0,0,0), swapRB=True)
recognizer.setInput(blob)
return recognizer.forward().flatten().tolist()
分布式 JOIN 优化
# 原始数据加载
df_images = spark.read.parquet("hdfs:///dataset/faces")
df_metadata = spark.read.parquet("hdfs:///dataset/metadata")
# 优化后的 JOIN 操作
optimized_join = df_images.repartition(1024, "image_id") \
.join(df_metadata.repartition(1024, "image_id"), "image_id") \
.persist(StorageLevel.MEMORY_AND_DISK)
性能对比
测试环境:AWS c5.4xlarge (16vCPU, 32GB 内存)
| 处理阶段 | 单机耗时 | 10 节点集群耗时 | 加速比 |
|---|---|---|---|
| 数据加载 | 48min | 3min | 16x |
| 特征提取 | 34h | 1.8h | 18.8x |
| 索引构建 | 6h | 25min | 14.4x |
内存使用监控关键指标:
– GC 频率:<5 次 / 分钟
– Executor 内存利用率:75-85%
避坑指南
- 小文件合并:
- 使用 Spark 的
coalesce()代替repartition() -
设置
spark.sql.shuffle.partitions=2000 -
特征维度对齐:
# 错误示例:不同模型产生的维度不同 model1_feat = [0.1, 0.2] # 512 维 model2_feat = [0.3, 0.4] # 256 维 # 正确做法:统一使用相同模型 assert len(extract_face_features(img)) == 512 -
动态资源分配陷阱:
- 关闭
spark.dynamicAllocation.enabled - 固定
executor.instances= 节点数 *2
动手实验
# docker-compose.yml
version: '3'
services:
namenode:
image: bde2020/hadoop-namenode:2.0.0-hadoop3.2.1-java8
spark-master:
image: bitnami/spark:3.3.0
depends_on: [namenode]
spark-worker:
image: bitnami/spark:3.3.0
scale: 3
部署命令:
docker-compose up -d
docker exec spark-master spark-submit \
--master spark://spark-master:7077 \
--executor-memory 4G \
your_script.py
合规说明
本文实验数据采用 LFW(Labeled Faces in the Wild)公开数据集,所有示例代码已做脱敏处理。实际业务场景中请确保遵守当地数据隐私法规。
正文完
发表至: 未分类
近两天内
