【发布时间】:2021-11-20 13:28:07
【问题描述】:
我正在尝试使用spark.catalog.createTable 创建一个表。它需要有一个名为“id”的分区列。
基于 Scala 中的How can I use "spark.catalog.createTable" function to create a partitioned table?,我尝试了:
df = spark.range(10).withColumn("foo", F.lit("bar"))
spark.catalog.createTable("default.test_partition", schema=df.schema, **{"partitionColumnNames":"id"})
但它不起作用。 它在 Hive 中创建一个具有这些属性的表:
CREATE TABLE default.test_partition ( id BIGINT, foo STRING )
WITH SERDEPROPERTIES ('partitionColumnNames'='id' ...
表的DDL实际上应该是:
CREATE TABLE default.test_partition ( foo STRING )
PARTITIONED BY ( id BIGINT )
WITH SERDEPROPERTIES (...
方法的签名是:
Signature: spark.catalog.createTable(tableName, path=None, source=None, schema=None, **options)
所以,我相信**options 中有一个特殊的参数来创建分区,但我尝试了“partitionColumnNames”、“partitionBy”、“partition”......它们都不起作用。
你知道什么是正确的关键字吗?
编辑: 如果您想知道我为什么要使用这种方法,有两个原因:
- 出于个人好奇心,想知道什么可行或不可行
- 我想在 spark 中使用“动态分区”,这需要我使用
insertInto方法(参见Overwrite specific partitions in spark dataframe write method)。但是这种方法需要先创建表,我想用spark.catalog.createTable执行的操作,因为它看起来是正确的。
【问题讨论】:
-
我还没有使用 spark.catalog 但查看源代码 here ,看起来
optionskwarg 仅在未提供架构时使用。if schema is None: df = self._jcatalog.createTable(tableName, source, description, options)。看起来他们没有使用那个 kwarg 进行分区 -
还想知道您是否正在尝试自动化任何 ddl 进程?
-
@anky 我添加了一个编辑来解释原因。是的,这是某种自动化的 ddl 过程。但如果你有更好的想法,我很乐意。
-
我明白了。我已经根据我的理解添加了一个答案,但也可以根据需要进行定制:)