【发布时间】:2021-10-13 01:40:14
【问题描述】:
我有一个 csv 文件,想将它加载到我硬盘上的 parquet 文件中,然后使用 spark-sql CLI 对其运行 SQL 查询。是否有一两个 spark-sql 命令可以做到这一点?
标签: csv apache-spark-sql parquet
我有一个 csv 文件,想将它加载到我硬盘上的 parquet 文件中,然后使用 spark-sql CLI 对其运行 SQL 查询。是否有一两个 spark-sql 命令可以做到这一点?
标签: csv apache-spark-sql parquet
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|
// +---+
}
【讨论】: