【发布时间】:2020-05-17 15:41:05
【问题描述】:
基本上,我想在数据框的每一行上应用我的函数 countSimilarColumns 并将结果放在一个新列中。
我的代码如下
def main(args: Array[String]) = {
val customerID = "customer-1" //args(0)
val rawData = readFromResource("json", "/spark-test-data-copy.json")
val flattenData = rawData.select(flattenSchema(rawData.schema): _*)
val referenceCustomerRow = flattenData.transform(getCustomer(customerID)).first
}
def getCustomer(customerID: String)(dataFrame: DataFrame) = {
dataFrame.filter($"customer" === customerID)
}
def countSimilarColumns(first: Row, second: Row): Int = {
if (!(first.getAs[String]("customer").equals(second.getAs[String]("customer"))))
first.toSeq.zip(second.toSeq).count { case (x, y) => x == y }
else
-1
}
我想做如下的事情。但我不知道该怎么做。
flattenData
.withColumn(
"similarity_score",
flattenData.map(row => countSimilarColumns(row, referenceCustomerRow))
)
.show()
扁平化示例数据:
{"customer":"customer-1","att-a":"7","att-b":"3","att-c":"10","att-d":"10"}
{"customer":"customer-2","att-a":"9","att-b":"7","att-c":"12","att-d":"4"}
{"customer":"customer-3","att-a":"7","att-b":"3","att-c":"1","att-d":"10"}
{"customer":"customer-4","att-a":"9","att-b":"14","att-c":"10","att-d":"4"}
想要的输出:
+--------------------+-----------+
| customer | similarity_score |
+--------------------+-----------+
|customer-1 | -1 |
|customer-2 | 0 |
|customer-3 | 3 |
|customer-4 | 1 |
UDF 是唯一的方法吗?如果是,那么我想保持我的函数 countSimilarColumns 原样,这样它就可以测试了。怎么可能? 我是 Spark/Scala 的新手。
【问题讨论】:
-
你是如何分配相似度分数的,你有什么逻辑吗?假设如果客户 3 重复了 10 次,你是不是加了 10 作为相似度分数?
-
我在贴的时候已经添加了countSimilarColumns这个函数。它导致该行的相似度得分。该函数检查两列的值是否相同然后计数++。我只需要在 DF 中添加一个新列。
-
如果有错误请纠正我,假设客户列的值 customer-3 5 次,这意味着客户 3 - 5 的计数.. 是否正确?
-
错了!你所说的是一个简单的 count() 数据框。我比较两行而不是列。并且一行可以有 x 个列,因此比较每个索引列,然后计算它们是否相同。
-
简单来说,我想将我的函数 countSimilarColumns 应用于数据帧的每一行并将结果放入新列中
标签: scala apache-spark apache-spark-sql