Spark DataFrame中基于另一DataFrame验证city_id一致性的实现方案
Spark DataFrame数组列关联验证方案
数据样例
df1
+--------------------+---------------------------------------------------------+ |addressId |sorted_city_ids | +--------------------+---------------------------------------------------------+ |AkhRMTbGbPiUnrBdSlov|[732916, 734241] | |AkhTMKHi9Ui7DHcspbfg|[724985, 725983, 725603, 728894, 728896, 728943, 729422] | |AkhflxpBdi1tTZmmTf1w|[729103, 731732, 731736, 731738] | +--------------------+---------------------------------------------------------+
df2
+--------------------+---------+------------------------------+-------------------+ |addressId |city_id |city_master_id |dateInserted | +--------------------+---------+------------------------------+-------------------+ |AkhRMTbGbPiUnrBdSlov|732916 |ncLsR29SYDIeyc9aZp1cYojGvSkkWf|2021-11-04 05:30:00| |AkhTMKHi9Ui7DHcspbfg|9361852 |mt9WV7omxa8nxD1n7yOGFDtEOTWPdq|2021-10-19 05:30:00| |AkhflxpBdi1tTZmmTf1w|729103 |kZiy4ZeULOqqdO8yKDmsDA13RdegC2|2022-11-04 20:50:57| |AkhflxpBdi1tTZmmTf1w|731732 |msEfdZemxa8nxD1n7yOGFDtEOTWPdq|2022-11-04 20:50:57| |AkhflxpBdi1tTZmmTf1w|731736 |kZiy4ZeULOqqdO8yKDmsDA13RdegC2|2022-11-04 20:50:57| |AkhflxpBdi1tTZmmTf1w|731738 |kZiy4ZeULOqqdO8yKDmsDA13RdegC2|2022-11-05 20:50:57| +--------------------+---------+------------------------------+-------------------+
需求说明
处理df1每一行,基于sorted_city_ids数组列,从df2中获取每个city_id对应的city_master_id和dateInserted的日期部分,验证该行所有city_id是否与数组第一个city_id的city_master_id和日期一致,筛选出不匹配的city_id。
示例预期
以df1第三行为例,需将731732、731736、731738分别与第一个city_id(729103)对比,最终筛选出不匹配的:
731732 731738
- 731732因
city_master_id不匹配被筛选 - 731738因日期不匹配被筛选
遇到的问题
尝试用collect.toMap()创建cityId -> List(city_master_id, date)的映射,但数据量过大导致Driver堆内存不足,需要高效的分布式处理方案。
解决方案
步骤1:预处理df2,提取日期部分
先对df2做预处理,提取dateInserted的日期部分,简化后续对比逻辑:
import org.apache.spark.sql.functions._ val df2_processed = df2.withColumn("date_part", to_date(col("dateInserted")))
步骤2:拆分df1数组并关联基准信息
将df1的sorted_city_ids数组拆分为单个city_id,提取数组第一个元素作为基准base_city_id,再通过addressId关联df2中基准城市的city_master_id和日期:
// 拆分数组并保留addressId,提取基准city_id val df1_exploded = df1 .withColumn("base_city_id", element_at(col("sorted_city_ids"), 1)) .withColumn("city_id", explode(col("sorted_city_ids"))) // 关联基准城市的对比参考值 val df1_with_base = df1_exploded.join( df2_processed.select( col("addressId"), col("city_id").alias("base_city_id"), col("city_master_id").alias("base_master_id"), col("date_part").alias("base_date") ), Seq("addressId", "base_city_id"), "left" )
步骤3:关联当前城市信息并筛选不匹配项
将结果与df2_processed关联,获取每个city_id的实际信息,对比基准值后筛选出不匹配的记录:
val result = df1_with_base.join( df2_processed.select(col("addressId"), col("city_id"), col("city_master_id"), col("date_part")), Seq("addressId", "city_id"), "left" ) // 排除基准自身,筛选master_id或日期不匹配的city_id .filter(col("city_id") =!= col("base_city_id")) .filter( col("city_master_id") =!= col("base_master_id") || col("date_part") =!= col("base_date") ) .select("city_id")
方案优势
- 完全基于Spark分布式计算,避免将大量数据拉取到Driver端,彻底解决内存溢出问题
- 所有操作在Executor节点并行处理,适配大数据量场景
- 分步骤关联基准与对比数据,逻辑清晰且性能可控
内容的提问来源于stack exchange,提问作者Abhinay
相关产品推荐
相关产品推荐

