Repository navigation
[SPARK-54760][SQL] DelegatingCatalogExtension as session catalog supports both V1 and V2 functions - #53531
[SPARK-54760][SQL] DelegatingCatalogExtension as session catalog supports both V1 and V2 functions#53531pan3793 wants to merge 6 commits into
Conversation
| "errorClass" : "IDENTIFIER_TOO_MANY_NAME_PARTS", | ||
| "sqlState" : "42601", | ||
| "errorClass" : "REQUIRES_SINGLE_PART_NAMESPACE", | ||
| "sqlState" : "42K05", |
There was a problem hiding this comment.
either before or after this change, the error condition of the query VALUES(IDENTIFIER('a.b.c.d')()) makes sense.
before stacktrace
org.apache.spark.sql.AnalysisException: [IDENTIFIER_TOO_MANY_NAME_PARTS] `a`.`b`.`c`.`d` is not a valid identifier as it has more than 2 name parts. SQLSTATE: 42601; line 1 pos 7
at org.apache.spark.sql.errors.QueryCompilationErrors$.identifierTooManyNamePartsError(QueryCompilationErrors.scala:2330)
at org.apache.spark.sql.connector.catalog.CatalogV2Implicits$IdentifierHelper.asFunctionIdentifier(CatalogV2Implicits.scala:197)
at org.apache.spark.sql.catalyst.analysis.FunctionResolution.$anonfun$resolveFunction$2(FunctionResolution.scala:56)
at scala.Option.getOrElse(Option.scala:201)
at org.apache.spark.sql.catalyst.analysis.FunctionResolution.$anonfun$resolveFunction$1(FunctionResolution.scala:52)
at org.apache.spark.sql.catalyst.analysis.package$.withPosition(package.scala:104)
at org.apache.spark.sql.catalyst.analysis.FunctionResolution.resolveFunction(FunctionResolution.scala:52)
...
after stacktrace
org.apache.spark.sql.AnalysisException: [REQUIRES_SINGLE_PART_NAMESPACE] spark_catalog requires a single-part namespace, but got `a`.`b`.`c`. SQLSTATE: 42K05; line 1 pos 7
at org.apache.spark.sql.errors.QueryCompilationErrors$.requiresSinglePartNamespaceError(QueryCompilationErrors.scala:1535)
at org.apache.spark.sql.execution.datasources.v2.V2SessionCatalog$TableIdentifierHelper.asFunctionIdentifier(V2SessionCatalog.scala:381)
at org.apache.spark.sql.execution.datasources.v2.V2SessionCatalog.loadFunction(V2SessionCatalog.scala:478)
at org.apache.spark.sql.connector.catalog.DelegatingCatalogExtension.loadFunction(DelegatingCatalogExtension.java:182)
at org.apache.spark.sql.connector.InMemoryTableSessionCatalog.org$apache$spark$sql$connector$TestV2SessionCatalogBase$$super$loadFunction(DataSourceV2DataFrameSessionCatalogSuite.scala:101)
at org.apache.spark.sql.connector.TestV2SessionCatalogBase.loadFunction(TestV2SessionCatalogBase.scala:139)
at org.apache.spark.sql.connector.TestV2SessionCatalogBase.loadFunction$(TestV2SessionCatalogBase.scala:135)
at org.apache.spark.sql.connector.InMemoryTableSessionCatalog.loadFunction(DataSourceV2DataFrameSessionCatalogSuite.scala:101)
at org.apache.spark.sql.catalyst.analysis.FunctionResolution.$anonfun$resolveFunction$2(FunctionResolution.scala:54)
at scala.Option.getOrElse(Option.scala:201)
at org.apache.spark.sql.catalyst.analysis.FunctionResolution.$anonfun$resolveFunction$1(FunctionResolution.scala:51)
at org.apache.spark.sql.catalyst.analysis.package$.withPosition(package.scala:104)
at org.apache.spark.sql.catalyst.analysis.FunctionResolution.resolveFunction(FunctionResolution.scala:51)
...
| } | ||
| } | ||
|
|
||
| private case object StrLenDefault extends ScalarFunction[Int] { |
There was a problem hiding this comment.
all changes in this file are just exposing those V2 functions to allow reusing in more test cases
| sessionCatalog.createFunction(Identifier.of(Array("ns"), "strlen"), StrLen(StrLenDefault)) | ||
| checkAnswer( | ||
| sql("SELECT char_length('Hello') as v1, ns.strlen('Spark') as v2"), | ||
| Row(5, 5)) |
There was a problem hiding this comment.
it's the effective test case, previously it could not load the v2 func ns.strlen
| if (CatalogV2Util.isSessionCatalog(catalog)) { | ||
| resolveV1Function(ident.asFunctionIdentifier, u.arguments, u) | ||
| } else { | ||
| resolveV2Function(catalog.asFunctionCatalog, ident, u.arguments, u) | ||
| catalog.asFunctionCatalog.loadFunction(ident) match { | ||
| case V1Function(_) => | ||
| resolveV1Function(ident.asFunctionIdentifier, u.arguments, u) | ||
| case unboundV2Func => | ||
| resolveV2Function(unboundV2Func, u.arguments, u) |
There was a problem hiding this comment.
here is the key change - the session catalog also uses V2 catalog API to load the function, the only difference is it may return V1Function, which is a special V2 UnboundFunction and must be resolved differently.
|
@HyukjinKwon, it would be nice if you can consider including this bugfix in the upcoming 4.1.1 |
| resolveV2Function(catalog.asFunctionCatalog, ident, u.arguments, u) | ||
| catalog.asFunctionCatalog.loadFunction(ident) match { | ||
| case V1Function(_) => | ||
| resolveV1Function(ident.asFunctionIdentifier, u.arguments, u) |
There was a problem hiding this comment.
does it mean we look up v1 function twice now? Checking V2SessionCatalog#loadFunction, it already looks up the function.
There was a problem hiding this comment.
@cloud-fan yes, the second time lookup looks cheap, do you have any concerns about that?
There was a problem hiding this comment.
@cloud-fan, as discussed offline, the second time resolution is just a memory operation and won't produce an expensive RPC call to external catalog, I added a comment to explain it.
|
thanks, merging to master/4.1! |
…orts both V1 and V2 functions ### What changes were proposed in this pull request? This PR fixes a bug that occurs when the user uses a custom `DelegatingCatalogExtension` as the session catalog, Spark can not load the v2 function properly provided by the catalog. A typical use case is Iceberg's `SparkSessionCatalog` ``` $ spark-sql \ --conf spark.sql.catalog.spark_catalog=org.apache.iceberg.spark.SparkSessionCatalog \ ... ``` ``` spark-sql (default)> SELECT spark_catalog.system.iceberg_version(); [ROUTINE_NOT_FOUND] The routine `system`.`iceberg_version` cannot be found. Verify the spelling and correctness of the schema and catalog. If you did not qualify the name with a schema and catalog, verify the current_schema() output, or qualify the name with the correct schema and catalog. To tolerate the error on drop use DROP ... IF EXISTS. SQLSTATE: 42883; line 1 pos 7 ``` ### Why are the changes needed? Fix bug. ### Does this PR introduce _any_ user-facing change? Yes, it fixes a bug. ### How was this patch tested? Add new UT. Also manually tested with Iceberg. ``` spark-sql (default)> SELECT spark_catalog.system.iceberg_version(); 1.10.0 Time taken: 1.715 seconds, Fetched 1 row(s) ``` ### Was this patch authored or co-authored using generative AI tooling? No. Closes #53531 from pan3793/SPARK-54760. Authored-by: Cheng Pan <chengpan@apache.org> Signed-off-by: Wenchen Fan <wenchen@databricks.com> (cherry picked from commit 512099b) Signed-off-by: Wenchen Fan <wenchen@databricks.com>
|
Thank you, @pan3793 and @cloud-fan . To @pan3793 , could you make backporting PRs for |
|
@dongjoon-hyun, I think @cloud-fan made a fair decision to backport it only to 4.1, because this is a feature-missing issue instead of a regression or correctness issue. Does the Spark have a clear backport policy, e.g., backport all bug fixes to all active branches? |
|
'branch-4.0' is open for 4.1.1 not a 4.1.0. So, in case of applicable bug fix patches, all live release branches are maintained in the same way. It's only a matter of time when we backport a bug fix from 'master' to 'branch-4.1', 'branch-4.1' to 'branch-4.0', 'branch-4.0' to 'branch-3.5'. It's clear. |
|
It's a feature gap rather than a bug. I did the 4.1 backport so that Iceberg doesn't need to wait for another 6 months, but it's not a serious bug that we need to backport to all the active branches. |
|
This actually adds an extra function lookup request. Let me revert it from 4.1 first. I'm working on a refactor for master branch to avoid the extra lookup. |
### What changes were proposed in this pull request? This PR adds `loadPersistentScalarFunction` method to `SessionCatalog`, to unify the scalar function lookup and resolution between v1 and v2 catalogs. The key difference between v1 and v2 catalogs about function is: - v1 catalog provides different APIs for describing function and invoking function to do function lookup: `lookupPersistentFunction` and `resolvePersistentFunction`. - v2 catalog only has one single lookup API: `loadFunction`. This API returns `UnboundFunction`, for invocation, it needs to enter the second phase and call `UnboundFunction.bind`. The newly added `loadPersistentScalarFunction` method follows the v2 catalog style and handles function invocation (load resources and create function builder) in a second phase. ### Why are the changes needed? Code cleanup, also also avoid redundant function lookup from catalog introduced by #53531 ### Does this PR introduce _any_ user-facing change? No ### How was this patch tested? existing tests ### Was this patch authored or co-authored using generative AI tooling? cursor 2.2.44 Closes #53627 from cloud-fan/func. Authored-by: Wenchen Fan <wenchen@databricks.com> Signed-off-by: Ruifeng Zheng <ruifengz@apache.org>
What changes were proposed in this pull request?
This PR fixes a bug that occurs when the user uses a custom
DelegatingCatalogExtensionas the session catalog, Spark can not load the v2 function properly provided by the catalog. A typical use case is Iceberg'sSparkSessionCatalogWhy are the changes needed?
Fix bug.
Does this PR introduce any user-facing change?
Yes, it fixes a bug.
How was this patch tested?
Add new UT. Also manually tested with Iceberg.
Was this patch authored or co-authored using generative AI tooling?
No.