EMR Serverless Spark PB级文本语义去重4倍加速的技术方案解读

简介: 针对大模型语料清洗中文本去重面临的性能瓶颈,某企业迁移至阿里云emr serverless spark后实现突破。新方案通过minhash-lsh内置函数将算法下沉引擎层,减少40%代码量;结合fusion engine向量化加速与shuffle优化,消除python udf跨进程开销并解决数据倾斜问题。实测去重性能提升4倍,任务耗时从天级降至小时级,且实现零shuffle失败与免运维。该实践验证了serverless架构在pb级数据预处理中的高效性与稳定性,显著加速模型迭代并降低计算成本。

大模型训练,数据是燃料,质量是引擎。

在语料准备的整条链路中,文本去重是最基础也最绕不开的环节。重复语料不仅浪费算力,还会导致模型过拟合,直接影响生成质量。然而,当数据规模达到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_lshbuild_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 都是值得考虑的方案。

了解更多:

相关文章
|
存储 缓存 文件存储
如何保证分布式文件系统的数据一致性
分布式文件系统需要向上层应用提供透明的客户端缓存,从而缓解网络延时现象,更好地支持客户端性能水平扩展,同时也降低对文件服务器的访问压力。当考虑客户端缓存的时候,由于在客户端上引入了多个本地数据副本(Replica),就相应地需要提供客户端对数据访问的全局数据一致性。
33254 201
如何保证分布式文件系统的数据一致性
|
设计模式 存储 监控
设计模式(C++版)
看懂UML类图和时序图30分钟学会UML类图设计原则单一职责原则定义:单一职责原则,所谓职责是指类变化的原因。如果一个类有多于一个的动机被改变,那么这个类就具有多于一个的职责。而单一职责原则就是指一个类或者模块应该有且只有一个改变的原因。bad case:IPhone类承担了协议管理(Dial、HangUp)、数据传送(Chat)。good case:里式替换原则定义:里氏代换原则(Liskov 
36823 22
设计模式(C++版)
|
存储 编译器 C语言
抽丝剥茧C语言(初阶 下)(下)
抽丝剥茧C语言(初阶 下)
|
机器学习/深度学习 人工智能 自然语言处理
带你简单了解Chatgpt背后的秘密:大语言模型所需要条件(数据算法算力)以及其当前阶段的缺点局限性
带你简单了解Chatgpt背后的秘密:大语言模型所需要条件(数据算法算力)以及其当前阶段的缺点局限性
24905 16
|
机器学习/深度学习 弹性计算 监控
重生之---我测阿里云U1实例(通用算力型)
阿里云产品全线降价的一力作,2023年4月阿里云推出新款通用算力型ECS云服务器Universal实例,该款服务器的真实表现如何?让我先测为敬!
36824 15
重生之---我测阿里云U1实例(通用算力型)
|
SQL 存储 弹性计算
Redis性能高30%,阿里云倚天ECS性能摸底和迁移实践
Redis在倚天ECS环境下与同规格的基于 x86 的 ECS 实例相比,Redis 部署在基于 Yitian 710 的 ECS 上可获得高达 30% 的吞吐量优势。成本方面基于倚天710的G8y实例售价比G7实例低23%,总性价比提高50%;按照相同算法,相对G8a,性价比为1.4倍左右。
|
存储 算法 Java
【分布式技术专题】「分布式技术架构」手把手教你如何开发一个属于自己的限流器RateLimiter功能服务
随着互联网的快速发展,越来越多的应用程序需要处理大量的请求。如果没有限制,这些请求可能会导致应用程序崩溃或变得不可用。因此,限流器是一种非常重要的技术,可以帮助应用程序控制请求的数量和速率,以保持稳定和可靠的运行。
29949 52

热门文章

最新文章