sparklyr使用spark_apply对全列应用自定义标准化函数报错如何解决
报错根因说明
之前的运行报错和输出为空的核心原因有两个:
spark_apply默认按分区处理数据,你在函数内部计算的mean和sd是分区内的统计量,不是全局值,不符合标准化的业务要求;- 测试用的样本量只有4行,数据可能被分配到空分区,且返回值格式不符合
spark_apply的要求,导致报错或无输出。
推荐优先使用Spark原生算子实现,完全避开spark_apply的性能和稳定性问题,适配100+列的批量处理场景:
方案1:sparklyr原生dplyr风格实现(最适配你的原有开发习惯)
library(sparklyr) library(dplyr) library(rlang) # 连接Spark、导入数据 sc <- spark_connect(method = "databricks") iris_sdf <- sdf_copy_to(sc, iris, overwrite = T) # 1. 筛选需要标准化的数值列,排除非数值列如Species num_cols <- setdiff(colnames(iris_sdf), "Species") # 2. 计算所有数值列的全局均值、标准差,拉取到Driver端 col_stats <- iris_sdf %>% summarise(across(all_of(num_cols), list(mean = mean, sd = sd), na.rm = TRUE)) %>% collect() # 3. 批量对所有数值列执行标准化 normalized_sdf <- iris_sdf %>% mutate( across( all_of(num_cols), ~ (. - !!sym(paste0(cur_column(), "_mean"))) / !!sym(paste0(cur_column(), "_sd")), # 如果需要直接覆盖原列,删除下面的.names参数即可 .names = "{.col}_norm" ) ) # 查看结果 head(normalized_sdf, 10)
方案2:sparklyr调用MLlib StandardScaler实现
用Spark官方内置的标准化算子,代码更简洁:
library(sparklyr) library(dplyr) sc <- spark_connect(method = "databricks") iris_sdf <- sdf_copy_to(sc, iris, overwrite = T) num_cols <- setdiff(colnames(iris_sdf), "Species") normalized_sdf <- iris_sdf %>% # 把数值列组装为特征向量 ft_vector_assembler(input_cols = num_cols, output_col = "raw_features") %>% # 执行标准化,with_mean=TRUE表示去中心化,with_std=TRUE表示除以标准差 ml_standard_scaler( input_col = "raw_features", output_col = "scaled_features", with_mean = TRUE, with_std = TRUE ) %>% # 把标准化后的向量拆回独立列 ft_vector_to_array(input_col = "scaled_features", output_col = "scaled_arr") %>% mutate(across(all_of(num_cols), ~ scaled_arr[which(num_cols == cur_column())])) %>% # 删除中间过程列 select(-raw_features, -scaled_features, -scaled_arr)
方案3:PySpark实现
from pyspark.sql import SparkSession from pyspark.ml.feature import VectorAssembler, StandardScaler from pyspark.sql.functions import element_at # 初始化Spark spark = SparkSession.builder.appName("col_normalize").getOrCreate() iris_df = spark.createDataFrame(iris) num_cols = [c for c in iris_df.columns if c != "Species"] # 组装特征向量 assembler = VectorAssembler(inputCols=num_cols, outputCol="raw_features") assembled_df = assembler.transform(iris_df) # 训练标准化模型并转换 scaler = StandardScaler(inputCol="raw_features", outputCol="scaled_features", withMean=True, withStd=True) scaler_model = scaler.fit(assembled_df) scaled_df = scaler_model.transform(assembled_df) # 把向量拆为独立列 for idx, col_name in enumerate(num_cols): scaled_df = scaled_df.withColumn(f"{col_name}_norm", element_at("scaled_features", idx + 1)) # 删除中间列 result_df = scaled_df.drop("raw_features", "scaled_features") result_df.show(10)
方案4:Scala实现
逻辑和PySpark完全一致,调用MLlib StandardScaler即可:
import org.apache.spark.ml.feature.{StandardScaler, VectorAssembler} import org.apache.spark.sql.functions.element_at val spark = SparkSession.builder().appName("col_normalize").getOrCreate() import spark.implicits._ // 假设已加载DataFrame为irisDf val numCols = irisDf.columns.filter(_ != "Species") val assembler = new VectorAssembler().setInputCols(numCols).setOutputCol("rawFeatures") val assembledDf = assembler.transform(irisDf) val scaler = new StandardScaler().setInputCol("rawFeatures").setOutputCol("scaledFeatures") .setWithMean(true).setWithStd(true) val scalerModel = scaler.fit(assembledDf) val scaledDf = scalerModel.transform(assembledDf) // 拆分向量为独立列 val resultDf = numCols.zipWithIndex.foldLeft(scaledDf) { (df, colInfo) => val (colName, idx) = colInfo df.withColumn(s"${colName}_norm", element_at($"scaledFeatures", idx + 1)) }.drop("rawFeatures", "scaledFeatures") resultDf.show(10)
内容的提问来源于stack exchange,提问作者Cyrus Mohammadian
相关产品推荐
相关产品推荐

