From 61f0632f8aaf8f11d08a11d9c0241c1922504dad Mon Sep 17 00:00:00 2001 From: gatorsmile Date: Wed, 1 Jun 2016 10:36:27 -0700 Subject: [PATCH 1/5] fix. --- .../sql/catalyst/parser/AstBuilder.scala | 11 ++++ .../sql/catalyst/parser/PlanParserSuite.scala | 14 +++-- .../sql/hive/InsertIntoHiveTableSuite.scala | 56 +++++++++++++++++++ 3 files changed, 77 insertions(+), 4 deletions(-) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/parser/AstBuilder.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/parser/AstBuilder.scala index 3473feec3209c..62efcfa10a478 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/parser/AstBuilder.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/parser/AstBuilder.scala @@ -171,6 +171,17 @@ class AstBuilder extends SqlBaseBaseVisitor[AnyRef] with Logging { val tableIdent = visitTableIdentifier(ctx.tableIdentifier) val partitionKeys = Option(ctx.partitionSpec).map(visitPartitionSpec).getOrElse(Map.empty) + if (ctx.EXISTS != null && ctx.partitionSpec == null) { + throw new ParseException( + s"IF NOT EXISTS is not expected when Partition Spec is not specified.", ctx) + } + + val dynamicPartitionKeys = partitionKeys.filter { case (k, v) => v.isEmpty } + if (ctx.EXISTS != null && dynamicPartitionKeys.nonEmpty) { + throw new ParseException(s"Dynamic partitions do not support IF NOT EXISTS. Specified " + + "partitions with value: " + dynamicPartitionKeys.keys.mkString("[", ",", "]"), ctx) + } + InsertIntoTable( UnresolvedRelation(tableIdent, None), partitionKeys, diff --git a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/parser/PlanParserSuite.scala b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/parser/PlanParserSuite.scala index a6fad2d8a0398..bcc43408a775b 100644 --- a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/parser/PlanParserSuite.scala +++ b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/parser/PlanParserSuite.scala @@ -182,14 +182,12 @@ class PlanParserSuite extends PlanTest { // Single inserts assertEqual(s"insert overwrite table s $sql", insert(Map.empty, overwrite = true)) - assertEqual(s"insert overwrite table s if not exists $sql", - insert(Map.empty, overwrite = true, ifNotExists = true)) + assertEqual(s"insert overwrite table s partition (e = 1) if not exists $sql", + insert(Map("e" -> Option("1")), overwrite = true, ifNotExists = true)) assertEqual(s"insert into s $sql", insert(Map.empty)) assertEqual(s"insert into table s partition (c = 'd', e = 1) $sql", insert(Map("c" -> Option("d"), "e" -> Option("1")))) - assertEqual(s"insert overwrite table s partition (c = 'd', x) if not exists $sql", - insert(Map("c" -> Option("d"), "x" -> None), overwrite = true, ifNotExists = true)) // Multi insert val plan2 = table("t").where('x > 5).select(star()) @@ -200,6 +198,14 @@ class PlanParserSuite extends PlanTest { table("u"), Map.empty, plan2, overwrite = false, ifNotExists = false))) } + test ("insert with if not exists") { + val sql = "select * from t" + intercept(s"insert overwrite table s partition (e = 1, x) if not exists $sql", + "Dynamic partitions do not support IF NOT EXISTS. Specified partitions with value: [x]") + intercept(s"insert overwrite table s if not exists $sql", + "IF NOT EXISTS is not expected when Partition Spec is not specified") + } + test("aggregation") { val sql = "select a, b, sum(c) as c from d group by a, b" diff --git a/sql/hive/src/test/scala/org/apache/spark/sql/hive/InsertIntoHiveTableSuite.scala b/sql/hive/src/test/scala/org/apache/spark/sql/hive/InsertIntoHiveTableSuite.scala index fae59001b98e1..f42cc904ae916 100644 --- a/sql/hive/src/test/scala/org/apache/spark/sql/hive/InsertIntoHiveTableSuite.scala +++ b/sql/hive/src/test/scala/org/apache/spark/sql/hive/InsertIntoHiveTableSuite.scala @@ -166,6 +166,62 @@ class InsertIntoHiveTableSuite extends QueryTest with TestHiveSingleton with Bef sql("DROP TABLE tmp_table") } + test("INSERT OVERWRITE - partition IF NOT EXISTS") { + val tmpDir = Utils.createTempDir() + val selQuery = "select c1, p1, p2 from table_with_partition" + sql( + s""" + |CREATE TABLE table_with_partition(c1 string) + |PARTITIONED by (p1 string,p2 string) + |location '${tmpDir.toURI.toString}' + """.stripMargin) + sql( + """ + |INSERT OVERWRITE TABLE table_with_partition + |partition (p1='a',p2='b') + |SELECT 'blarr' + """.stripMargin) + + checkAnswer( + sql(selQuery), + Row("blarr", "a", "b")) + + sql( + """ + |INSERT OVERWRITE TABLE table_with_partition + |partition (p1='a',p2='b') + |SELECT 'blarr2' + """.stripMargin) + + checkAnswer( + sql(selQuery), + Row("blarr2", "a", "b")) + + val e = intercept[AnalysisException] { + sql( + """ + |INSERT OVERWRITE TABLE table_with_partition + |partition (p1='a',p2) IF NOT EXISTS + |SELECT 'blarr3' + """.stripMargin) + } + assert(e.getMessage.contains( + "Dynamic partitions do not support IF NOT EXISTS. Specified partitions with value: [p2]")) + + // If the partition already exists, the insert will overwrite the data + // unless users specify IF NOT EXISTS + sql( + """ + |INSERT OVERWRITE TABLE table_with_partition + |partition (p1='a',p2='b') IF NOT EXISTS + |SELECT 'blarr3' + """.stripMargin) + + checkAnswer( + sql(selQuery), + Row("blarr2", "a", "b")) + } + test("Insert ArrayType.containsNull == false") { val schema = StructType(Seq( StructField("a", ArrayType(StringType, containsNull = false)))) From 70d9cad6d2653f5adfdc77913f23fe366ab56ebe Mon Sep 17 00:00:00 2001 From: gatorsmile Date: Wed, 1 Jun 2016 10:59:44 -0700 Subject: [PATCH 2/5] style fix. --- .../apache/spark/sql/hive/InsertIntoHiveTableSuite.scala | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/sql/hive/src/test/scala/org/apache/spark/sql/hive/InsertIntoHiveTableSuite.scala b/sql/hive/src/test/scala/org/apache/spark/sql/hive/InsertIntoHiveTableSuite.scala index f42cc904ae916..075f1542eb547 100644 --- a/sql/hive/src/test/scala/org/apache/spark/sql/hive/InsertIntoHiveTableSuite.scala +++ b/sql/hive/src/test/scala/org/apache/spark/sql/hive/InsertIntoHiveTableSuite.scala @@ -200,10 +200,10 @@ class InsertIntoHiveTableSuite extends QueryTest with TestHiveSingleton with Bef val e = intercept[AnalysisException] { sql( """ - |INSERT OVERWRITE TABLE table_with_partition - |partition (p1='a',p2) IF NOT EXISTS - |SELECT 'blarr3' - """.stripMargin) + |INSERT OVERWRITE TABLE table_with_partition + |partition (p1='a',p2) IF NOT EXISTS + |SELECT 'blarr3' + """.stripMargin) } assert(e.getMessage.contains( "Dynamic partitions do not support IF NOT EXISTS. Specified partitions with value: [p2]")) From be187d87584dee82f6dc9241bafd631907d67c74 Mon Sep 17 00:00:00 2001 From: gatorsmile Date: Thu, 2 Jun 2016 11:41:59 -0700 Subject: [PATCH 3/5] address comments --- .../antlr4/org/apache/spark/sql/catalyst/parser/SqlBase.g4 | 2 +- .../org/apache/spark/sql/catalyst/parser/AstBuilder.scala | 7 +------ .../apache/spark/sql/catalyst/parser/PlanParserSuite.scala | 3 +-- 3 files changed, 3 insertions(+), 9 deletions(-) diff --git a/sql/catalyst/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBase.g4 b/sql/catalyst/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBase.g4 index b0e71c7e7c7d1..013fab7ebc8e1 100644 --- a/sql/catalyst/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBase.g4 +++ b/sql/catalyst/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBase.g4 @@ -199,7 +199,7 @@ query ; insertInto - : INSERT OVERWRITE TABLE tableIdentifier partitionSpec? (IF NOT EXISTS)? + : INSERT OVERWRITE TABLE tableIdentifier (partitionSpec (IF NOT EXISTS)?)? | INSERT INTO TABLE? tableIdentifier partitionSpec? ; diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/parser/AstBuilder.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/parser/AstBuilder.scala index 62efcfa10a478..a37076efb6613 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/parser/AstBuilder.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/parser/AstBuilder.scala @@ -171,12 +171,7 @@ class AstBuilder extends SqlBaseBaseVisitor[AnyRef] with Logging { val tableIdent = visitTableIdentifier(ctx.tableIdentifier) val partitionKeys = Option(ctx.partitionSpec).map(visitPartitionSpec).getOrElse(Map.empty) - if (ctx.EXISTS != null && ctx.partitionSpec == null) { - throw new ParseException( - s"IF NOT EXISTS is not expected when Partition Spec is not specified.", ctx) - } - - val dynamicPartitionKeys = partitionKeys.filter { case (k, v) => v.isEmpty } + val dynamicPartitionKeys = partitionKeys.filter(_._2.isEmpty) if (ctx.EXISTS != null && dynamicPartitionKeys.nonEmpty) { throw new ParseException(s"Dynamic partitions do not support IF NOT EXISTS. Specified " + "partitions with value: " + dynamicPartitionKeys.keys.mkString("[", ",", "]"), ctx) diff --git a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/parser/PlanParserSuite.scala b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/parser/PlanParserSuite.scala index bcc43408a775b..9d8cba854d53a 100644 --- a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/parser/PlanParserSuite.scala +++ b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/parser/PlanParserSuite.scala @@ -202,8 +202,7 @@ class PlanParserSuite extends PlanTest { val sql = "select * from t" intercept(s"insert overwrite table s partition (e = 1, x) if not exists $sql", "Dynamic partitions do not support IF NOT EXISTS. Specified partitions with value: [x]") - intercept(s"insert overwrite table s if not exists $sql", - "IF NOT EXISTS is not expected when Partition Spec is not specified") + intercept[ParseException](parsePlan(s"insert overwrite table s if not exists $sql")) } test("aggregation") { From bbaad666250532d80c8ce57a33b6b94433bcef76 Mon Sep 17 00:00:00 2001 From: gatorsmile Date: Thu, 2 Jun 2016 13:33:30 -0700 Subject: [PATCH 4/5] address comments. --- .../spark/sql/catalyst/plans/logical/basicLogicalOperators.scala | 1 + 1 file changed, 1 insertion(+) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/basicLogicalOperators.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/basicLogicalOperators.scala index 898784dab1d98..6c3eb3a5a28ab 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/basicLogicalOperators.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/basicLogicalOperators.scala @@ -377,6 +377,7 @@ case class InsertIntoTable( } assert(overwrite || !ifNotExists) + assert(partition.values.forall(_.nonEmpty) || !ifNotExists) override lazy val resolved: Boolean = childrenResolved && table.resolved && expectedColumns.forall { expected => child.output.size == expected.size && child.output.zip(expected).forall { From 84ed71ad89f3571f317694fd5803ce341a3b46b8 Mon Sep 17 00:00:00 2001 From: gatorsmile Date: Sat, 11 Jun 2016 09:00:58 -0700 Subject: [PATCH 5/5] updated the test cases --- .../sql/hive/InsertIntoHiveTableSuite.scala | 116 ++++++++++-------- 1 file changed, 64 insertions(+), 52 deletions(-) diff --git a/sql/hive/src/test/scala/org/apache/spark/sql/hive/InsertIntoHiveTableSuite.scala b/sql/hive/src/test/scala/org/apache/spark/sql/hive/InsertIntoHiveTableSuite.scala index 075f1542eb547..3bf45ced75e0a 100644 --- a/sql/hive/src/test/scala/org/apache/spark/sql/hive/InsertIntoHiveTableSuite.scala +++ b/sql/hive/src/test/scala/org/apache/spark/sql/hive/InsertIntoHiveTableSuite.scala @@ -167,59 +167,71 @@ class InsertIntoHiveTableSuite extends QueryTest with TestHiveSingleton with Bef } test("INSERT OVERWRITE - partition IF NOT EXISTS") { - val tmpDir = Utils.createTempDir() - val selQuery = "select c1, p1, p2 from table_with_partition" - sql( - s""" - |CREATE TABLE table_with_partition(c1 string) - |PARTITIONED by (p1 string,p2 string) - |location '${tmpDir.toURI.toString}' - """.stripMargin) - sql( - """ - |INSERT OVERWRITE TABLE table_with_partition - |partition (p1='a',p2='b') - |SELECT 'blarr' - """.stripMargin) - - checkAnswer( - sql(selQuery), - Row("blarr", "a", "b")) - - sql( - """ - |INSERT OVERWRITE TABLE table_with_partition - |partition (p1='a',p2='b') - |SELECT 'blarr2' - """.stripMargin) - - checkAnswer( - sql(selQuery), - Row("blarr2", "a", "b")) - - val e = intercept[AnalysisException] { - sql( - """ - |INSERT OVERWRITE TABLE table_with_partition - |partition (p1='a',p2) IF NOT EXISTS - |SELECT 'blarr3' - """.stripMargin) + withTempDir { tmpDir => + val table = "table_with_partition" + withTable(table) { + val selQuery = s"select c1, p1, p2 from $table" + sql( + s""" + |CREATE TABLE $table(c1 string) + |PARTITIONED by (p1 string,p2 string) + |location '${tmpDir.toURI.toString}' + """.stripMargin) + sql( + s""" + |INSERT OVERWRITE TABLE $table + |partition (p1='a',p2='b') + |SELECT 'blarr' + """.stripMargin) + checkAnswer( + sql(selQuery), + Row("blarr", "a", "b")) + + sql( + s""" + |INSERT OVERWRITE TABLE $table + |partition (p1='a',p2='b') + |SELECT 'blarr2' + """.stripMargin) + checkAnswer( + sql(selQuery), + Row("blarr2", "a", "b")) + + var e = intercept[AnalysisException] { + sql( + s""" + |INSERT OVERWRITE TABLE $table + |partition (p1='a',p2) IF NOT EXISTS + |SELECT 'blarr3', 'newPartition' + """.stripMargin) + } + assert(e.getMessage.contains( + "Dynamic partitions do not support IF NOT EXISTS. Specified partitions with value: [p2]")) + + e = intercept[AnalysisException] { + sql( + s""" + |INSERT OVERWRITE TABLE $table + |partition (p1='a',p2) IF NOT EXISTS + |SELECT 'blarr3', 'b' + """.stripMargin) + } + assert(e.getMessage.contains( + "Dynamic partitions do not support IF NOT EXISTS. Specified partitions with value: [p2]")) + + // If the partition already exists, the insert will overwrite the data + // unless users specify IF NOT EXISTS + sql( + s""" + |INSERT OVERWRITE TABLE $table + |partition (p1='a',p2='b') IF NOT EXISTS + |SELECT 'blarr3' + """.stripMargin) + checkAnswer( + sql(selQuery), + Row("blarr2", "a", "b")) + } } - assert(e.getMessage.contains( - "Dynamic partitions do not support IF NOT EXISTS. Specified partitions with value: [p2]")) - - // If the partition already exists, the insert will overwrite the data - // unless users specify IF NOT EXISTS - sql( - """ - |INSERT OVERWRITE TABLE table_with_partition - |partition (p1='a',p2='b') IF NOT EXISTS - |SELECT 'blarr3' - """.stripMargin) - - checkAnswer( - sql(selQuery), - Row("blarr2", "a", "b")) } test("Insert ArrayType.containsNull == false") {