老话题,但我认为如果没有正确回答,即使是老话题也很重要。
在 spark 版本中 >=2 csv 包已经包含在内,您需要将 databricks csv 包导入到您的工作中,例如“--packages com.databricks:spark-csv_2.10:1.5.0”。
示例 csv:
id,name,date
1,pete,2017-10-01 16:12
2,paul,2016-10-01 12:23
3,steve,2016-10-01 03:32
4,mary,2018-10-01 11:12
5,ann,2018-10-02 22:12
6,rudy,2018-10-03 11:11
7,mike,2018-10-04 10:10
首先,您需要创建 hivetable,以便 spark 写入的数据与 hive 架构兼容。 (在未来的版本中可能不再需要)
创建表:
create table part_parq_table (
id int,
name string
)
partitioned by (date string)
stored as parquet
完成此操作后,您可以轻松读取 csv 并将数据框保存到该表中。第二步用“yyyy-mm-dd”之类的日期格式覆盖列日期。将为每个值创建一个文件夹,其中包含特定的行。
SCALA Spark-Shell 示例:
spark.sqlContext.setConf("hive.exec.dynamic.partition", "true")
spark.sqlContext.setConf("hive.exec.dynamic.partition.mode", "nonstrict")
前两行是 hive 配置,需要创建一个尚不存在的分区文件夹。
var df=spark.read.format("csv").option("header","true").load("/tmp/test.csv")
df=df.withColumn("date",substring(col("date"),0,10))
df.show(false)
df.write.format("parquet").mode("append").insertInto("part_parq_table")
插入完成后,您可以直接查询表,如“select * from part_parq_table”。
这些文件夹将在默认 cloudera 的 tablefolder 中创建,例如hdfs:///users/hive/warehouse/part_parq_table
希望有所帮助
BR