Spark如何获取Scala中两个JSONS中更改的键数?

use*_*267 10 json scala apache-spark apache-spark-sql

我有两个数据框,我正试图找出它们之间的区别。2个数据帧包含struct数组。我不需要该结构中的1个键。因此,我首先将其删除,然后转换为JSON字符串。比较时,我需要知道该数组(Json)中更改了多少个元素。有办法做到这一点吗?

双方base_data_settarget_data_set包含IDKEYKEY是一个array<Struct>

root
 |-- id: string (nullable = true)
 |-- result: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- key1: integer (nullable = true)
 |    |    |-- key3: string (nullable = false)
 |    |    |-- key2: string (nullable = true)
 |    |    |-- key4: string (nullable = true)

val temp_base = base_data_set
    .withColumn("base_result", explode(base_data_set(RESULT)))
    .withColumn("base",
        struct($"base_result.key1", $"base_result.key2", $"base_result.key3"))
    .groupBy(ID)
    .agg(to_json(collect_list("base")).as("base_picks"))

val temp_target = target_data_set
    .withColumn("target_result", explode(target_data_set(RESULT)))
    .withColumn("target",
        struct($"target_result.key1", $"target_result.key2", $"target_result.key3"))
    .groupBy(ID)
    .agg(to_json(collect_list("target")).as("target_picks"))    


val common_keys = temp_base
    .join(temp_target, temp_base(ID) ===  temp_target(ID))
    .drop(temp_target(ID))
    .withColumn("isModified", $"base_picks" =!= $"target_picks") 
Run Code Online (Sandbox Code Playgroud)

即使有1个项目更改,它也会返回false,但只有在更改了n个(例如n = 3)元素(在数组中)时,我才需要返回false。有人可以建议我如何实现这一目标吗?

moe*_*moe 5

我不确定这是否是您的意思,因为问题的某些部分不容易理解(至少对我而言)。

我使用了两个json文件来模拟您的架构。他们看起来像这样:

base_data_set:
{ "id": 1,  "result": [ {"key1":  23, "key2": "qwerty", "key3":  "abc"}, {"key1":  24, "key2": "asdf", "key3":  "abc"},  {"key1":  25, "key2": "xcv", "key3":  "abc"}]}
{ "id": 2,  "result": [ {"key1":  23, "key2": "qwerty", "key3":  "abc"}, {"key1":  24, "key2": "asdf", "key3":  "abc"},  {"key1":  25, "key2": "xcv", "key3":  "abc"}]}
{ "id": 3,  "result": [ {"key1":  "1", "key2": "2", "key3":  "3"}, {"key1":  "4", "key2": "5", "key3":  "6"},  {"key1":  "7", "key2": "8", "key3":  "9"}]}
{ "id": 4,  "result": [ {"key1":  "4", "key2": "5", "key3":  "6"}, {"key1":  "1", "key2": "2", "key3":  "3"},   {"key1":  "7", "key2": "8", "key3":  "9"}]}

target_data_set:
{ "id": 1,  "result": [ {"key1":  24, "key2": "qwerty", "key3":  "abc"}, {"key1":  24, "key2": "asdf", "key3":  "abc"},  {"key1":  25, "key2": "xcv", "key3":  "abc"}]}
{ "id": 2,  "result": [ {"key1":  23, "key2": "qwertu", "key3":  "abc"}, {"key1":  24, "key2": "asdfg", "key3":  "abc"},  {"key1":  25, "key2": "xcvv", "key3":  "abc"}]}
{ "id": 3,  "result": [ {"key1":  "1", "key2": "2", "key3":  "3"}, {"key1":  "4", "key2": "5", "key3":  "6"},  {"key1":  "7", "key2": "8", "key3":  "9"}]}
{ "id": 4,  "result": [ {"key1":  "1", "key2": "2", "key3":  "3"}, {"key1":  "4", "key2": "5", "key3":  "6"},  {"key1":  "7", "key2": "8", "key3":  "9"}]}
Run Code Online (Sandbox Code Playgroud)

如您所见,第一行仅在结果数组中的一个结构中有所不同,而第二行中的所有结构均不同。第3行和第4行显示了一种情况,如果您认为这是更改,我不清楚。两个表之间的结构相同,但是它们的顺序在第4行中有所变化。

从您的初始转换开始,我删除了to_json函数,因为它将结构化元素转换为字符串,这使得比较变得困难:

val temp_base = base_data_set
  .withColumn("base_result", explode(base_data_set("result")))
  .withColumn("base",
    struct($"base_result.key1", $"base_result.key2", $"base_result.key3"))
  .groupBy("id")
  .agg(collect_list("base").as("base_picks"))


val temp_target = target_data_set
  .withColumn("target_result", explode(target_data_set(RESULT)))
  .withColumn("target",
    struct($"target_result.key1", $"target_result.key2", $"target_result.key3"))
  .groupBy(ID)
  .agg(collect_list("target").as("target_picks"))


val common_keys = temp_base
  .join(temp_target, temp_base(ID) ===  temp_target(ID))
  .drop(temp_target(ID))
  .withColumn("isModified", $"base_picks" =!= $"target_picks")
Run Code Online (Sandbox Code Playgroud)

之后,您可以使用用户定义的函数比较的结果collect_list。它采用两列的内容,并计算有多少不同的元素:

  val numChangedStruct = udf {
  (left: mutable.WrappedArray[Object], right: mutable.WrappedArray[Object]) =>
    left.zip(right).count(x => !x._1.equals(x._2))
}
Run Code Online (Sandbox Code Playgroud)

并应用:

common_keys.withColumn("numChangedStruct", numChangedStruct($"base_picks", $"target_picks")).show(20, false)

+---+----------------------------------------------+------------------------------------------------+----------+----------------+
|id |base_picks                                    |target_picks                                    |isModified|numChangedStruct|
+---+----------------------------------------------+------------------------------------------------+----------+----------------+
|1  |[[23,qwerty,abc], [24,asdf,abc], [25,xcv,abc]]|[[24,qwerty,abc], [24,asdf,abc], [25,xcv,abc]]  |true      |1               |
|3  |[[1,2,3], [4,5,6], [7,8,9]]                   |[[1,2,3], [4,5,6], [7,8,9]]                     |false     |0               |
|2  |[[23,qwerty,abc], [24,asdf,abc], [25,xcv,abc]]|[[23,qwertu,abc], [24,asdfg,abc], [25,xcvv,abc]]|true      |3               |
|4  |[[4,5,6], [1,2,3], [7,8,9]]                   |[[1,2,3], [4,5,6], [7,8,9]]                     |true      |2               |
+---+----------------------------------------------+------------------------------------------------+----------+----------------+
Run Code Online (Sandbox Code Playgroud)

但是,此解决方案取决于“结果”中元素的顺序,如ID 3和4的行所示。