Spark读取可变列数整数矩阵CSV后,列类型转换无效的问题求助
解决Spark CSV列批量转换为Integer类型的问题
我明白你遇到的问题了——Spark DataFrame是不可变对象,你原来的foreach循环里虽然调用了withColumn、drop这些方法,但这些操作只会生成新的DataFrame,而你没有把这些新的实例赋值回原来的df变量,所以原DataFrame的Schema完全没变化。
下面给你两种可行的解决方案:
方法一:使用foldLeft累积转换(推荐)
foldLeft可以帮你遍历所有列,逐步更新DataFrame,每一步都把转换后的新DataFrame传递下去,最后得到所有列都是Integer类型的结果:
import org.apache.spark.sql.types.IntegerType // 初始读取CSV var df = spark.read.csv(originalPath) // 用foldLeft遍历所有列,批量转换类型 df = df.columns.foldLeft(df) { (currentDf, colName) => currentDf.withColumn(colName, currentDf.col(colName).cast(IntegerType)) } // 现在打印Schema就能看到所有列都是Integer类型了 df.printSchema()
这个方法的核心是:每次转换列类型后,都把新的DataFrame赋值给currentDf,最终覆盖原来的df变量,这样所有列的类型转换都会生效。
方法二:读取时直接指定动态Schema(适合带表头的CSV)
如果你的CSV文件包含表头行,还可以先读取表头,动态构建包含IntegerType的Schema,再读取文件,这样一步到位:
import org.apache.spark.sql.types.{StructType, StructField, IntegerType} // 第一步:读取表头行 val headerLine = spark.read.text(originalPath).first().getString(0) val columnNames = headerLine.split(",") // 第二步:构建动态Schema,所有列都设为IntegerType val schemaFields = columnNames.map(name => StructField(name, IntegerType, nullable = true)).toList val dynamicSchema = StructType(schemaFields) // 第三步:用指定的Schema读取CSV val df = spark.read.schema(dynamicSchema).csv(originalPath)
为什么你的原代码不生效?
再强调一下:Spark的DataFrame是不可变的,任何转换操作(比如withColumn、drop)都不会修改原对象,而是返回一个全新的DataFrame实例。你原来的foreach循环里,每次操作都生成了新的DataFrame,但没有把它赋值给df,所以原df始终是最初读取的String类型版本。
内容的提问来源于stack exchange,提问作者Hoori M.
相关产品推荐
相关产品推荐

