Spark Scala中按cuid聚合并求和Sparse Vector及dt_diff字段
Hey there! Let's work through this aggregation problem where we need to group data by cuid, sum up the dt_diff integers, and also combine the SparseVector features by summing their corresponding indices, values, and total sizes (as shown in your example).
Step-by-Step Solution
Since Spark doesn't have a built-in function for this specific SparseVector combination, we'll create a custom UDF (User-Defined Function) to handle the vector merging, then use standard aggregation functions for the rest.
1. Import Required Libraries
First, make sure you have the necessary imports for Spark's SparseVector and SQL functions:
import org.apache.spark.ml.linalg.SparseVector import org.apache.spark.sql.functions.{udf, sum, collect_list}
2. Define the Custom SparseVector Merging UDF
This UDF takes a list of SparseVectors for a single cuid, sums their total sizes, and adds corresponding indices and values together (matching your expected output):
val mergeMatchingSparseVectors = udf((vectors: Seq[SparseVector]) => { if (vectors.isEmpty) null else { // Calculate total size by summing all vector sizes val totalSize = vectors.map(_.size).sum // Transpose index arrays to sum corresponding positions val summedIndices = vectors.map(_.indices).transpose.map(_.sum).toArray // Transpose value arrays to sum corresponding positions val summedValues = vectors.map(_.values).transpose.map(_.sum).toArray new SparseVector(totalSize, summedIndices, summedValues) } })
3. Perform the Aggregation
Now group the data by cuid, sum dt_diff, and apply our UDF to aggregate the features:
val aggregatedDF = originalDF .groupBy("cuid") .agg( sum("dt_diff").alias("dt_diff"), mergeMatchingSparseVectors(collect_list("features")).alias("features") )
Input vs. Expected Output
Let's verify with your sample data:
Input Data
| cuid | features | dt_diff |
|---|---|---|
| 12654467 | (2013492,[1743933,2013491],[2.0,2.0]) | 4 |
| 12654467 | (1876451,[1000000,1876451],[5.0,7.0]) | 10 |
Aggregated Output
| cuid | features | dt_diff |
|---|---|---|
| 12654467 | (3889943,[2743933,3889942],[7.0,9.0]) | 14 |
This matches exactly what you're looking for: dt_diff is summed to 14, and the features vector has its size, indices, and values all summed from the original two vectors.
内容的提问来源于stack exchange,提问作者Jerdy

