Nak*_*euh 1 scala apache-spark
我有两个数据帧df1,df2它们具有以下结构:
print(df1)
+-------+------------+-------------+---------+
| id| vector| start_time | end_time|
+-------+------------+-------------+---------+
| 1| [0,0,0,0,0]| 000| 200|
| 2| [1,1,1,1,1]| 200| 500|
| 3| [0,1,0,1,0]| 100| 500|
+-------+------------+-------------+---------+
print(df2)
+-------+------------+-------+
| id| vector| time|
+-------+------------+-------+
| A| [0,1,1,1,0]| 050|
| B| [1,0,0,1,1]| 150|
| C| [1,1,1,1,1]| 250|
| D| [1,0,1,0,1]| 350|
| E| [1,1,1,1,1]| 450|
| F| [1,0,5,0,0]| 550|
+-------+------------+-------+
Run Code Online (Sandbox Code Playgroud)
我要的是:对于每一个数据df1,得到的所有数据来自df2于该time之间start_time,并end_time与所有这些数据计算两个向量之间的欧氏距离。
我从以下代码开始,但我在计算距离的过程中遇到了困难:
val joined_DF = kafka_DF.crossJoin(
hdfs_DF.withColumnRenamed("id","id2").withColumnRenamed("vector","vector2")
)
.filter(col("time")>= col("start_time") &&
col("time")<= col("end_time"))
.withColumn("distance", ???) // Euclidean distance element-wise between columns vector and column vector2
Run Code Online (Sandbox Code Playgroud)
以下是示例数据的预期输出:
+-------+------------+-------------+---------+-------+------------+------+----------+
| id| vector| start_time | end_time| id2| vector2| time| distance |
+-------+------------+-------------+---------+-------+------------+------+----------+
| 1| [0,0,0,0,0]| 000| 200| A| [0,1,1,1,0]| 050| 1.73205|
| 1| [0,0,0,0,0]| 000| 200| B| [1,0,0,1,1]| 150| 1.73205|
| 2| [1,1,1,1,1]| 200| 500| C| [1,1,1,1,1]| 250| 0|
| 2| [1,1,1,1,1]| 200| 500| D| [1,0,1,0,1]| 350| 1.41421|
| 2| [1,1,1,1,1]| 200| 500| E| [1,1,1,1,1]| 450| 0|
| 3| [0,1,0,1,0]| 100| 500| B| [1,0,0,1,1]| 150| 1.73205|
| 3| [0,1,0,1,0]| 100| 500| C| [1,1,1,1,1]| 250| 1.73205|
| 3| [0,1,0,1,0]| 100| 500| D| [1,0,1,0,1]| 350| 2.23606|
| 3| [0,1,0,1,0]| 100| 500| E| [1,1,1,1,1]| 450| 1.73205|
+-------+------------+-------------+---------+-------+------------+------+----------+
Run Code Online (Sandbox Code Playgroud)
注意事项:
df1 总是会有少量数据,所以 crossJoin 不会冒着破坏我的记忆的风险。Audf应该在这种情况下工作。
import math._
import org.apache.spark.ml.linalg.Vector
import org.apache.spark.ml.linalg.Vectors
//input two vectors of length n, but must be equal length
//output euclidean distance between the vectors
val euclideanDistance = udf { (v1: Vector, v2: Vector) =>
sqrt(Vectors.sqdist(v1, v2))
}
Run Code Online (Sandbox Code Playgroud)
udf像这样使用你的新:
joined_DF.withColumn("distance", euclideanDistance($"vector", $"vector2"))
Run Code Online (Sandbox Code Playgroud)
| 归档时间: |
|
| 查看次数: |
2392 次 |
| 最近记录: |