【问题标题】:Spark SQLContext Query with header带有标头的 Spark SQLContext 查询
【发布时间】:2018-12-20 00:50:36
【问题描述】:

我正在使用 SQLContext 读取这样的 CSV 文件:

val csvContents = sqlContext.read.sql("SELECT * FROM 
                  csv.`src/test/resources/afile.csv` WHERE firstcolumn=21")

但它将第一列打印为_c0 并包括其下方的标题。如何设置标头并使用 SQL 查询?我见过这个解决方案:

 val df = spark.read
         .option("header", "true") //reading the headers
         .csv("file.csv")

但这不允许我使用WHERE 子句进行SELECT 查询。有没有办法指定 CSV 标头并进行 SQL SELECT 查询?

【问题讨论】:

  • 它无法从我的 where 子句中识别列名,即使该名称存在。我做了:.select("*").where("col_id=22") 但它给出了一个例外:cannot resolve '`col_id`' given input columns: [col_id col_name
  • 显示src/test/resources/afile.csv。看起来 CSV 数据源无法识别标题,并且您最终得到了一个长列名称,例如col_id col_name(注意列名之间的空格)。
  • 顺便说一句,您使用什么 Spark 版本?

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


【解决方案1】:

您可以在从数据框创建视图后使用 sql 查询。像这样的。

val df = spark.read
  .option("header", "true") //reading the headers
  .csv("file.csv")

df.createOrReplaceTempView("table")

val sqlDf = spark.sql("SELECT * FROM table WHERE firstcolumn=21")

希望这会有所帮助。

【讨论】:

    【解决方案2】:

    首先,如果您使用 Spark 2.0 或以后尝试开始使用 SparkSession 而不是 SparkContext,那么如果您的列数很少,我建议您作为一个好习惯作为另一种选择

    import org.apache.spark.sql.types._    
    val schema = StructType(
      StructField("firstcolumn", StringType, true), 
      StructField("secondcolumn", IntegerType, true)
    )
    
    val df = spark.
      read.
      option("header", true).
      schema(schema).
      csv("file.csv")
    

    因此您可以选择具有正确名称的列

    val etl = df.select("firstcolumn").where("secondcolumn=0")
    

    【讨论】:

    • 它无法从我的 where 子句中识别列名,即使该名称存在。我做了:.select("*").where("col_id=22") 但它给出了一个例外:cannot resolve '`col_id`' given input columns: [col_id col_name
    • 在阅读 CSV 后立即运行 df.show 之前,您能否尝试 df.where(col("col_id") = 22)df.where("col_id = 22")
    【解决方案3】:

    原来没有正确解析标题。 CSV 文件是制表符分隔的,所以我必须明确指定:

    val csvContents = sqlContext.read
            .option("delimiter", "\t")
            .option("header", "true")
            .csv(csvPath)
            .select("*")
            .where(s"col_id=22")
    

    【讨论】:

      【解决方案4】:
      1. 初始化 SparkSession
      2. val fileDF = spark.read.format("csv").option("header",true).load("file.csv")
      3. 发布此内容您可以访问列
           import spark.implicits._  
           fileDF.select($"columnName").where(conditions)
      

      【讨论】:

        猜你喜欢
        • 2015-03-02
        • 2017-01-29
        • 2015-08-02
        • 2016-11-24
        • 1970-01-01
        • 2019-06-05
        • 2016-07-28
        • 2011-10-05
        • 1970-01-01
        相关资源
        最近更新 更多