EMR Serverless Spark 基于 MinHash-LSH 实现 PB 级文本语义去重 4 倍加速
大模型训练数据是燃料质量是引擎。在语料准备的整条链路中文本去重是最基础也最绕不开的环节。重复语料不仅浪费算力还会导致模型过拟合直接影响生成质量。然而当数据规模达到PB 级别时去重任务本身就变成了一个计算密集型的性能黑洞——跑一整夜还跑不完数据倾斜一出现整个任务就卡死这些都是数据工程师再熟悉不过的场景。某企业此前在原有云平台上使用开源 Spark 集群进行大规模文本去重。迁移至阿里云 EMR Serverless Spark 后借助 MinHash-LSH 内置函数与 Fusion Engine 向量化加速去重性能提升 4 倍数据准备周期从天级降至小时级。本文拆解这一迁移实践背后的技术方案与关键优化点。文本去重大模型语料清洗的关键一环在大型语言模型LLM的训练流程中语料数据的质量直接决定模型的最终效果。训练数据中的重复内容会导致三个核心问题计算资源浪费重复文本被反复处理消耗额外的 GPU/CPU 算力模型过拟合风险模型对重复内容产生记忆效应降低泛化能力评估失真测试集与训练集存在重复时评估指标虚高文本去重的本质是相似度检测——在海量文本中找出内容相同或高度相似的文档只保留代表性副本。当数据规模较小时精确比较尚可应对但当文档数量达到亿级别两两比较的 O(n²) 复杂度便成为不可承受之重这就需要更高效的算法方案。原有架构的瓶颈开源 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 超时/失败运维成本自管集群需持续投入人力维护扩展弹性需手动扩缩容响应滞后技术方案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: ARRAYSTRING,-- 分词后的词元数组perms_a: ARRAYBIGINT,-- MinHash 哈希函数组的乘数参数perms_b: ARRAYBIGINT,-- MinHash 哈希函数组的加数参数hash_ranges: ARRAYINT,-- Band 划分边界 [0, R, 2R, ..., B*R]ngram_size:INT,-- n-gram 大小长文本建议 5-9min_length:INT-- 输入 tokens 最小长度)-- 返回 ARRAYSTRING每个元素为对应 band 的十六进制哈希值build_lsh_edges函数对落入同一 LSH 桶的文档 ID基于最小节点连接策略生成边集用于后续图连通分量分析以聚类重复文档。build_lsh_edges(doc_ids: ARRAYBIGINT)-- 返回 ARRAYSTRUCTsrc: LONG, dst: LONG-- 示例桶内 ID 为 [1003, 1001, 1005] → 取最小 1001-- 生成边 (1001,1003) 和 (1001,1005)两个函数将算法逻辑下沉到引擎层开发者无需自行实现复杂的哈希逻辑代码量减少约 40%。minhash_lsh 和 build_lsh_edges 只是 Serverless Spark 内置函数生态的一部分。平台还内置了 ai_queryLLM 调用、ai_embedding_multimodal多模态 Embedding等 AI 函数可在同一 Spark SQL 会话、Spark任务中直接调用无需额外搭建推理服务。3.2 Fusion EngineSpark 原生向量化计算加速Serverless Spark 内置 Fusion EngineSpark Native Engine这是阿里云优化的向量化执行引擎相对开源版本性能提升 300%。在文本去重场景中Fusion Engine 带来三个关键优势向量化哈希计算MinHash 签名生成的大规模哈希运算在列式内存上按批向量化执行摊薄逐行处理的固定开销单条文档处理耗时显著降低消除 Python UDF 跨进程开销哈希逻辑以 C 内置函数在引擎内直接执行不再经 JVM 与 Python 进程间传输数据彻底省去跨进程序列化/反序列化成本Shuffle 稳定性优化针对存算分离架构进行了专门的 Shuffle 优化有效解决数据倾斜场景下的性能瓶颈在该客户的实际业务场景中取得了 4 倍性能提升的实测结果。迁移实践三步完成平滑迁移该企业的迁移过程分三个阶段稳步推进整体迁移成本可控。4.1 数据迁移将原始文本数据从原有云存储迁移至阿里云 OSS。Serverless Spark 原生支持 OSS-HDFS 协议完全兼容 HDFS 的云上存储确保数据访问的透明性与一致性。迁移过程中通过 checksum 校验确保数据完整性。4.2 代码迁移得益于 Spark API 的完全兼容性原有 PySpark 去重脚本迁移成本极低。核心改动仅需将数据读写路径替换为 OSS并引入minhash_lsh和build_lsh_edges内置函数替代原有自行实现的哈希逻辑。以下是关键代码片段# 1. 读取数据并生成 MinHash 签名hash_dfdf \.select(index_column,sf.split(sf.lower(text_column),patternSPLIT_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_dfhash_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. 图连通分量分析聚类重复文档assignmentGraphFrame(vertices_df,edges_df).connectedComponents()# 4. 保留每个连通分量中 ID 最小的代表文档dfdf.join(assignment.select(sf.col(id).alias(index_column),sf.col(component).alias(__component__)),onindex_column,howleft)\.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.partitions10001TB 以下每增加 1TB 加 1000防止单 task 数据倾斜或 OOMspark.sql.files.maxPartitionBytes256MB控制读取阶段分片大小spark.rdd.ensureConfigConsistencytrue必填项确保 RDD 配置一致性spark.executor.cores / memory4 核 / 14GB 2GB overhead建议 4:16 的 CPU 与内存比例效果验证4 倍性能提升的业务价值迁移完成后该企业对同一批文本数据进行了去重性能对比测试。测试使用相同的 MinHash-LSH 算法参数num_perm256, threshold0.8, ngram_size5。指标原有架构Serverless Spark 新架构提升幅度去重任务总耗时1-2天数小时4-5 倍提升Shuffle 失败率频繁失败零失败稳定性大幅改善运维投入需专职团队近零运维Serverless 免运维以阿里云官方文档中的 fineweb-edu 数据集为例进行验证使用 sample/10BT 子集2.15GB727,000 条文档在 Serverless Spark 上运行 MinHash-LSH 去重最终去除 2,191 条重复项保留 724,809 条文档去重过程高效且准确。业务价值加速模型迭代语料清洗耗时缩短 75%数据准备周期从天级降至小时级降低计算成本Serverless 按量计费模式避免了闲置资源浪费释放团队精力无需关注集群运维团队聚焦数据质量优化与模型效果提升弹性应对峰值Serverless 架构可秒级弹性扩容无需提前规划容量FAQQ1MinHash-LSH 内置函数支持哪些引擎版本支持 esr-4.xesr-4.1.1 及之后、esr-3.xesr-3.1.1 及之后、esr-2.xesr-2.5.1 及之后版本。建议使用最新版本以获得最佳性能。Q2从原有云平台迁移到 Serverless Spark 的代码改动量大吗改动量很小。Spark API 完全兼容主要改动集中在数据读写路径替换和引入内置函数替代自行实现的哈希逻辑。根据该客户实践代码量减少约 40%。Q3Serverless Spark 适合多大规模的文本去重任务Serverless Spark 采用弹性伸缩架构可从 GB 级到 PB 级灵活适配。1TB 以下数据建议配置 1000 个 shuffle 分区每增加 1TB 增加 1000 个分区。Q4MinHash-LSH 去重的精度如何控制通过 num_perm签名长度、threshold相似度阈值和 LSH 的 B/R 参数组合控制。num_perm256 threshold0.8 是推荐的平衡配置可在召回率和精度之间取得良好平衡。Q5除了文本去重Serverless Spark 还能用于哪些大模型数据预处理场景Serverless Spark 面向 DataAI 场景设计还支持数据清洗、特征工程、向量计算、多模态数据处理等场景。内置 AI Function 能力允许在 Spark 作业中直接调用大模型实现端到端的数据处理流水线。总结文本去重是大模型语料清洗的核心环节也是数据质量保障的基础。该企业从原有云平台迁移至阿里云 EMR Serverless Spark 的实践表明通过 MinHash-LSH 内置函数与 Fusion Engine 向量化加速的深度协同文本去重性能可获得 4 倍提升同时彻底释放运维负担。核心优势概括如下少写代码— MinHash-LSH 算法逻辑封装为内置函数开发者无需自行实现哈希与图分析逻辑代码量减少 40%少调集群— Serverless 架构免运维无需关注版本升级、资源调度和故障排查跑得更快— Fusion Engine 向量化加速 Shuffle 稳定性优化实测性能提升 4 倍用得更稳— 零 Shuffle 失败弹性扩缩容应对数据峰值阿里云 EMR Serverless Spark 作为面向 DataAI 的高性能 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

相关新闻

多无人机协同作业算法在农业植保中的优化与应用

多无人机协同作业算法在农业植保中的优化与应用

1. 项目背景与核心价值多无人机系统在农业植保领域的应用已经成为精准农业的重要技术支撑。2022年发表在BE SCI二区Top期刊的这项研究,针对作物保护场景中的无人机协同作业问题,提出了创新的任务分配算法。我在实际农业无人机项目中发现,传统…

2026/9/23 22:41:43 阅读更多 →
半迭代探索:平衡确定性与灵活性的工程实践

半迭代探索:平衡确定性与灵活性的工程实践

1. 项目概述:什么是半迭代探索半迭代探索(Semi-Iterative Exploration)是一种介于完全随机探索和系统化探索之间的实验方法。在我的工程实践中,这种技术特别适用于资源有限但需要快速验证假设的场景。与传统的瀑布式开发或纯敏捷开…

2026/9/18 0:30:05 阅读更多 →
3分钟学会使用Balena Etcher:安全烧录SD卡和USB驱动器的终极指南

3分钟学会使用Balena Etcher:安全烧录SD卡和USB驱动器的终极指南

3分钟学会使用Balena Etcher:安全烧录SD卡和USB驱动器的终极指南 【免费下载链接】etcher Flash OS images to SD cards & USB drives, safely and easily. 项目地址: https://gitcode.com/GitHub_Trending/et/etcher 你是否曾经因为烧录系统镜像而丢失重…

2026/9/24 15:17:47 阅读更多 →

最新新闻

2026年落地窗源头直供按需定制行业发展现状与市场占有率及排名研究分析报告

2026年落地窗源头直供按需定制行业发展现状与市场占有率及排名研究分析报告

2026年落地窗源头直供按需定制行业发展现状与市场占有率及排名研究分析报告 湖南皓思门窗有限公司皓思门窗落地窗源头直供品牌厂家、落地窗源头制造企业、落地窗源头实力厂商 长沙本土家装门窗需求升级,落地窗定制成核心痛点解决方案当下家装消费中,封阳…

2026/9/25 22:10:45 阅读更多 →
AI时代程序员:告别“被取代焦虑”,迎接“协作工程师”新身份

AI时代程序员:告别“被取代焦虑”,迎接“协作工程师”新身份

这一现象指向了一个更加深层的规律, 它暗示出效率的大幅提升并不必然会引发岗位的消灭, 反过来, 却有可能推动市场规模的整体扩充。随着编码成本不断降低, 企业将会更有勇气去尝试多样化的数字化创新举措, 这些创新又进一步催生了针对软件领域的旺盛需求。这与银行业里自动取款…

2026/9/25 22:10:45 阅读更多 →
文生视频提示词教程完整版|全套可直接复制提示词库(商用通用)

文生视频提示词教程完整版|全套可直接复制提示词库(商用通用)

文章目录一、万能基础模板(全模型通用)1. 通用正向提示词(基础稳定版)2. 通用负面提示词(所有场景必加)二、MiniMax H3 数字人口播专属模板(量产定稿)1. 口播正向提示词(…

2026/9/25 22:09:45 阅读更多 →
校园论文选题系统开发实战:Laravel+uniapp+微信小程序

校园论文选题系统开发实战:Laravel+uniapp+微信小程序

毕业论文选题,每年春季都是高校信息部门最头疼的环节。纸质表格传阅、Excel来回汇总、学生线下找老师签字协调,一套流程走下来少说两周,还免不了各种重复和错漏。后来我接手了一个校园团队的项目,用 Thinkphp/Laravel 作为后端、u…

2026/9/25 22:08:45 阅读更多 →
zvec-grep混合搜索原理揭秘:BM25、向量检索与ripgrep如何用RRF融合排名

zvec-grep混合搜索原理揭秘:BM25、向量检索与ripgrep如何用RRF融合排名

zvec-grep混合搜索原理揭秘:BM25、向量检索与ripgrep如何用RRF融合排名 【免费下载链接】zvec-grep Local-first search across your workspace, built for humans and AI agents. 项目地址: https://gitcode.com/gh_mirrors/zv/zvec-grep zvec-grep&#xf…

2026/9/25 22:08:45 阅读更多 →
SpringBoot+Vue 实现办公用品管理系统|计算机毕设源码讲解

SpringBoot+Vue 实现办公用品管理系统|计算机毕设源码讲解

💖💖作者:计算机毕业设计小明哥 💙💙个人简介:曾长期从事计算机专业培训教学,本人也热爱上课教学,语言擅长Java、微信小程序、Python、Golang、安卓Android等,开发项目包…

2026/9/25 22:07:44 阅读更多 →

日新闻

AI元人文:从工具使用到思维重构的深度探索

AI元人文:从工具使用到思维重构的深度探索

最近半年我一直在琢磨一件事:AI元人文到底是什么?说白了,就是“用元视角重新审视人与AI的关系”,也在“探索AI如何反向逼着我们发现自己的思考边界”。标题里的“元探索”,在我看就是一层套一层的追问——当你用AI解决…

2026/9/25 0:00:41 阅读更多 →
Python+CNN车牌识别实战:从数据预处理到模型训练与部署

Python+CNN车牌识别实战:从数据预处理到模型训练与部署

简介:基于Python与卷积神经网络的车牌识别项目,面向计算机视觉初学者及智能交通开发者,目标是帮助用户掌握从数据预处理、模型构建到实际部署的完整流程。压缩包共25个文件,包含jpg/png图像样本、py训练脚本、md说明文档、dat数据…

2026/9/25 0:00:41 阅读更多 →
Vim基础操作全攻略:保存退出、模式切换与高频命令实战

Vim基础操作全攻略:保存退出、模式切换与高频命令实战

1. 项目概述1.1 核心需求解析今天聊聊Vim。写这个题目的原因是:几乎每个后端开发者、运维人员、数据工程师某天都会遇到一个场景——深夜加班,服务器登录界面只有黑底白字,编辑器只有vi/vim,你必须在五分钟内完成一次配置修改并保…

2026/9/25 0:00:41 阅读更多 →

周新闻

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

直接铺开项目本身吧。这几个月我一直在折腾一件事:用Flutter给OpenHarmony做一款游戏集合类的App,说白了就是把若干小游戏塞进一个壳里,用统一入口分发。这个方向本身不算新鲜,真正让我花了不少心思的,是首页那堆游戏卡…

2026/9/25 19:27:14 阅读更多 →
Word表格编号全攻略:从列表编号到题注交叉引用

Word表格编号全攻略:从列表编号到题注交叉引用

写Word文档,最让人头疼的往往是那些“看起来不起眼”的小问题。比如表格编号这事:今天在表后面多加了两个空白行,明天给客户交稿前发现整个章节的编号全部错位,光是挨个改序号就能耗掉大半个下午。我前阵子帮人整理一份上百页的技…

2026/9/25 11:15:26 阅读更多 →
从第一个站到第二个站:独立开发者的静态网站选型与落地实践

从第一个站到第二个站:独立开发者的静态网站选型与落地实践

1. 项目概述1.1 核心需求解析做独立开发者这几年,说实话,第一个网站上线的那天晚上我兴奋得没睡着。但等它跑了半年,流量惨淡、功能臃肿、代码自己都懒得看第二遍之后,我才慢慢琢磨明白一个道理:第一个网站是练手&…

2026/9/25 20:29:09 阅读更多 →

月新闻

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能分类:[AI/大模型]细分主题:AI 增强型 CI/CD 流水线自动化与 GitOps 实践:Agent 工作流、工具调用与任务拆解:从原型到生产的验收清单很多团队在尝试用大…

2026/9/25 20:29:43 阅读更多 →
容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场分类:[工程技术]细分主题:Kubernetes 生产环境运维与排障实战:可复制的项目复盘模板与决策记录大部分团队的事故复盘报告,最后都变成了躺在 Confluence 或钉…

2026/9/25 20:29:31 阅读更多 →
容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步分类:[工程技术]细分主题:Docker 容器化技术与镜像安全管理:核心链路的逐步实现与关键代码取舍面对一个积累了五六年历史包袱的单体架构应用(包含 Web 接口、后台…

2026/9/25 19:27:26 阅读更多 →