【问题标题】:Spark doesn't read columns with null values in first rowSpark 不读取第一行中具有空值的列
【发布时间】:2018-01-18 05:43:07
【问题描述】:

以下是我的 csv 文件中的内容:

A1,B1,C1
A2,B2,C2,D1
A3,B3,C3,D2,E1
A4,B4,C4,D3
A5,B5,C5,,E2

所以,有 5 列,但第一行只有 3 个值。

我使用以下命令阅读它:

val csvDF : DataFrame = spark.read
.option("header", "false")
.option("delimiter", ",")
.option("inferSchema", "false")
.csv("file.csv") 

以下是我使用 csvDF.show() 得到的结果

+---+---+---+
|_c0|_c1|_c2|
+---+---+---+
| A1| B1| C1|
| A2| B2| C2|
| A3| B3| C3|
| A4| B4| C4|
| A5| B5| C5|
+---+---+---+

如何读取所有列中的所有数据?

【问题讨论】:

  • 是否可以将所有 5 列添加到每一行?就像第 1 行而不是 A1,B1,C1 一样,它是 A1,B1,C1,,
  • 这只是一种解决方法,如果 csv 由其他人管理,则将无法使用。
  • 手动指定架构
  • 如果我们不知道schema,csv中的内容事先不知道怎么办。
  • csv中的所有内容都可以指定为StringType

标签: apache-spark apache-spark-sql apache-spark-2.0 spark-csv


【解决方案1】:

基本上,您的 csv 文件格式不正确,因为每行中的列数不相等,如果您想使用 spark.read.csv 读取它,则需要这样做。但是,您可以改为使用 spark.read.textFile 读取它,然后解析每一行。

据我了解,您事先并不知道列数,因此您希望您的代码能够处理任意数量的列。为此,您需要确定数据集中的最大列数,因此您需要对数据集进行两次遍历。

对于这个特定的问题,我实际上会使用 RDD 而不是 DataFrames 或 Datasets,如下所示:

val data  = spark.read.textFile("file.csv").rdd

val rdd = data.map(s => (s, s.split(",").length)).cache
val maxColumns = rdd.map(_._2).max()

val x = rdd
  .map(row => {
    val rowData = row._1.split(",")
    val extraColumns = Array.ofDim[String](maxColumns - rowData.length)
    Row((rowData ++ extraColumns).toList:_*)
  })

希望有帮助:)

【讨论】:

    【解决方案2】:

    您可以将其读取为只有一列的数据集(例如使用另一个分隔符):

    var df = spark.read.format("csv").option("delimiter",";").load("test.csv")
    df.show()
    
    +--------------+
    |           _c0|
    +--------------+
    |      A1,B1,C1|
    |   A2,B2,C2,D1|
    |A3,B3,C3,D2,E1|
    |   A4,B4,C4,D3|
    |  A5,B5,C5,,E2|
    +--------------+
    

    然后您可以使用this answer 手动将您的列一分为五,这将在元素不存在时添加 null 值:

    var csvDF = df.withColumn("_tmp",split($"_c0",",")).select(
        $"_tmp".getItem(0).as("col1"),
        $"_tmp".getItem(1).as("col2"),
        $"_tmp".getItem(2).as("col3"),
        $"_tmp".getItem(3).as("col4"),
        $"_tmp".getItem(4).as("col5")
    )
    csvDF.show()
    
    +----+----+----+----+----+
    |col1|col2|col3|col4|col5|
    +----+----+----+----+----+
    |  A1|  B1|  C1|null|null|
    |  A2|  B2|  C2|  D1|null|
    |  A3|  B3|  C3|  D2|  E1|
    |  A4|  B4|  C4|  D3|null|
    |  A5|  B5|  C5|    |  E2|
    +----+----+----+----+----+
    

    【讨论】:

      【解决方案3】:

      如果dataTypes 列和列数已知,那么您可以定义schema 并在将csv 文件读取为dataframe 时应用schema。下面我将所有五列定义为stringType

      val schema = StructType(Seq(
        StructField("col1", StringType, true),
        StructField("col2", StringType, true),
        StructField("col3", StringType, true),
        StructField("col4", StringType, true),
        StructField("col5", StringType, true)))
      
      val csvDF : DataFrame = sqlContext.read
        .option("header", "false")
        .option("delimiter", ",")
        .option("inferSchema", "false")
        .schema(schema)
        .csv("file.csv")
      

      你应该得到dataframe

      +----+----+----+----+----+
      |col1|col2|col3|col4|col5|
      +----+----+----+----+----+
      |A1  |B1  |C1  |null|null|
      |A2  |B2  |C2  |D1  |null|
      |A3  |B3  |C3  |D2  |E1  |
      |A4  |B4  |C4  |D3  |null|
      |A5  |B5  |C5  |null|E2  |
      +----+----+----+----+----+
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2023-03-31
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2020-07-29
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多