【问题标题】:Spark Join for Each Item in List列表中每个项目的 Spark Join
【发布时间】:2022-09-22 00:27:09
【问题描述】:

我有一个 Spark 数据集

+----------+-------+----+---+--------------+
|        _1|     _2|  _3| _4|            _5|
+----------+-------+----+---+--------------+
|      null|1111111|null| 15|       [98765]|
|      null|2222222|null| 16|[97008, 98765]|
|6436334664|3333333|null| 15|       [97008]|
|2356242642|4444444|null| 11|       [97008]|
+----------+-------+----+---+--------------+

其中第五列是与该行关联的邮政编码列表。我有另一个表,其中每个邮政编码都有唯一的行以及相应的经度和纬度。我想创建一个像

+----------+-------+----+---+--------------+-----------------------------------
|        _1|     _2|  _3| _4|            _5|                                _6|
+----------+-------+----+---+--------------+----------------------------------+
|3572893528|1111111|null| 15|       [98765]| [(54.12,-80.53)]                 |
|5325232523|2222222|null| 16|[98765, 97008]| [(54.12,-80.53), (44.12,-75.11)] |
|6436334664|3333333|null| 15|       [97008]| [(54.12,-80.53)]                 | 
|2356242642|4444444|null| 11|       [97008]| [(54.12,-80.53)]                 |
+----------+-------+----+---+--------------+----------------------------------+

其中第六列是第五列序列中拉链的坐标。

每次我需要坐标时,我都尝试过滤邮政编码表,但我得到了一个 NPE,我认为是因为this 问题中详述的类似原因。如果我在过滤之前尝试收集邮政编码表,我的内存就会用完。

我正在使用 Scala,并且在 Spark 作业中使用 Spark SQL 获得了原始数据集。任何解决方案将不胜感激,谢谢。

  • 您的示例是否有点误导或这是您真正想要的?因为您将98765(54.12,-80.53)(44.12,-75.11) 相关联 - 前两行?它必须是一对一的映射对吗?这意味着98765(54.12,-80.53)97008(44.12,-75.11) 相关?
  • @vilalabinot感谢您的澄清,这就是我的意思,映射是1对1。我已经更新了问题

标签: scala apache-spark join apache-spark-sql


【解决方案1】:

假设(对您问题的评论成立,并且)我们有两个数据集(简化您的示例),分别为 dsds2

+---+--------------+
|_1 |_2            |
+---+--------------+
|15 |[98765]       |
|16 |[97008, 98765]|
|15 |[97008]       |
|15 |[97008]       |
+---+--------------+
+-----+---------------+
|_2   |_3             |
+-----+---------------+
|98765|{54.12, -80.53}|
|97008|{44.12, -75.11}|
+-----+---------------+

想法是创建一个唯一 ID(以便我们稍后加入),explode 数据集,然后join 以获取每个唯一 ID 的坐标,最后再次加入表。

创建唯一 ID:

ds = ds.withColumn("id", monotonically_increasing_id())

然后创建包含id 和您的邮政编码的映射表:

val map = ds
  .withColumn("_2", explode(col("_2")))
  .join(ds2, Seq("_2"), "left")
  .groupBy("id").agg(collect_set(col("_3")))

最后加入主表:

ds = ds.join(map, Seq("id"))

最终输出:

+---+--------------+----------------------------------+
|_1 |_2            |collect_set(_3)                   |
+---+--------------+----------------------------------+
|15 |[98765]       |[{54.12, -80.53}]                 |
|16 |[97008, 98765]|[{54.12, -80.53}, {44.12, -75.11}]|
|15 |[97008]       |[{44.12, -75.11}]                 |
|15 |[97008]       |[{44.12, -75.11}]                 |
+---+--------------+----------------------------------+

祝你好运!

【讨论】:

  • 此方法效果很好,但邮政编码与坐标的顺序不匹配。
  • 我害怕这种情况,让我尝试解决它
  • 我不认为你可以做很多事情,除了保存key 本身,如:ds2 = ds2.withColumn("_3", struct("_2", "_3")),那么你收集的集合看起来像:[{98765, {54.12, -80.53}}, {97008, {44.12, -75.11}}]
猜你喜欢
  • 2022-11-07
  • 2021-12-22
  • 2014-03-21
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-07-03
  • 1970-01-01
相关资源
最近更新 更多