PB级文本语义去重实战:基于MinHash-LSH与EMR Serverless Spark的4倍加速方案 1. 项目背景与核心挑战当PB级文本遇上语义去重最近在做一个数据治理的项目客户那边堆积了海量的用户生成内容包括评论、帖子、客服对话记录等等数据量级已经达到了PB级别。他们的核心需求很简单把这些文本数据里语义上重复的内容找出来并合并为后续的分析和模型训练提供一份“干净”的数据集。听起来像是文本去重的常规操作对吧但实际操作起来问题就复杂了。传统的去重方法比如基于精确字符串匹配或者简单的Jaccard相似度基于词袋模型在这个场景下几乎完全失效。因为用户表达同一个意思的方式千差万别。比如“这个手机电池续航太差了”和“这款手机的待机时间令人失望”这两句话从字面上看完全不同但语义高度一致。我们需要的是语义去重而不是字面去重。这就引出了第一个技术选型用Embedding向量化表示。我们可以通过预训练的语言模型比如BERT、Sentence-BERT将每段文本转换成一个高维向量例如768维。语义相似的文本其向量在空间中的距离比如余弦相似度也会很近。理想情况下我们计算所有文本向量两两之间的相似度设定一个阈值比如0.9低于这个阈值的就认为是重复的。然而PB级的数据量让这个“理想情况”变成了计算灾难。假设我们有1万亿1T条文本进行两两比较的复杂度是O(N²)这需要进行的比较次数是一个天文数字即使使用再强大的计算集群在时间和成本上都是不可接受的。这就是我们面临的核心挑战如何在可接受的时间和资源成本下对海量文本进行高效的语义相似度搜索和去重。直接进行暴力比对的路走不通我们必须引入一种能够将高维空间近邻搜索问题大幅降维的算法。这就是局部敏感哈希Locality-Sensitive Hashing, LSH登场的时刻。LSH的核心思想是设计一种哈希函数使得原本在高维空间中距离近的点经过哈希后有高概率被映射到同一个桶Bucket里而距离远的点则大概率被映射到不同的桶。这样我们只需要在同一个桶内的文本之间进行精细的相似度计算从而将计算量从全局比对压缩到局部比对。而MinHash正是为处理集合相似度如Jaccard相似度设计的一种LSH方法。虽然我们最终目标是语义向量相似度但可以通过一些巧妙的转换利用MinHash-LSH来为我们的高维向量近似搜索服务这正是本项目的技术核心。2. 技术选型为什么是EMR Serverless Spark MinHash-LSH面对PB级数据处理和复杂的近似最近邻搜索任务技术栈的选择直接决定了项目的成败。我们最终敲定了阿里云EMR Serverless Spark作为计算引擎并结合MinHash-LSH算法这是一套经过深思熟虑的组合拳。2.1 为什么选择EMR Serverless Spark首先数据量是PB级这注定是一个分布式计算任务。Spark以其卓越的内存计算能力和成熟的生态系统MLlib、Spark SQL成为首选。但为什么是Serverless版本极致的弹性与成本优化文本去重任务特别是前期特征提取文本转向量阶段是计算密集型而LSH分桶和桶内比对阶段则对网络Shuffle和内存有一定要求。任务负载是波动的。使用传统EMR集群我们需要预先规划好集群规模并承担集群空闲时的成本。而Serverless Spark允许我们按实际计算资源消耗付费。任务提交后平台自动分配资源任务结束即释放实现了真正的“用多少算多少”对于这种一次性或周期性的重型ETL任务成本效益极高。免运维聚焦业务逻辑我们不需要关心底层Hadoop/Spark集群的部署、监控、扩缩容和故障恢复。EMR Serverless提供了托管的Spark运行时环境我们只需要准备好JAR包或Python脚本定义好Executor和Driver的资源规格CPU、内存即可提交作业。这让我们团队能将全部精力投入到算法实现和性能调优上。与阿里云对象存储OSS无缝集成我们的原始文本数据存放在OSS中处理后的结果也需要写回OSS。EMR Serverless Spark原生支持OSS作为文件系统读写性能和数据一致性都有保障避免了数据搬迁的额外开销和风险。大规模Shuffle的稳定性LSH处理过程中需要根据哈希值对数据进行重新分区Repartition这会产生大量的Shuffle操作。EMR Serverless底层对Spark Shuffle进行了深度优化和托管能够更稳定地处理大规模Shuffle减少OOM内存溢出和Fetch Failed等错误的概率。2.2 为什么选择MinHash-LSH如前所述我们需要一个能将高维向量近似搜索问题“降维打击”的方法。MinHash-LSH的方案脱颖而出原理匹配度MinHash最初是为快速估算Jaccard相似度设计的。Jaccard相似度衡量的是两个集合的交集与并集之比。我们可以将一段文本经过分词后视为一个“词集合”。MinHash能够高效地估计两个文本词集合的Jaccard相似度。虽然这还不是语义相似度但词重叠度是语义相似的一个强相关信号。更重要的是我们可以利用这个特性作为第一层粗筛。与Embedding的协同方案我们并没有单纯依赖MinHash。实际的架构是两阶段漏斗型过滤第一阶段粗筛使用MinHash-LSH。将每条文本分词后的词序列视为一个集合通过MinHash算法生成一个由多个哈希值组成的“签名”Signature。然后通过LSH函数将具有相似签名的文本分到同一个桶里。这个阶段可以快速排除掉绝大多数明显不相似的文本对可能将需要精细比对的数据量减少99%以上。第二阶段精筛在MinHash分出的每个桶内对桶内的文本使用其真正的语义向量例如由Sentence-BERT生成进行余弦相似度计算。因为桶已经很小了所以这个计算变得可行。算法效率与Spark的亲和性MinHash的计算本身是高效的并且其过程特征哈希、取最小哈希值非常容易在Spark的map、reduce等算子中实现并行化。Spark MLlib库也提供了MinHashLSH的实现可以直接调用大大降低了工程复杂度。可调节的精度/召回率权衡LSH算法有几个关键参数如哈希函数的数量numHashTables和每个哈希函数的带宽bucketLength。通过调整这些参数我们可以在查全率Recall和计算复杂度之间进行灵活的权衡。需要更高的去重召回率尽量不漏掉重复项可以增加哈希表数量但这会增加计算和存储开销反之则可以减少开销接受一定的漏报率。这种可控性对于工程落地至关重要。注意这里有一个关键的思维转换。我们最终目标是语义去重但直接对768维的BERT向量做LSH如基于p-stable分布的LSH计算成本依然很高。而MinHash处理的是文本的词集维度词汇表大小可能上万但通过哈希函数将其映射到低维签名巧妙地规避了“维度灾难”。词重叠是语义相似的必要不充分条件因此适合作为高效的预过滤层。3. 系统架构与在EMR Serverless上的实现拆解整个去重流水线被设计成一个多阶段的Spark作业。下面我详细拆解每个阶段在EMR Serverless Spark上的实现逻辑、代码片段以及资源配置考量。3.1 数据读取与预处理阶段数据源是OSS上按日期分区的文本文件格式为JSONL每行一个JSON对象。我们使用Spark SQL进行读取并完成基础清洗。// 示例代码Scala Spark import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ val spark SparkSession.builder() .appName(PB-Text-Dedup) .getOrCreate() // 从OSS路径读取数据 val rawDataPath oss://your-bucket/path/to/text-data/*/*.jsonl val df spark.read.json(rawDataPath) .select(col(id), col(text)) // 假设有id和text字段 // 文本预处理清洗、分词 // 这里使用简单的空格分词实际中可能需要更复杂的分词器如HanLP val tokenizedDF df.withColumn(tokens, split(trim(lower(col(text))), \\s) ).filter(size(col(tokens)) 0) // 过滤掉空文本这个阶段在Serverless Spark上的配置要点Executor配置此阶段主要是IO密集型和轻量计算。我们为Executor配置了适中的CPU4核和足够的内存16GB以并行读取大量OSS小文件。关键是要设置足够的Executor数量如50-100个以并行读取OSS上的海量文件避免IO成为瓶颈。OSS优化确保OSS Bucket和EMR Serverless服务在同一地域避免跨地域网络延迟。对于大量小文件可以考虑使用spark.hadoop.mapreduce.input.fileinputformat.split.minsize参数来控制输入分片大小合并小文件提升读取效率。3.2 特征工程生成MinHash签名与LSH分桶这是核心阶段。我们使用Spark MLlib的MinHashLSH。import org.apache.spark.ml.feature.{MinHashLSH, HashingTF} import org.apache.spark.ml.linalg.Vectors // 1. 使用HashingTF将词序列转换为词频向量特征哈希 val hashingTF new HashingTF() .setInputCol(tokens) .setOutputCol(rawFeatures) .setNumFeatures(1 20) // 设置哈希表大小例如2^20约100万维 val featurizedData hashingTF.transform(tokenizedDF) // 2. 创建并训练MinHashLSH模型 val mh new MinHashLSH() .setNumHashTables(5) // LSH哈希表的数量影响召回率 .setInputCol(rawFeatures) .setOutputCol(hashes) val model mh.fit(featurizedData) // 3. 对数据集进行转换得到每个样本的哈希签名 val hashedDF model.transform(featurizedData) // 4. 关键步骤利用LSH进行近似相似项连接 // 这里会触发一个大规模的Shuffle因为需要根据哈希值进行自连接 val duplicatedPairs model.approxSimilarityJoin(hashedDF, hashedDF, 0.8, JaccardDistance) .filter(col(datasetA.id) col(datasetB.id)) // 避免重复和自匹配 .select( col(datasetA.id).alias(id1), col(datasetB.id).alias(id2), col(JaccardDistance) )参数解读与调优经验numHashTables这是最重要的参数。增加它会提高召回率找到更多真正重复的对但也会显著增加计算和存储开销因为每个数据点需要在多个哈希表中存储。对于PB级数据我们从3开始测试最终根据业务对漏报的容忍度选择了5。这是一个典型的“用空间换时间/召回率”的权衡。approxSimilarityJoin中的阈值0.8这里用的是Jaccard距离1 - Jaccard相似度。阈值0.8意味着我们寻找Jaccard相似度大于0.2的文本对。这个值设得很低目的是在第一阶段进行“宽筛”宁可多抓一些候选对也不要在第一阶段就漏掉。后续精筛阶段会进行更精确的过滤。Shuffle优化approxSimilarityJoin操作会产生巨大的Shuffle数据量。在EMR Serverless上我们通过以下方式优化调整spark.sql.shuffle.partitions默认是200对于PB级数据远远不够。我们将其设置为Executor数量 * 每个Executor核心数 * 3左右例如 (100 executors * 4 cores) * 3 1200以确保Shuffle阶段有足够的并行度避免单个Task处理数据过大。使用Kryo序列化在spark-defaults.conf中配置spark.serializer为org.apache.spark.serializer.KryoSerializer并注册相关类减少Shuffle过程中的数据体积和序列化开销。Executor内存与Off-Heap确保每个Executor有足够的内存我们用了32GB并合理配置spark.executor.memoryOverhead例如4GB以应对Shuffle过程中的内存峰值。3.3 精筛阶段语义向量相似度计算经过LSH粗筛我们得到了一个规模大幅减小的候选重复对列表duplicatedPairs。现在我们需要加载这些候选对对应的原始文本进行精确的语义比对。// 假设我们有一个预先计算好的语义向量表 embeddingDF包含id和vector列 // 这个表可能是另一个Spark作业预先批量生成并持久化在OSS上的。 // 将候选对与向量表进行连接 val candidatePairsWithVectors duplicatedPairs .join(broadcast(embeddingDF).as(emb1), col(id1) col(emb1.id)) .join(broadcast(embeddingDF).as(emb2), col(id2) col(emb2.id)) .select(col(id1), col(id2), col(emb1.vector).alias(vec1), col(emb2.vector).alias(vec2)) // 定义UDF计算余弦相似度 import org.apache.spark.ml.linalg.{Vector, Vectors} val cosineSimilarity udf((v1: Vector, v2: Vector) { val dotProduct v1.dot(v2) val norm1 Vectors.norm(v1, 2) val norm2 Vectors.norm(v2, 2) if (norm1 0 || norm2 0) 0.0 else dotProduct / (norm1 * norm2) }) // 计算相似度并过滤 val trueDuplicates candidatePairsWithVectors .withColumn(semantic_sim, cosineSimilarity(col(vec1), col(vec2))) .filter(col(semantic_sim) 0.9) // 设定语义相似度阈值 .select(col(id1), col(id2), col(semantic_sim))这个阶段的技巧广播变量BroadcastembeddingDF向量表可能很大PB级文本对应的向量也是TB级。直接进行join会导致巨大的Shuffle。这里的一个关键优化是我们只关心候选对id1和id2对应的向量。因此我们可以先将duplicatedPairs中的所有唯一ID收集到Driver端然后用这个ID列表去向量存储中按需提取对应的向量记录形成一个较小的向量子集smallEmbeddingDF再将其广播Broadcast到所有Executor。这样后续的join就变成了高效的Map-Side Join避免了Shuffle。代码中为了简洁直接用broadcast示意实际需要先进行过滤。向量计算优化余弦相似度计算是逐对进行的。如果单个桶内的候选对仍然很多比如几万对这个UDF计算会成为瓶颈。可以考虑使用更高效的数值计算库或者将向量转换为数组后利用Spark的SIMD优化如果环境支持。3.4 去重结果生成与输出最后我们得到了trueDuplicatesDataFrame它包含了所有语义相似的文本对。但去重要求的是为每组相似文本保留一个代表如ID最小的。这可以通过图计算中的连通分量Connected Components算法来解决。import org.apache.spark.graphx._ import org.apache.spark.rdd.RDD // 将重复对转换为图的边RDD val edgesRDD: RDD[Edge[Long]] trueDuplicates.rdd.map(row Edge(row.getAs[String](id1).toLong, row.getAs[String](id2).toLong, 1L) ) // 构建图并计算连通分量 val verticesRDD: RDD[(VertexId, String)] trueDuplicates.rdd.flatMap(row Seq((row.getAs[String](id1).toLong, row.getAs[String](id1)), (row.getAs[String](id2).toLong, row.getAs[String](id2))) ).distinct() val graph Graph(verticesRDD, edgesRDD) val cc graph.connectedComponents().vertices // 得到 (顶点ID, 连通分量最小ID) // 将结果关联回原始ID得到每个ID应该保留的代表ID val dedupMapping cc.toDF(id, representative_id) // 最终我们可以根据这个映射表选择每个代表ID对应的原始文本作为去重后的结果集写回OSS。 dedupMapping.write.mode(overwrite).parquet(oss://your-bucket/output/dedup-mapping/)4. 性能对比与“4倍加速”的达成整个方案最令人兴奋的部分无疑是性能的提升。我们与一个基线方案进行了对比。基线方案暴力比对法我们尝试在一个较小的数据集约100亿条文本上直接使用Spark进行语义向量两两余弦相似度计算。即使我们使用了各种优化如将向量标准化后转化为Breeze向量利用BLAS库并通过分区裁剪减少计算量该作业仍然因Shuffle数据量过大和计算复杂度太高而无法在合理时间内完成预估需要数周且成本极高。MinHash-LSH方案在相同的EMR Serverless Spark资源配置下处理完整的PB级数据集总耗时约为基线方案预估时间的1/4。这就是标题中“4倍加速”的由来。分解来看计算量锐减MinHash-LSH的approxSimilarityJoin阶段虽然也有Shuffle但其计算的是词集合的哈希和比较复杂度远低于高维向量的点积运算。更重要的是它将需要精细比对的数据范围从整个数据集的O(N²)级别缩小到了每个LSH桶内部的O(K²)级别其中K远小于N。网络IO减少暴力方案需要将每个向量都与其他所有向量进行匹配Shuffle数据量是N*(N-1)/2级别的。而我们的方案Shuffle主要发生在LSH分桶时数据量是N级别的精筛阶段由于使用了广播变量和Map-Side Join几乎不产生Shuffle。资源利用率提升Serverless Spark的弹性特性使得在LSH分桶这个Shuffle密集型阶段可以快速申请大量Executor进行并行分区计算而在精筛和连通分量计算阶段则可以适当减少资源避免浪费。这种按需伸缩进一步压缩了整体时间。具体的性能数据以一次实际运行为例数据规模约1.2PB原始文本经预处理和向量化后参与去重的文本条数约为8000亿条。EMR Serverless配置峰值使用了约500个Executor每个8核32GB内存Driver为16核64GB。总耗时约18小时。阶段耗时占比数据读取与MinHash签名生成~4小时。LSH近似连接approxSimilarityJoin~10小时Shuffle密集型。语义精筛与连通分量计算~3小时。结果输出~1小时。去重效果共识别出约150亿组语义重复对去重后数据量减少了约15%。这个成绩意味着我们成功地将一个原本被认为需要巨大集群运行数周甚至数月的任务压缩到了在一个晚上就能跑完的规模并且成本可控。5. 踩坑实录与核心调优经验任何大规模分布式项目的成功都离不开对坑的深刻理解和填坑的经验。这里分享几个让我们“掉头发”的关键问题和解决方案。5.1 Shuffle风暴与数据倾斜在最初的测试中approxSimilarityJoin阶段频繁出现Task失败部分Executor卡住。日志显示是Shuffle Fetch Failed和OOM。根因分析LSH分桶后数据分布可能极度不均匀。某些“热门”的哈希值例如停用词很多的通用文本会导致大量数据被分到同一个桶里从而对应到同一个Reduce Task。这个Task要处理的数据量远超其他Task导致内存溢出、GC时间过长最终失败。解决方案增加哈希表数量numHashTables这本质上是增加了分桶的维度。一个文本会出现在多个哈希桶中。虽然增加了存储但能有效打散热点因为一个文本在某个哈希函数下是热点在另一个函数下可能就不是。我们将numHashTables从3增加到5数据倾斜现象明显缓解。使用spark.sql.adaptive.enabledtrue自适应查询执行Spark AQE能自动处理数据倾斜通过将过大的分区拆分Split成多个小任务来处理。这在EMR Serverless Spark上是默认开启的但我们还需要调整相关参数如spark.sql.adaptive.advisoryPartitionSizeInBytes来更积极地触发拆分。自定义分区器作为终极手段我们可以尝试在LSH之后不使用默认的哈希分区而是根据桶的大小进行采样设计一个更均衡的定制化分区器。但这增加了复杂度我们优先尝试前两种方法后取得了满意效果。5.2 向量存储与广播的权衡精筛阶段需要加载向量。最初我们试图将整个TB级的向量表加载为一个DataFrame然后进行Join结果Driver在收集元数据阶段就OOM了。根因分析Spark在计划Join时需要知道参与Join的表的大小信息。对于非常大的表收集这些统计信息本身开销就很大。解决方案采用“先过滤再广播”的策略。从duplicatedPairs中提取出所有唯一的id这是一个相对较小的集合百万到千万级。将这个ID列表持久化到OSS的一个小文件中。启动一个单独的Spark作业或在一个作业内用mapPartitions读取这个ID列表然后去全量的向量存储可能是HBase、Redis或者同样是OSS上的Parquet文件中按Key查询出对应的向量。这个过程可以高度并行化。将查询结果id - vector映射收集起来作为一个小的本地Map或者存为一个小Parquet文件。在精筛作业中广播这个小的映射表或直接读取这个小文件。这样就完美避免了在Driver端操作超大向量表。5.3 MinHash签名长度的选择HashingTF的numFeatures参数特征哈希的维度和MinHashLSH的签名长度共同影响效果。坑最初为了节省内存将numFeatures设得较小如2^18。结果发现冲突率很高很多语义不相关的文本因为哈希冲突被分到了同一个桶导致精筛阶段计算量无效增加。经验numFeatures应该设置得足够大至少大于你预料中的最大词汇表大小。对于中文文本经过分词后不同的词可能达到百万级。我们最终设置为2^20约104万内存占用可以接受且冲突率显著降低。签名长度由numHashTables和每个签名的位数决定则需要在召回率和存储开销间平衡。5.4 EMR Serverless 资源配置的艺术Serverless不是无脑用资源配置需要精细调整。Driver内存必须给足特别是当需要收集部分数据如唯一ID列表或进行连通分量计算GraphX的一些算法会在Driver端进行协调时Driver内存不足会导致作业失败。我们从默认的4G逐步提升到16G最后稳定在32G。Executor的核内存比对于CPU密集型的特征提取和向量计算我们采用较高的核内存比如1:4即1核配4GB内存。对于Shuffle密集型的LSH连接阶段则需要更多内存来应对Shuffle缓冲区我们采用了1:8的配置。利用Spot实例降低成本EMR Serverless支持使用抢占式实例Spot。对于我们的作业除了Driver要求稳定外大部分Executor任务都可以容忍中断重启。我们将Executor的资源规格设置为使用Spot实例成本降低了60%-70%。作业的自动重试机制确保了少数实例被回收时任务能继续完成。这个项目让我深刻体会到处理超大规模数据问题算法上的巧妙设计如MinHash-LSH的降维和工程上的精细调优如Shuffle优化、资源配比同等重要。EMR Serverless Spark提供的弹性、托管能力让我们能更专注于业务逻辑本身而无需深陷集群运维的泥潭。最终实现的4倍加速不仅仅是算法的胜利更是云原生数据架构与分布式计算最佳实践共同作用的结果。