Spark:将RDD[String]转换为DataFrame遇列合并异常
解决RDD[String]转DataFrame时拆分逗号分隔值为多列的问题
我来帮你搞定这个问题!你目前的代码把数组当成了单列,是因为toDF()默认会把Array[String]类型整体作为一个列处理,而不是拆分里面的元素。咱们分两种场景来解决:
场景1:已知列数和列名
如果你的每条数据拆分后列数固定,而且知道要给每列起什么名字,最简单的方法是把数组转换成元组,再转成DataFrame:
// 先把RDD[Array[String]]转成RDD[Tuple],这里假设每条数据拆分后有3个元素 val allNewData_tuple = allNewData_split.map(arr => (arr(0), arr(1), arr(2))) // 直接指定列名生成DataFrame val df_newData = allNewData_tuple.toDF("col1", "col2", "col3")
这样生成的DataFrame就会把每个数组元素对应到单独的列啦。
场景2:列数动态未知
如果不确定每条数据拆分后的列数,或者列数是动态变化的,就需要自定义Schema来实现:
import org.apache.spark.sql.types.{StringType, StructField, StructType} import org.apache.spark.sql.Row // 第一步:生成Schema。这里假设所有数据拆分后的长度一致,取第一条数据的长度来生成字段 val sampleArray = allNewData_split.take(1)(0) val schema = StructType( sampleArray.indices.map(index => StructField(s"column_${index + 1}", StringType, nullable = true) ) ) // 第二步:把RDD[Array[String]]转成RDD[Row] val allNewData_rows = allNewData_split.map(arr => Row.fromSeq(arr)) // 第三步:用自定义Schema创建DataFrame val df_newData = spark.createDataFrame(allNewData_rows, schema)
关键说明
你之前的代码之所以只生成单列,是因为Spark对RDD[Array[String]]调用toDF()时,会自动推断出一个ArrayType的单列(默认列名为value),而不会自动拆分数组元素。只有把数组转换成Tuple(固定列数)或者Row+自定义Schema(动态列数),才能让每个元素对应单独的列。
注意哦:如果你的数据里存在拆分后长度不一致的数组,上面的两种方法都会报错,所以建议先做数据校验,过滤掉长度不符合预期的记录~
内容的提问来源于stack exchange,提问作者diens
相关产品推荐
相关产品推荐

