From 6068e44a0ca7128105f95bc438ed958b57c11f55 Mon Sep 17 00:00:00 2001 From: Herman van Hovell Date: Tue, 18 Oct 2016 19:16:33 -0700 Subject: [PATCH 1/3] Fix unqualified catalog.getFunction --- .../sql/catalyst/expressions/ExpressionInfo.java | 14 ++++++++++++-- .../sql/catalyst/analysis/FunctionRegistry.scala | 2 +- .../sql/catalyst/catalog/SessionCatalog.scala | 10 ++++++++-- .../apache/spark/sql/internal/CatalogImpl.scala | 6 +++--- .../apache/spark/sql/internal/CatalogSuite.scala | 15 ++++++++++++--- 5 files changed, 36 insertions(+), 11 deletions(-) diff --git a/sql/catalyst/src/main/java/org/apache/spark/sql/catalyst/expressions/ExpressionInfo.java b/sql/catalyst/src/main/java/org/apache/spark/sql/catalyst/expressions/ExpressionInfo.java index ba8e9cb4be28b..2710c672430fc 100644 --- a/sql/catalyst/src/main/java/org/apache/spark/sql/catalyst/expressions/ExpressionInfo.java +++ b/sql/catalyst/src/main/java/org/apache/spark/sql/catalyst/expressions/ExpressionInfo.java @@ -25,6 +25,7 @@ public class ExpressionInfo { private String usage; private String name; private String extended; + private String db; public String getClassName() { return className; @@ -42,14 +43,23 @@ public String getExtended() { return extended; } - public ExpressionInfo(String className, String name, String usage, String extended) { + public String getDb() { + return db; + } + + public ExpressionInfo(String className, String name, String usage, String extended, String db) { this.className = className; this.name = name; this.usage = usage; this.extended = extended; + this.db = db; } public ExpressionInfo(String className, String name) { - this(className, name, null, null); + this(className, name, null, null, null); + } + + public ExpressionInfo(String className, String db, String name) { + this(className, name, null, null, db); } } diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/FunctionRegistry.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/FunctionRegistry.scala index b05f4f61f6a3e..ce45641964f83 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/FunctionRegistry.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/FunctionRegistry.scala @@ -495,7 +495,7 @@ object FunctionRegistry { val clazz = scala.reflect.classTag[T].runtimeClass val df = clazz.getAnnotation(classOf[ExpressionDescription]) if (df != null) { - new ExpressionInfo(clazz.getCanonicalName, name, df.usage(), df.extended()) + new ExpressionInfo(clazz.getCanonicalName, name, df.usage(), df.extended(), null) } else { new ExpressionInfo(clazz.getCanonicalName, name) } diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/catalog/SessionCatalog.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/catalog/SessionCatalog.scala index fe41c41a6eb20..b038ecd86d53e 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/catalog/SessionCatalog.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/catalog/SessionCatalog.scala @@ -915,7 +915,10 @@ class SessionCatalog( requireDbExists(db) if (externalCatalog.functionExists(db, name.funcName)) { val metadata = externalCatalog.getFunction(db, name.funcName) - new ExpressionInfo(metadata.className, qualifiedName.unquotedString) + new ExpressionInfo( + metadata.className, + qualifiedName.database.orNull, + qualifiedName.identifier) } else { failFunctionLookup(name.funcName) } @@ -972,7 +975,10 @@ class SessionCatalog( // catalog. So, it is possible that qualifiedName is not exactly the same as // catalogFunction.identifier.unquotedString (difference is on case-sensitivity). // At here, we preserve the input from the user. - val info = new ExpressionInfo(catalogFunction.className, qualifiedName.unquotedString) + val info = new ExpressionInfo( + catalogFunction.className, + qualifiedName.database.orNull, + qualifiedName.funcName) val builder = makeFunctionBuilder(qualifiedName.unquotedString, catalogFunction.className) createTempFunction(qualifiedName.unquotedString, info, builder, ignoreIfExists = false) // Now, we need to create the Expression. diff --git a/sql/core/src/main/scala/org/apache/spark/sql/internal/CatalogImpl.scala b/sql/core/src/main/scala/org/apache/spark/sql/internal/CatalogImpl.scala index f6c297e91b7c5..44fd38dfb96f6 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/internal/CatalogImpl.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/internal/CatalogImpl.scala @@ -133,11 +133,11 @@ class CatalogImpl(sparkSession: SparkSession) extends Catalog { private def makeFunction(funcIdent: FunctionIdentifier): Function = { val metadata = sessionCatalog.lookupFunctionInfo(funcIdent) new Function( - name = funcIdent.identifier, - database = funcIdent.database.orNull, + name = metadata.getName, + database = metadata.getDb, description = null, // for now, this is always undefined className = metadata.getClassName, - isTemporary = funcIdent.database.isEmpty) + isTemporary = metadata.getDb == null) } /** diff --git a/sql/core/src/test/scala/org/apache/spark/sql/internal/CatalogSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/internal/CatalogSuite.scala index 214bc736bd4de..89ec162c8ed52 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/internal/CatalogSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/internal/CatalogSuite.scala @@ -386,15 +386,24 @@ class CatalogSuite createFunction("fn2", Some(db)) // Find a temporary function - assert(spark.catalog.getFunction("fn1").name === "fn1") + val fn1 = spark.catalog.getFunction("fn1") + assert(fn1.name === "fn1") + assert(fn1.database === null) + assert(fn1.isTemporary) // Find a qualified function - assert(spark.catalog.getFunction(db, "fn2").name === "fn2") + val fn2 = spark.catalog.getFunction(db, "fn2") + assert(fn2.name === "fn2") + assert(fn2.database === db) + assert(!fn2.isTemporary) // Find an unqualified function using the current database intercept[AnalysisException](spark.catalog.getFunction("fn2")) spark.catalog.setCurrentDatabase(db) - assert(spark.catalog.getFunction("fn2").name === "fn2") + val unqualified = spark.catalog.getFunction("fn2") + assert(unqualified.name === "fn2") + assert(unqualified.database === db) + assert(!unqualified.isTemporary) } } } From 5597d254ff36d18adf775b50f18bc255c9bbee97 Mon Sep 17 00:00:00 2001 From: Herman van Hovell Date: Tue, 18 Oct 2016 22:42:39 -0700 Subject: [PATCH 2/3] Fix describe function --- .../org/apache/spark/sql/execution/command/functions.scala | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/command/functions.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/command/functions.scala index 26593d2918a6e..f270f95118a60 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/command/functions.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/command/functions.scala @@ -118,14 +118,16 @@ case class DescribeFunctionCommand( case _ => try { val info = sparkSession.sessionState.catalog.lookupFunctionInfo(functionName) + val db = if (info.getDb != null) info.getDb + "." else "" + val name = db + info.getName val result = - Row(s"Function: ${info.getName}") :: + Row(s"Function: $name") :: Row(s"Class: ${info.getClassName}") :: Row(s"Usage: ${replaceFunctionName(info.getUsage, info.getName)}") :: Nil if (isExtended) { result :+ - Row(s"Extended Usage:\n${replaceFunctionName(info.getExtended, info.getName)}") + Row(s"Extended Usage:\n${replaceFunctionName(info.getExtended, name)}") } else { result } From e8d1a7e15b14e1b9799fd65d89797a337ef63c59 Mon Sep 17 00:00:00 2001 From: Herman van Hovell Date: Tue, 1 Nov 2016 10:25:47 +0100 Subject: [PATCH 3/3] Code Review --- .../spark/sql/catalyst/expressions/ExpressionInfo.java | 8 ++++---- .../spark/sql/catalyst/analysis/FunctionRegistry.scala | 2 +- .../apache/spark/sql/execution/command/functions.scala | 3 +-- 3 files changed, 6 insertions(+), 7 deletions(-) diff --git a/sql/catalyst/src/main/java/org/apache/spark/sql/catalyst/expressions/ExpressionInfo.java b/sql/catalyst/src/main/java/org/apache/spark/sql/catalyst/expressions/ExpressionInfo.java index 2710c672430fc..4565ed44877a5 100644 --- a/sql/catalyst/src/main/java/org/apache/spark/sql/catalyst/expressions/ExpressionInfo.java +++ b/sql/catalyst/src/main/java/org/apache/spark/sql/catalyst/expressions/ExpressionInfo.java @@ -47,19 +47,19 @@ public String getDb() { return db; } - public ExpressionInfo(String className, String name, String usage, String extended, String db) { + public ExpressionInfo(String className, String db, String name, String usage, String extended) { this.className = className; + this.db = db; this.name = name; this.usage = usage; this.extended = extended; - this.db = db; } public ExpressionInfo(String className, String name) { - this(className, name, null, null, null); + this(className, null, name, null, null); } public ExpressionInfo(String className, String db, String name) { - this(className, name, null, null, db); + this(className, db, name, null, null); } } diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/FunctionRegistry.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/FunctionRegistry.scala index ce45641964f83..3e836ca375e2e 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/FunctionRegistry.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/FunctionRegistry.scala @@ -495,7 +495,7 @@ object FunctionRegistry { val clazz = scala.reflect.classTag[T].runtimeClass val df = clazz.getAnnotation(classOf[ExpressionDescription]) if (df != null) { - new ExpressionInfo(clazz.getCanonicalName, name, df.usage(), df.extended(), null) + new ExpressionInfo(clazz.getCanonicalName, null, name, df.usage(), df.extended()) } else { new ExpressionInfo(clazz.getCanonicalName, name) } diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/command/functions.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/command/functions.scala index f270f95118a60..24d825f5cb33a 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/command/functions.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/command/functions.scala @@ -118,8 +118,7 @@ case class DescribeFunctionCommand( case _ => try { val info = sparkSession.sessionState.catalog.lookupFunctionInfo(functionName) - val db = if (info.getDb != null) info.getDb + "." else "" - val name = db + info.getName + val name = if (info.getDb != null) info.getDb + "." + info.getName else info.getName val result = Row(s"Function: $name") :: Row(s"Class: ${info.getClassName}") ::