【问题标题】:Columns value comparison in Spark data frameSpark数据框中的列值比较
【发布时间】:2018-09-02 02:06:51
【问题描述】:

我有一个包含大量记录的数据框。在该 DF 中,一条记录可以重复多次,并且每次更新时,最后更新的字段都将具有修改记录的日期。

我们有一组列,我们希望在这些列上比较相似 id 的行。在此比较期间,我们希望捕获从先前记录到当前记录的字段/列发生了哪些变化,并将其捕获到更新记录的“updated_columns”列中。将此第二条记录与第三条记录进行比较并识别更新的列并在第三条记录的“updated_columns”字段中捕获它,继续相同直到该 id 的最后一条记录,并对具有多个条目的每个 id 执行相同的操作.

最初,我们将列分组并从该组列中创建一个哈希值,并与下一行的哈希值进行比较,这样可以帮助我识别有更新的记录,但想要更新的列。

我在这里分享一些数据,这是预期的结果,这就是添加更新列后最终数据的样子(我可以说,使用列 Col1、Col2、Col3、col4 和 Col5 来比较两行):

希望以有效的方式做到这一点。任何人都尝试过这样的事情。

寻求帮助!

~克里什。

【问题讨论】:

  • 图形表示总是有帮助的!
  • 我已经用例子更新了这个问题,谢谢。

标签: apache-spark apache-spark-sql


【解决方案1】:

可以使用window

思路是按ID对数据进行分组,按LAST-UPDATED排序,将上一行(如果存在的话)的值复制到当前行然后将复制的数据与当前值进行比较。

val data = ... //the dataframe has the columns ID,Col1,Col2,Col3,Col4,Col5,LAST_UPDATED,IS_DELETED

val fieldNames = data.schema.fieldNames.dropRight(1) //1
val columns = fieldNames.map(f => col(f))
val windowspec = Window.partitionBy("ID").orderBy("LAST_UPDATED") //2
def compareArrayUdf() = ... //3

val result = data
  .withColumn("cur", array(columns: _*)) //4
  .withColumn("prev", lag($"cur", 1).over(windowspec)) //5
  .withColumn("updated_columns", compareArrayUdf()($"cur", $"prev")) //6
  .drop("cur", "prev") //7
  .orderBy("LAST_UPDATED")

备注:

  1. 创建所有要比较的字段的列表。使用除最后一个 (LAST-UPDATED) 之外的所有字段
  2. 创建一个按ID分区的窗口,每个分区按LAST-UPDATED排序
  3. 创建一个比较两个数组并将发现的差异映射到字段名称的 udf,代码见下文
  4. 创建一个包含所有应该比较的值的新列
  5. 创建一个新列,其中包含应该比较的上一个行的所有值(通过使用lag-函数)。上一行是具有相同 ID 且最大 LAST-UPDATED 小于当前行的行。该字段可以为空
  6. 比较两个新列并将结果放入updated-columns
  7. 删除在步骤 3 和 4 中创建的两个中间列

compareArraysUdf

def compareArray(cur: mutable.WrappedArray[String], prev: mutable.WrappedArray[String]): String = {
  if (prev == null || cur == null) return ""
  val res = new StringBuilder
  for (i <- cur.indices) {
    if (!cur(i).contentEquals(prev(i))) {
      if (res.nonEmpty) res.append(",")
      res.append(fieldNames(i))
    }
  }
  res.toString()
}
def compareArrayUdf() = udf[String, mutable.WrappedArray[String], mutable.WrappedArray[String]](compareArray)

【讨论】:

  • 即使我们有一个非常相似的用例,但是我们的容量非常大,有 1.5 亿条记录,当我执行 Window.partitionBy("ID") 时,它会创建太多分区并且失败了。有没有办法绕过这个,或者我做错了什么?期待一些帮助!
【解决方案2】:

您可以将 DataFrame 或 DataSet 连接到自身,连接两行中 id 相同的行,左行的版本为i,右行的版本为i+1。这是一个例子

case class T(id: String, version: Int, data: String)

val data = Seq(T("1", 1, "d1-1"), T("1", 2, "d1-2"), T("2", 1, "d2-1"), T("2", 2, "d2-2"), T("2", 3, "d2-3"), T("3", 1, "d3-1"))
data: Seq[T] = List(T(1,1,d1-1), T(1,2,d1-2), T(2,1,d2-1), T(2,2,d2-2), T(2,3,d2-3), T(3,1,d3-1))

val ds = data.toDS

val joined = ds.as("ds1").join(ds.as("ds2"), $"ds1.id" === $"ds2.id" && (($"ds1.version"+1) === $"ds2.version"))

然后您可以引用新 DataFrame/DataSet 中的列,例如 $"ds1.data$"ds2.data 等。

要查找数据从一个版本更改为另一个版本的行,您可以这样做

joined.filter($"ds1.data" !== $"ds2.data")

【讨论】:

  • 这是一个自连接,但是有了这个,我们如何才能知道值从前一个记录更改为当前记录的列的列表?我在这里有什么遗漏吗?
  • 我更新了答案并展示了如何将一行的值与另一行的值进行比较的示例。希望这能给你一个想法
  • 感谢您分享您的观点,我已经解决了这个问题,但这需要逐列比较,然后我可以捕获更新的列并填充更新的列。我正在寻找一些有效的方法来比较字段并派生更新数据的字段名称,然后我可以将其填充到更新的列中。
猜你喜欢
  • 1970-01-01
  • 2022-01-20
  • 1970-01-01
  • 1970-01-01
  • 2022-08-18
  • 2015-03-23
  • 2017-03-08
  • 1970-01-01
  • 2020-08-13
相关资源
最近更新 更多