大模型训练,数据是燃料,质量是引擎。
在语料准备的整条链路中,文本去重是最基础也最绕不开的环节。重复语料不仅浪费算力,还会导致模型过拟合,直接影响生成质量。然而,当数据规模达到PB 级别时,去重任务本身就变成了一个计算密集型的"性能黑洞"——跑一整夜还跑不完,数据倾斜一出现整个任务就卡死,这些都是数据工程师再熟悉不过的场景。
某企业此前在原有云平台上使用开源 Spark 集群进行大规模文本去重。迁移至阿里云 EMR Serverless Spark 后,借助 MinHash-LSH 内置函数与 Fusion Engine 向量化加速,去重性能提升 4 倍,数据准备周期从天级降至小时级。
本文拆解这一迁移实践背后的技术方案与关键优化点。
01 文本去重:大模型语料清洗的关键一环
在大型语言模型(LLM)的训练流程中,语料数据的质量直接决定模型的最终效果。训练数据中的重复内容会导致三个核心问题:
- 计算资源浪费:重复文本被反复处理,消耗额外的 GPU/CPU 算力
- 模型过拟合风险:模型对重复内容产生记忆效应,降低泛化能力
- 评估失真:测试集与训练集存在重复时,评估指标虚高
文本去重的本质是相似度检测——在海量文本中找出内容相同或高度相似的文档,只保留代表性副本。当数据规模较小时,精确比较尚可应对;但当文档数量达到亿级别,两两比较的 O(n²) 复杂度便成为不可承受之重,这就需要更高效的算法方案。
02 原有架构的瓶颈:开源 Spark 去重之困
该企业此前在原有云平台上构建了数据处理平台,使用开源 Spark 集群运行 MinHash-LSH 去重算法。随着数据规模持续增长,三个瓶颈逐渐显现。
2.1 计算效率瓶颈
开源 Spark 的执行引擎基于 JVM 采用行式(row-based)迭代模型,逐行处理数据。在大规模 n-gram 分词和多组哈希这类计算密集型操作上,逐行执行带来大量虚函数调用和对象封装/拆箱开销,CPU cache 利用率低,存在天然的性能瓶颈。更关键的是,原方案的哈希逻辑以 Python UDF 实现,数据需在 JVM 与 Python 进程之间跨进程传输、序列化与反序列化,进一步放大了计算开销,导致 CPU 算力无法充分释放。
2.2 Shuffle 稳定性问题
MinHash-LSH 算法中的 LSH 分桶和图连通分量计算阶段涉及大量 Shuffle 操作。在开源 Spark 环境下, Shuffle 稳定性问题经常发生——尤其在数据倾斜场景下,容易出现任务超时甚至失败,需要人工介入调优。
2.3 运维成本高昂
维护自管 Spark 集群意味着持续投入人力进行版本升级、资源调度和故障排查。随着业务规模扩大,运维成本占比逐年攀升,团队希望将精力聚焦于业务逻辑而非基础设施管理。
痛点维度 |
原有架构表现 |
计算效率 |
开源 Spark 无向量化加速,Python UDF 跨进程开销大 |
Shuffle 稳定性 |
Shuffle 超时/失败 |
运维成本 |
自管集群,需持续投入人力维护 |
扩展弹性 |
需手动扩缩容,响应滞后 |
03 技术方案:MinHash-LSH 内置函数 + Fusion Engine
迁移至阿里云 EMR Serverless Spark 后,该企业采用了一套全新的文本去重技术方案。核心由两大能力支撑:将 MinHash-LSH 算法深度集成到 Spark Dataframe/SQL 引擎的内置函数,以及提供向量化加速和Shuffle稳定性的 Fusion Engine。
3.1 MinHash-LSH:给每篇文档生成"指纹身份证"
MinHash-LSH 是一种经典的近似相似性检测算法组合,广泛应用于大规模集合相似度计算(如 Jaccard 相似度)。其核心分为两步:
第一步:MinHash——生成签名向量
将文本转换为 n-gram 集合后,通过多组哈希函数生成紧凑的签名向量(Signature)。可以理解为给每篇文档发一张"指纹身份证"——原始文本可能数 KB,但签名向量固定长度(如 256 位),保留了集合间的相似性特征,后续比较只需对比签名而非全文。
第二步:LSH——分诊台快速分流
将签名向量划分为多个"band",每个 band 单独哈希。高相似度的文档更可能落入同一哈希桶中,只有落入同一桶的文档对才需要进一步比较。这相当于在医院分诊台快速将相似症状的患者分流到同一科室,将 O(n²) 的全量比较降为近线性复杂度。
Serverless Spark 通过两个内置函数实现这一能力:
minhash_lsh 函数:将输入文本分词后生成 MinHash 签名,并按 bands 划分生成对应的哈希值列表。
minhash_lsh( tokens: ARRAY<STRING>, -- 分词后的词元数组 perms_a: ARRAY<BIGINT>, -- MinHash 哈希函数组的乘数参数 perms_b: ARRAY<BIGINT>, -- MinHash 哈希函数组的加数参数 hash_ranges: ARRAY<INT>, -- Band 划分边界 [0, R, 2R, ..., B*R] ngram_size: INT, -- n-gram 大小,长文本建议 5-9 min_length: INT -- 输入 tokens 最小长度 ) -- 返回 ARRAY<STRING>,每个元素为对应 band 的十六进制哈希值
build_lsh_edges 函数:对落入同一 LSH 桶的文档 ID,基于"最小节点连接"策略生成边集,用于后续图连通分量分析以聚类重复文档。
build_lsh_edges(doc_ids: ARRAY<BIGINT>) -- 返回 ARRAY<STRUCT<src: LONG, dst: LONG>> -- 示例:桶内 ID 为 [1003, 1001, 1005] → 取最小 1001 -- 生成边 (1001,1003) 和 (1001,1005)
两个函数将算法逻辑下沉到引擎层,开发者无需自行实现复杂的哈希逻辑,代码量减少约 40%。
minhash_lsh 和 build_lsh_edges 只是 Serverless Spark 内置函数生态的一部分。平台还内置了 ai_query(LLM 调用)、ai_embedding_multimodal(多模态 Embedding)等 AI 函数,可在同一 Spark SQL 会话、Spark任务中直接调用,无需额外搭建推理服务。
3.2 Fusion Engine:Spark 原生向量化计算加速
Serverless Spark 内置 Fusion Engine(Spark Native Engine),这是阿里云优化的向量化执行引擎,相对开源版本性能提升 300%。在文本去重场景中,Fusion Engine 带来三个关键优势:
- 向量化哈希计算:MinHash 签名生成的大规模哈希运算在列式内存上按批(向量化)执行,摊薄逐行处理的固定开销,单条文档处理耗时显著降低
- 消除 Python UDF 跨进程开销:哈希逻辑以 C++ 内置函数在引擎内直接执行,不再经 JVM 与 Python 进程间传输数据,彻底省去跨进程序列化/反序列化成本
- Shuffle 稳定性优化:针对存算分离架构进行了专门的 Shuffle 优化,有效解决数据倾斜场景下的性能瓶颈
在该客户的实际业务场景中,取得了 4 倍性能提升的实测结果。
04 迁移实践:三步完成平滑迁移
该企业的迁移过程分三个阶段稳步推进,整体迁移成本可控。
4.1 数据迁移
将原始文本数据从原有云存储迁移至阿里云 OSS。Serverless Spark 原生支持 OSS-HDFS 协议,完全兼容 HDFS 的云上存储,确保数据访问的透明性与一致性。迁移过程中通过 checksum 校验确保数据完整性。
4.2 代码迁移
得益于 Spark API 的完全兼容性,原有 PySpark 去重脚本迁移成本极低。核心改动仅需将数据读写路径替换为 OSS,并引入 minhash_lsh 和 build_lsh_edges 内置函数替代原有自行实现的哈希逻辑。以下是关键代码片段:
# 1. 读取数据并生成 MinHash 签名 hash_df = df \ .select(index_column, sf.split(sf.lower(text_column), pattern=SPLIT_PATTERN.pattern).alias("tokens")) \ .select(index_column, sf.minhash_lsh("tokens", a.tolist(), b.tolist(), HASH_RANGES_SLICE, ngram_size, min_length).alias("hashes")) \ .select(index_column, sf.posexplode("hashes").alias("band_idx", "band_hash")) # 2. 对同一 LSH 桶的文档生成边集 edges_df = hash_df.groupBy("band_idx", "band_hash") \ .agg(sf.count(index_column).alias("cnt"), sf.collect_list(index_column).alias("doc_ids")) \ .filter(sf.col("cnt") > 1) \ .select(sf.build_lsh_edges("doc_ids").alias("edges")) \ .select(sf.explode("edges").alias("edge")) \ .selectExpr("edge.src as src", "edge.dst as dst") # 3. 图连通分量分析,聚类重复文档 assignment = GraphFrame(vertices_df, edges_df).connectedComponents() # 4. 保留每个连通分量中 ID 最小的代表文档 df = df.join(assignment.select( sf.col("id").alias(index_column), sf.col("component").alias("__component__")), on=index_column, how="left") \ .filter(sf.col("__component__").isNull() | (sf.col("__component__") == sf.col(index_column))) \ .drop("__component__")
4.3 资源配置
迁移至 Serverless 架构后,无需再维护固定大小的集群。按需配置 executor 资源(建议 4 CPU : 16 GB 内存比例),系统自动完成资源的弹性分配与回收。
配置项 |
推荐设置 |
说明 |
spark.sql.shuffle.partitions |
1000(1TB 以下),每增加 1TB 加 1000 |
防止单 task 数据倾斜或 OOM |
spark.sql.files.maxPartitionBytes |
256MB |
控制读取阶段分片大小 |
spark.rdd.ensureConfigConsistency |
true |
必填项,确保 RDD 配置一致性 |
spark.executor.cores / memory |
4 核 / 14GB + 2GB overhead |
建议 4:16 的 CPU 与内存比例 |
05 效果验证:4 倍性能提升的业务价值
迁移完成后,该企业对同一批文本数据进行了去重性能对比测试。测试使用相同的 MinHash-LSH 算法参数(num_perm=256, threshold=0.8, ngram_size=5)。
指标 |
原有架构 |
Serverless Spark 新架构 |
提升幅度 |
去重任务总耗时 |
1-2天 |
数小时 |
4-5 倍提升 |
Shuffle 失败率 |
频繁失败 |
零失败 |
稳定性大幅改善 |
运维投入 |
需专职团队 |
近零运维 |
Serverless 免运维 |
以阿里云官方文档中的 fineweb-edu 数据集为例进行验证:使用 sample/10BT 子集(2.15GB,727,000 条文档),在 Serverless Spark 上运行 MinHash-LSH 去重,最终去除 2,191 条重复项,保留 724,809 条文档,去重过程高效且准确。
业务价值:
- 加速模型迭代:语料清洗耗时缩短 75%,数据准备周期从天级降至小时级
- 降低计算成本:Serverless 按量计费模式避免了闲置资源浪费
- 释放团队精力:无需关注集群运维,团队聚焦数据质量优化与模型效果提升
- 弹性应对峰值:Serverless 架构可秒级弹性扩容,无需提前规划容量
06 FAQ
Q1:MinHash-LSH 内置函数支持哪些引擎版本?
支持 esr-4.x(esr-4.1.1 及之后)、esr-3.x(esr-3.1.1 及之后)、esr-2.x(esr-2.5.1 及之后)版本。建议使用最新版本以获得最佳性能。
Q2:从原有云平台迁移到 Serverless Spark 的代码改动量大吗?
改动量很小。Spark API 完全兼容,主要改动集中在数据读写路径替换和引入内置函数替代自行实现的哈希逻辑。根据该客户实践,代码量减少约 40%。
Q3:Serverless Spark 适合多大规模的文本去重任务?
Serverless Spark 采用弹性伸缩架构,可从 GB 级到 PB 级灵活适配。1TB 以下数据建议配置 1000 个 shuffle 分区,每增加 1TB 增加 1000 个分区。
Q4:MinHash-LSH 去重的精度如何控制?
通过 num_perm(签名长度)、threshold(相似度阈值)和 LSH 的 B/R 参数组合控制。num_perm=256 + threshold=0.8 是推荐的平衡配置,可在召回率和精度之间取得良好平衡。
Q5:除了文本去重,Serverless Spark 还能用于哪些大模型数据预处理场景?
Serverless Spark 面向 Data+AI 场景设计,还支持数据清洗、特征工程、向量计算、多模态数据处理等场景。内置 AI Function 能力允许在 Spark 作业中直接调用大模型,实现端到端的数据处理流水线。
07 总结
文本去重是大模型语料清洗的核心环节,也是数据质量保障的基础。该企业从原有云平台迁移至阿里云 EMR Serverless Spark 的实践表明,通过 MinHash-LSH 内置函数与 Fusion Engine 向量化加速的深度协同,文本去重性能可获得 4 倍提升,同时彻底释放运维负担。
核心优势概括如下:
- 少写代码 — MinHash-LSH 算法逻辑封装为内置函数,开发者无需自行实现哈希与图分析逻辑,代码量减少 40%
- 少调集群 — Serverless 架构免运维,无需关注版本升级、资源调度和故障排查
- 跑得更快 — Fusion Engine 向量化加速 + Shuffle 稳定性优化,实测性能提升 4 倍
- 用得更稳 — 零 Shuffle 失败,弹性扩缩容应对数据峰值
阿里云 EMR Serverless Spark 作为面向 Data+AI 的高性能 Lakehouse 产品,在 TPC-DS 100TB 基准测试中表现优异。无论是大模型语料清洗、数据湖分析还是 AI 数据预处理,Serverless Spark 都是值得考虑的方案。
了解更多:
- 产品文档:https://help.aliyun.com/zh/emr/emr-serverless-spark/
- MinHash-LSH 去重方案:https://help.aliyun.com/zh/emr/emr-serverless-spark/use-cases/minhash-lsh-based-large-scale-text-duplication-scheme