Spark MLlib KMeans执行报错:IllegalArgumentException求助
IllegalArgumentException: requirement failed in Spark MLlib Hey there! Let's break down the error you're facing and fix it step by step.
What's Causing the Error?
Looking at your error stack trace, the failure happens in MLUtils$.fastSquaredDistance at the require check. Digging into Spark's source code, this specific require ensures that two vectors being compared have identical dimensions. That means your trainingData RDD contains Vector instances with mismatched sizes—and KMeans can't handle that during training.
Why Are Vector Dimensions Mismatched?
Your dataset has 15-minute interval readings, so each day should have exactly 96 entries (24h * 4 intervals/hour). But it's likely that some dates are missing records, causing the collect_list(values) aggregation to return arrays with length != 96. When you convert these arrays to Vectors, you end up with vectors of varying dimensions—leading to the error when KMeans tries to compute distances between these vectors and the cluster centers.
Fixing the Issue
The solution is to filter out any groups where the collected values list doesn't have exactly 96 elements before converting to Vectors. Here's how to adjust your code:
// Original code val df2 = rd.groupBy("data").agg(collect_list("values")) // Add this line to filter invalid groups val df2Filtered = df2.filter(size($"collect_list(values)") === 96) // Continue with your existing logic using df2Filtered instead of df2 val convertUDF = udf((array : Seq[Double]) => { Vectors.dense(array.toArray) }) val withVector = df2Filtered.withColumn("collect_list(values)", convertUDF($"collect_list(values)"))
Bonus Checks & Improvements
Verify Vector Dimensions: To confirm the fix works, add a quick debug print before training:
trainingData.foreach(vec => println(s"Vector dimension: ${vec.size}"))All outputs should show
96if the filter works correctly.Use Modern Spark ML API: You're using the older
mllibAPI (RDD-based). Consider migrating to themlAPI (DataFrame-based), which has better built-in validation and is more user-friendly. For example:import org.apache.spark.ml.clustering.KMeans import org.apache.spark.ml.feature.VectorAssembler // Your preprocessing to get a DataFrame with a "features" column val kmeans = new KMeans().setK(4).setMaxIter(20) val model = kmeans.fit(trainingDataDF)Replace Deprecated Method:
unionAllis deprecated in newer Spark versions—useunioninstead for your DataFrame concatenation.
Note on Your "Local Machine Capacity" Suspicions
This error isn't related to memory or processing power (that would throw an OutOfMemoryError). It's a logical issue with inconsistent data dimensions. The fix above should resolve it even in local mode.
内容的提问来源于stack exchange,提问作者Lucas Peñalver

