use*_*267 10 json scala apache-spark apache-spark-sql
我有两个数据框,我正试图找出它们之间的区别。2个数据帧包含struct数组。我不需要该结构中的1个键。因此,我首先将其删除,然后转换为JSON字符串。比较时,我需要知道该数组(Json)中更改了多少个元素。有办法做到这一点吗?
双方base_data_set并target_data_set包含ID和KEY。KEY是一个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。有人可以建议我如何实现这一目标吗?
我不确定这是否是您的意思,因为问题的某些部分不容易理解(至少对我而言)。
我使用了两个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的行所示。
| 归档时间: |
|
| 查看次数: |
193 次 |
| 最近记录: |