I'm a trying to create a table with spark.catalog.createTable. It needs to have a partition column named "id".
Based on How can I use "spark.catalog.createTable" function to create a partitioned table? which is in Scala, I tried :
df = spark.range(10).withColumn("foo", F.lit("bar"))
spark.catalog.createTable("default.test_partition", schema=df.schema, **{"partitionColumnNames":"id"})
But it is not working. It creates a table in Hive with these properties :
CREATE TABLE default.test_partition ( id BIGINT, foo STRING )
WITH SERDEPROPERTIES ('partitionColumnNames'='id' ...
The DDL of the table should actually be:
CREATE TABLE default.test_partition ( foo STRING )
PARTITIONED BY ( id BIGINT )
WITH SERDEPROPERTIES (...
The signature of the method is :
Signature: spark.catalog.createTable(tableName, path=None, source=None, schema=None, **options)
So, I believe there is a special argument in **options to create the partition, but I tried "partitionColumnNames", "partitionBy", "partition" ... none of them is working.
Do you know what is the proper keyword for that ?
EDIT: If you are wondering why I want to use this method, there are 2 reasons :
insertInto method (see Overwrite specific partitions in spark dataframe write method). But this method needs the table to be created first, action that I want to perform with the spark.catalog.createTable because it seems right.Looking at the source code for spark.catalog here , it looks like that the keyword argument options is an alternative for schema and is only used when the schema parameter is not passed. This can be seen below:
"Optionally, a schema can be provided as the schema of the returned" #for options
if path is not None:
options["path"] = path
...................
...........
if schema is None: #this line and the line below
df = self._jcatalog.createTable(tableName, source, description, options)
else:
if not isinstance(schema, StructType):
raise TypeError("schema should be StructType")
scala_datatype = self._jsparkSession.parseDataType(schema.json())
df = self._jcatalog.createTable(
tableName, source, scala_datatype, description, options)
return DataFrame(df, self._sparkSession._wrapped)
However, if you are looking to create an automated DDL process, something along the lines of the below function might help you:
def mycreateTable(tablename,schema,partitioncols):
schema_json = schema.json()
ddlstring = (spark.sparkContext._jvm.org.apache.spark.sql.types.
DataType.fromJson(schema_json).toDDL())
#if you dont want to DROP the table when it exists change the below line
spark.sql(f"""DROP TABLE IF EXISTS {tablename} ;""")
spark.sql(f"""
CREATE TABLE {tablename} ({ddlstring}) partitioned by ({','.join(partitioncols)}) """)
Now executing the below should work:
mycreateTable("default.test_partition",df.schema,['id'])

If you love us? You can donate to us via Paypal or buy me a coffee so we can maintain and grow! Thank you!
Donate Us With