【发布时间】:2019-06-19 16:04:32
【问题描述】:
我正在尝试使用 bucketing 技术对 Spark 作业进行一些性能优化。我正在阅读.parquet 和.csv 文件并进行一些转换。在我进行分桶并加入两个 DataFrame 之后。然后我将加入 DF 写入镶木地板,但我有一个空文件 ~500B 而不是 500Mb。
- Cloudera (cdh5.15.1)
- Spark 2.3.0
-
斑点
val readParquet = spark.read.parquet(inputP) readParquet .write .format("parquet") .bucketBy(23, "column") .sortBy("column") .mode(SaveMode.Overwrite) .saveAsTable("bucketedTable1") val firstTableDF = spark.table("bucketedTable1") val readCSV = spark.read.csv(inputCSV) readCSV .filter(..) .ordrerBy(someColumn) .write .format("parquet") .bucketBy(23, "column") .sortBy("column") .mode(SaveMode.Overwrite) .saveAsTable("bucketedTable2") val secondTableDF = spark.table("bucketedTable2") val resultDF = secondTableDF .join(firstTableDF, Seq("column"), "fullouter") . . resultDF .coalesce(1) .write .mode(SaveMode.Overwrite) .parquet(output)
当我在命令行中使用 ssh 启动 Spark 作业时,我得到了正确的结果,~500Mb parquet 文件,我可以使用 Hive 看到该文件。如果我使用 oozie 工作流运行相同的作业,我有一个空文件 (~500 Bytes)。
当我在resultDF 上执行.show() 时,我可以看到数据,但我有空的镶木地板文件。
+-----------+---------------+----------+
| col1| col2 | col3|
+-----------+---------------+----------+
|33601234567|208012345678910| LOL|
|33601234567|208012345678910| LOL|
|33601234567|208012345678910| LOL|
当我不将数据保存为表格时,写入 parquet 没有问题。它只发生在从表创建的 DF 中。
有什么建议吗?
提前感谢您的任何想法!
【问题讨论】:
-
空文件有 _SUCCESS 吗?
-
忘记现有的数据框,在您编写现有的有问题的镶木地板之前,在相同的流程中通过 oozie 写入此数据框
spark.sparkContext.parallelize(1 to 4).toDF .coalesce(1) .write.mode(SaveMode.Overwrite).parquet(destDir)看看会发生什么,不要认为这是火花问题。可能是其他问题 -
@RamGhadiyaram 是的,有 _SUCCESS 文件。没有分桶和
.saveAsTable()效果很好,我只是试图避免随机播放问题。 -
@RamGhadiyaram 您对
spark.sparkContext.parallelize(1 to 4).toDF .coalesce(1) .write.mode(SaveMode.Overwrite).parquet(destDir)的建议适用于 Oozie。我可以查询表中的数据。 -
Seq("coumn"), 可能是 Seq("column") 我认为这是剩余的类型对我来说很好
标签: apache-spark apache-spark-sql cloudera parquet apache-spark-2.3