You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.30 08:12:05