diff --git a/pramen/core/src/main/scala/za/co/absa/pramen/core/utils/hive/HiveHelper.scala b/pramen/core/src/main/scala/za/co/absa/pramen/core/utils/hive/HiveHelper.scala index 836a01f4..2c5f54e3 100644 --- a/pramen/core/src/main/scala/za/co/absa/pramen/core/utils/hive/HiveHelper.scala +++ b/pramen/core/src/main/scala/za/co/absa/pramen/core/utils/hive/HiveHelper.scala @@ -48,7 +48,8 @@ abstract class HiveHelper { partitionBy: Seq[String], partitionValues: Seq[String], databaseName: Option[String], - tableName: String): Unit + tableName: String, + location: String): Unit def repairHiveTable(databaseName: Option[String], tableName: String, diff --git a/pramen/core/src/main/scala/za/co/absa/pramen/core/utils/hive/HiveHelperSparkCatalog.scala b/pramen/core/src/main/scala/za/co/absa/pramen/core/utils/hive/HiveHelperSparkCatalog.scala index 819576ce..6a05e817 100644 --- a/pramen/core/src/main/scala/za/co/absa/pramen/core/utils/hive/HiveHelperSparkCatalog.scala +++ b/pramen/core/src/main/scala/za/co/absa/pramen/core/utils/hive/HiveHelperSparkCatalog.scala @@ -17,7 +17,6 @@ package za.co.absa.pramen.core.utils.hive import org.apache.spark.sql.SparkSession -import org.apache.spark.sql.execution.datasources.PartitioningUtils.PartitionValues import org.apache.spark.sql.types.StructType import org.slf4j.LoggerFactory import za.co.absa.pramen.api.CatalogTable @@ -108,7 +107,8 @@ class HiveHelperSparkCatalog(spark: SparkSession) extends HiveHelper { partitionBy: Seq[String], partitionValues: Seq[String], databaseName: Option[String], - tableName: String): Unit = { + tableName: String, + location: String): Unit = { if (partitionBy.length != partitionValues.length) { throw new IllegalArgumentException(s"Partition columns and values must have the same length. Columns: $partitionBy, values: $partitionValues") } @@ -128,8 +128,22 @@ class HiveHelperSparkCatalog(spark: SparkSession) extends HiveHelper { .replace("@partitionClause", partitionClause) .replace("@schema", schemaDDL) - log.info(s"Executing: $sql") - spark.sql(sql).collect() + try { + log.info(s"Executing: $sql") + spark.sql(sql).collect() + } catch { + case ex: Throwable if ex.getMessage != null && ex.getMessage.toLowerCase.contains("partition not found") => + log.info(s"Partition not found for $fullTableName, partition: $partitionClause. Adding partition...") + try { + addPartition(databaseName, tableName, partitionBy, partitionValues, location) + } catch { + case NonFatal(ex) => + log.warn(s"Failed to add partition for $fullTableName, partition: $partitionClause", ex) + } + + log.info(s"Executing: $sql") + spark.sql(sql).collect() + } } override def repairHiveTable(databaseName: Option[String], diff --git a/pramen/core/src/main/scala/za/co/absa/pramen/core/utils/hive/HiveHelperSql.scala b/pramen/core/src/main/scala/za/co/absa/pramen/core/utils/hive/HiveHelperSql.scala index 7803eba3..34077400 100644 --- a/pramen/core/src/main/scala/za/co/absa/pramen/core/utils/hive/HiveHelperSql.scala +++ b/pramen/core/src/main/scala/za/co/absa/pramen/core/utils/hive/HiveHelperSql.scala @@ -20,6 +20,8 @@ import org.apache.spark.sql.types.StructType import org.slf4j.LoggerFactory import za.co.absa.pramen.core.utils.SparkUtils +import scala.util.control.NonFatal + class HiveHelperSql(val queryExecutor: QueryExecutor, hiveConfig: HiveQueryTemplates, alwaysEscapeColumnNames: Boolean) extends HiveHelper { @@ -90,7 +92,8 @@ class HiveHelperSql(val queryExecutor: QueryExecutor, partitionBy: Seq[String], partitionValues: Seq[String], databaseName: Option[String], - tableName: String): Unit = { + tableName: String, + location: String): Unit = { if (partitionBy.length != partitionValues.length) { throw new IllegalArgumentException(s"Partition columns and values must have the same length. Columns: $partitionBy, values: $partitionValues") } @@ -100,8 +103,22 @@ class HiveHelperSql(val queryExecutor: QueryExecutor, log.info(s"Replacing partition schema for $fullTableName, partition: $partitionClause...") - val sql = applyPartitionTemplate(hiveConfig.replacePartitionSchemaTemplate, fullTableName, "", partitionClause, schemaDDL) - queryExecutor.execute(sql) + val sql = applyPartitionTemplate(hiveConfig.replacePartitionSchemaTemplate, fullTableName, location, partitionClause, schemaDDL) + + try { + queryExecutor.execute(sql) + } catch { + case ex: Throwable if ex.getMessage != null && ex.getMessage.toLowerCase.contains("partition not found") => + log.info(s"Partition not found for $fullTableName, partition: $partitionClause. Adding partition...") + try { + addPartition(databaseName, tableName, partitionBy, partitionValues, location) + } catch { + case NonFatal(ex) => + log.warn(s"Failed to add partition for $fullTableName, partition: $partitionClause", ex) + } + + queryExecutor.execute(sql) + } } override def repairHiveTable(databaseName: Option[String], diff --git a/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/utils/hive/HiveHelperSqlSuite.scala b/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/utils/hive/HiveHelperSqlSuite.scala index a82d3c77..d13b39f1 100644 --- a/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/utils/hive/HiveHelperSqlSuite.scala +++ b/pramen/core/src/test/scala/za/co/absa/pramen/core/tests/utils/hive/HiveHelperSqlSuite.scala @@ -177,7 +177,7 @@ class HiveHelperSqlSuite extends AnyWordSpec with SparkTestBase with TempDirFixt val hiveHelper = new HiveHelperSql(qe, defaultHiveConfig, true) val schema = spark.read.parquet(path).withColumn("b", lit(1)).schema - hiveHelper.replaceHivePartitionSchema(schema, "a" :: "b" :: Nil, Seq("AA", "22"), Some("db"), "tbl") + hiveHelper.replaceHivePartitionSchema(schema, "a" :: "b" :: Nil, Seq("AA", "22"), Some("db"), "tbl", "") val actual = qe.queries.mkString("\n")