【问题标题】:Use spark-sql cli to load csv data directly into parquet table使用 spark-sql cli 将 csv 数据直接加载到 parquet 表中
【发布时间】:2021-10-13 01:40:14
【问题描述】:

我有一个 csv 文件,想将它加载到我硬盘上的 parquet 文件中,然后使用 spark-sql CLI 对其运行 SQL 查询。是否有一两个 spark-sql 命令可以做到这一点?

标签: csv apache-spark-sql parquet


【解决方案1】:
package spark

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions.{col, trim}

object csv2parquet extends App {
  val spark = SparkSession.builder()
    .master("local")
    .appName("CSV-Parquet")
    .getOrCreate()

  import spark.implicits._
  val sourceFile = "/<path file>/test.csv" // bad data in file
  val targetFile = "/<path file>/testResult.parquet"
  // read csv file
  val df1 = spark.read.option("header", false).csv(sourceFile)
  df1.show(false)
  //    +-------+-------+----------+-----------+
  //    |_c0    |_c1    |_c2       |_c3        |
  //    +-------+-------+----------+-----------+
  //    |Header |TestApp|2020-01-01|null       |
  //    |name   | dept  | age      | batchDate |
  //    |john   | dept1 | 33       | 2020-01-01|
  //    |john   | dept1 | 33       | 2020-01-01|
  //    |john   | dept1 | 33       | 2020-01-01|
  //    |john   | dept1 | 33       | 2020-01-01|
  //    |Trailer|count  |4         |null       |
  //    +-------+-------+----------+-----------+

  // write data to parquet. 
  df1.write.mode("append").parquet(targetFile)

  val resDF = spark.read.parquet(targetFile)
  resDF.show(false)
  //          +-------+-------+----------+-----------+
  //          |_c0    |_c1    |_c2       |_c3        |
  //          +-------+-------+----------+-----------+
  //          |Header |TestApp|2020-01-01|null       |
  //          |name   | dept  | age      | batchDate |
  //          |john   | dept1 | 33       | 2020-01-01|
  //          |john   | dept1 | 33       | 2020-01-01|
  //          |john   | dept1 | 33       | 2020-01-01|
  //          |john   | dept1 | 33       | 2020-01-01|
  //          |Trailer|count  |4         |null       |
  //          +-------+-------+----------+-----------+
  // try sql
  resDF
    .filter(trim(col("_c2")).equalTo(33))
    .select(col("_c2"))
    .show(false)
    //          +---+
    //          |_c2|
    //          +---+
    //          | 33|
    //          | 33|
    //          | 33|
    //          | 33|
    //          +---+
}

【讨论】:

  • 您能否解释一下为什么要多次编写输出 parquet 文件?还是更新代码?谢谢。
  • 1.例如,仅用于收到许多行。 2. 你只需要一个命令 df1.write.mode("append").parquet(targetFile) 或者如果你使用 table 你可以使用 ....format("parquet").saveAsTable(....)
猜你喜欢
  • 2018-01-31
  • 2017-07-29
  • 2020-11-09
  • 1970-01-01
  • 2019-12-23
  • 1970-01-01
  • 2017-01-16
  • 1970-01-01
  • 2017-03-13
相关资源
最近更新 更多