【发布时间】:2018-04-06 12:58:30
【问题描述】:
上下文:在我们的数据管道中,我们使用 spark SQL 运行大量查询,这些查询由我们的最终用户作为文本文件提供,然后我们对其进行参数化。
情况:
我们的查询如下所示:
INSERT OVERWRITE TABLE ... PARTITION (...)
SELECT
stuff
FROM
sometable
问题是,当您查看此结果时,它不会创建一堆大小为最大块大小的文件,而是创建 200 个小文件(因为默认情况下 spark 创建 200 个分区)。 (对于某些查询,取决于输入数据和SELECT 查询,对于 200 百阅读数不胜数)。大量的小文件让我们不受系统管理员的欢迎。
已尝试修复(不起作用)
大量文档表明,在这种情况下,您应该使用 DISTRIBUTE BY 以确保给定分区的所有数据都进入同一个分区,所以让我们尝试以下操作:
INSERT OVERWRITE TABLE ... PARTITION (...)
SELECT
stuff
FROM
sometable
DISTRIBUTE BY
1
那么为什么这不起作用(在 spark 2.0 和 spark 2.2 上测试)?它确实成功地将所有数据发送到一个reducer - 所有实际数据都在一个大文件中。但它仍然会创建 200 个文件,其中 199 个是空的! (我知道我们可能应该 DISTRIBUTE BY 我们的分区列,但这是提供最简单的示例)
修复确实有效,但不适合我们的用例
可以通过使用coalesce 或partition 来做正确的事情,因此(在pyspark 语法中):
select = sqlContext.sql('''SELECT stuff FROM sometable''').coalesce(1)
select.write.insertInto(target_table, overwrite=True)
但我不想这样做,因为我们需要彻底改变用户向我们提供查询的方式。
我还看到我们可以设置:
conf.set("spark.sql.shuffle.partitions","1");
但我还没有尝试过,因为我不想强制(相当复杂的)查询中的所有计算都发生在一个 reducer 上,只发生在最后写入磁盘的那个上。 (如果我不应该担心这个,请告诉我!)
问题:
- 仅使用 spark SQL 语法,我如何编写一个尽可能少写入文件且不创建大量空/小文件的查询?
可能相关:
-
merge-multiple-small-files-into-few-larger-files-in-spark(没有解决方案必须是 SparkSQL 的限制,根据上述,
DISTRIBUTE BY实际上不起作用) - spark coalesce doesn't work(仅对我们而言,所以这不是问题)
【问题讨论】:
-
没有回答您的问题,但我们遇到了类似的问题。我们通过使用两个工作解决了这个问题。一个人用 spark 计算结果并将其写到“暂存”位置。下一个作业是一个触发的 bash 脚本,它拾取文件并将其连接在一起并传送到必要的位置,然后清理暂存区域。 bash 脚本非常快。不是很优雅,但绝对有效。
标签: apache-spark hdfs apache-spark-sql