我已经开始在Spark 1.4.0中使用Spark SQL和DataFrames.我想在Scala中定义DataFrame上的自定义分区程序,但是没有看到如何执行此操作.
我正在使用的一个数据表包含一个事务列表,按帐户,silimar到下面的示例.
Account Date Type Amount
1001 2014-04-01 Purchase 100.00
1001 2014-04-01 Purchase 50.00
1001 2014-04-05 Purchase 70.00
1001 2014-04-01 Payment -150.00
1002 2014-04-01 Purchase 80.00
1002 2014-04-02 Purchase 22.00
1002 2014-04-04 Payment -120.00
1002 2014-04-04 Purchase 60.00
1003 2014-04-02 Purchase 210.00
1003 2014-04-03 Purchase 15.00
Run Code Online (Sandbox Code Playgroud)
至少在最初,大多数计算将发生在帐户内的交易之间.所以我希望对数据进行分区,以便帐户的所有事务都在同一个Spark分区中.
但我没有看到定义这个的方法.DataFrame类有一个名为"repartition(Int)"的方法,您可以在其中指定要创建的分区数.但我没有看到任何方法可用于为DataFrame定义自定义分区程序,例如可以为RDD指定.
源数据存储在Parquet中.我确实看到在向Parquet编写DataFrame时,您可以指定要分区的列,因此我可以告诉Parquet通过"帐户"列对其数据进行分区.但是可能有数百万个帐户,如果我正确理解Parquet,它会为每个帐户创建一个独特的目录,因此这听起来不是一个合理的解决方案.
有没有办法让Spark分区这个DataFrame,以便一个帐户的所有数据都在同一个分区?
上下文
Spark 2.0.1,在集群模式下spark-submit.我正在读取hdfs的镶木地板文件:
val spark = SparkSession.builder
.appName("myApp")
.config("hive.metastore.uris", "thrift://XXX.XXX.net:9083")
.config("spark.sql.sources.bucketing.enabled", true)
.enableHiveSupport()
.getOrCreate()
val df = spark.read
.format("parquet")
.load("hdfs://XXX.XX.X.XX/myParquetFile")
Run Code Online (Sandbox Code Playgroud)
我保存df到50个桶的蜂巢表分组userid:
df0.write
.bucketBy(50, "userid")
.saveAsTable("myHiveTable")
Run Code Online (Sandbox Code Playgroud)
现在,当我查看hdfs的hive仓库时,/user/hive/warehouse有一个名为的文件夹myHiveTable.里面是一堆part-*.parquet文件.我希望有50个文件.但不,有3201个文件!!!! 每个分区有64个文件,为什么?对于我保存为hive表的不同文件,每个分区有不同数量的文件.所有文件都很小,每个只有几十Kb!
我要补充的,不同的,这个数字userid大约是1 000 000在myParquetFile.
题
为什么文件夹中有3201个文件而不是50个!这些是什么?
当我将此表读回DataFrame并打印分区数时:
val df2 = spark.sql("SELECT * FROM myHiveTable")
println(df2.rdd.getNumPartitions)
Run Code Online (Sandbox Code Playgroud)
分区数isIt正确50,我确认数据被正确分区userid.
对于我的一个大型数据集3Tb,我创建了一个包含1000个分区的表,这些分区创建了大约数百万个文件!这超出了目录项限制1048576并给出org.apache.hadoop.hdfs.protocol.FSLimitException$MaxDirectoryItemsExceededException
题
创建的文件数量取决于什么?
题
有没有办法限制创建的文件数量?
题
我应该担心这些文件吗?df2拥有所有这些文件会损害性能吗?总是说我们不应该创建太多分区,因为它有问题.
题
我发现这个信息HIVE动态分区提示文件的数量可能与映射器的数量有关.建议distribute by在插入蜂巢表时使用.我怎么能在Spark中做到这一点?
题 …
我想GROUP BY在正确分区上执行一个子句,DataFrame同时按作为分区键的列进行分组。显然,在这种情况下,实际上不需要改组,因为所有相等的键都已经位于相同的分区中。但是,我无法弄清楚如何真正避免这种洗牌以及是否有可能。我在 上尝试了分桶和分区选项DataFrameWriter,但随着我继续看到计划中的交流,这些选项似乎没有多大帮助。除了说,还有什么方法可以做类似的事情mapPartitions吗?
apache-spark ×3
bigdata ×1
dataframe ×1
group-by ×1
hive ×1
partitioning ×1
scala ×1
shuffle ×1
sql ×1