From 8301dc3c327167d00d7c145c1ec3deb1f24c337f Mon Sep 17 00:00:00 2001 From: PA Savalle Date: Thu, 27 Aug 2026 17:31:47 +0200 Subject: [PATCH 1/4] [SPARK-50569][CONNECT]: clean up data cached by individual sessions --- .../spark/sql/connect/config/Connect.scala | 8 ++ .../sql/connect/service/SessionHolder.scala | 4 + .../SparkConnectSessionManagerSuite.scala | 70 ++++++++++++++++- .../spark/sql/execution/CacheManager.scala | 77 ++++++++++++++++--- 4 files changed, 148 insertions(+), 11 deletions(-) diff --git a/sql/connect/server/src/main/scala/org/apache/spark/sql/connect/config/Connect.scala b/sql/connect/server/src/main/scala/org/apache/spark/sql/connect/config/Connect.scala index 10531e063224d..62f587c0ab01b 100644 --- a/sql/connect/server/src/main/scala/org/apache/spark/sql/connect/config/Connect.scala +++ b/sql/connect/server/src/main/scala/org/apache/spark/sql/connect/config/Connect.scala @@ -156,6 +156,14 @@ object Connect { .timeConf(TimeUnit.MILLISECONDS) .createWithDefaultString("30s") + val CONNECT_SESSION_MANAGER_CLEANUP_CACHED_DATA_ENABLED = + buildStaticConf("spark.connect.session.manager.cleanupCachedData.enabled") + .doc("When true, cached data persisted by an isolated session is removed when the session " + + "is closed. Cached data that is also persisted by another session is preserved.") + .version("4.4.0") + .booleanConf + .createWithDefault(false) + val CONNECT_EXECUTE_MANAGER_DETACHED_TIMEOUT = buildStaticConf("spark.connect.execute.manager.detachedTimeout") .internal() diff --git a/sql/connect/server/src/main/scala/org/apache/spark/sql/connect/service/SessionHolder.scala b/sql/connect/server/src/main/scala/org/apache/spark/sql/connect/service/SessionHolder.scala index e97bcb0c786e3..0c76cafc40f51 100644 --- a/sql/connect/server/src/main/scala/org/apache/spark/sql/connect/service/SessionHolder.scala +++ b/sql/connect/server/src/main/scala/org/apache/spark/sql/connect/service/SessionHolder.scala @@ -434,6 +434,10 @@ case class SessionHolder(userId: String, sessionId: String, session: SparkSessio // Clean up ML cache (only if ML models were created) mlCache.close() + if (SparkEnv.get.conf.get(Connect.CONNECT_SESSION_MANAGER_CLEANUP_CACHED_DATA_ENABLED)) { + session.sharedState.cacheManager.clearCache(session) + } + session.cleanupPythonWorkerLogs() eventManager.postClosed() diff --git a/sql/connect/server/src/test/scala/org/apache/spark/sql/connect/service/SparkConnectSessionManagerSuite.scala b/sql/connect/server/src/test/scala/org/apache/spark/sql/connect/service/SparkConnectSessionManagerSuite.scala index 5fae135bcf18a..e5bbc683d9084 100644 --- a/sql/connect/server/src/test/scala/org/apache/spark/sql/connect/service/SparkConnectSessionManagerSuite.scala +++ b/sql/connect/server/src/test/scala/org/apache/spark/sql/connect/service/SparkConnectSessionManagerSuite.scala @@ -21,14 +21,28 @@ import java.util.UUID import org.scalatest.time.SpanSugar._ -import org.apache.spark.SparkSQLException +import org.apache.spark.{SparkEnv, SparkSQLException} import org.apache.spark.sql.SparkSession +import org.apache.spark.sql.connect.config.Connect import org.apache.spark.sql.pipelines.graph.{DataflowGraph, PipelineUpdateContextImpl} import org.apache.spark.sql.pipelines.logging.PipelineEvent import org.apache.spark.sql.test.SharedSparkSession class SparkConnectSessionManagerSuite extends SharedSparkSession { + private def withSparkConf(pairs: (String, String)*)(f: => Unit): Unit = { + val conf = SparkEnv.get.conf + val previousValues = pairs.map { case (key, _) => key -> conf.getOption(key) } + pairs.foreach { case (key, value) => conf.set(key, value) } + try f + finally { + previousValues.foreach { + case (key, Some(value)) => conf.set(key, value) + case (key, None) => conf.remove(key) + } + } + } + override def beforeEach(): Unit = { super.beforeEach() SparkConnectService.sessionManager.invalidateAllSessions() @@ -177,6 +191,60 @@ class SparkConnectSessionManagerSuite extends SharedSparkSession { "pipeline execution was not removed") } + test("SPARK-50569: cached data cleanup on session close is configurable and isolated") { + Seq(false, true).foreach { cleanupCachedData => + withClue(s"cleanupCachedData=$cleanupCachedData") { + withSparkConf( + Connect.CONNECT_SESSION_MANAGER_CLEANUP_CACHED_DATA_ENABLED.key -> + cleanupCachedData.toString) { + val first = SparkConnectService.sessionManager.getOrCreateIsolatedSession( + SessionKey("user", UUID.randomUUID().toString), None) + val second = SparkConnectService.sessionManager.getOrCreateIsolatedSession( + SessionKey("user", UUID.randomUUID().toString), None) + val firstDataFrame = first.session.range(1) + second.session.range(1, 2).createTempView("second_view") + second.session.catalog.cacheTable("second_view") + val secondDataFrame = second.session.table("second_view") + + firstDataFrame.persist() + SparkConnectService.sessionManager.closeSession(first.key) + + assert( + first.session.sharedState.cacheManager.lookupCachedData(firstDataFrame).isDefined === + !cleanupCachedData) + assert( + second.session.sharedState.cacheManager.lookupCachedData(secondDataFrame).isDefined) + + SparkConnectService.sessionManager.closeSession(second.key) + assert( + second.session.sharedState.cacheManager.lookupCachedData(secondDataFrame).isDefined === + !cleanupCachedData) + spark.catalog.clearCache() + } + } + } + } + + test("SPARK-50569: cached data cleanup preserves entries persisted by another session") { + withSparkConf(Connect.CONNECT_SESSION_MANAGER_CLEANUP_CACHED_DATA_ENABLED.key -> "true") { + val first = SparkConnectService.sessionManager.getOrCreateIsolatedSession( + SessionKey("user", UUID.randomUUID().toString), None) + val second = SparkConnectService.sessionManager.getOrCreateIsolatedSession( + SessionKey("user", UUID.randomUUID().toString), None) + val firstDataFrame = first.session.range(1) + val secondDataFrame = second.session.range(1) + + firstDataFrame.persist() + secondDataFrame.persist() + SparkConnectService.sessionManager.closeSession(first.key) + + assert(second.session.sharedState.cacheManager.lookupCachedData(secondDataFrame).isDefined) + + SparkConnectService.sessionManager.closeSession(second.key) + assert(second.session.sharedState.cacheManager.lookupCachedData(secondDataFrame).isEmpty) + } + } + test("baseSession allows creating sessions after default session is cleared") { // Create a new session manager to test initialization val sessionManager = new SparkConnectSessionManager() diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/CacheManager.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/CacheManager.scala index 01fedf51602f6..6efc56b9f460e 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/CacheManager.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/CacheManager.scala @@ -17,6 +17,8 @@ package org.apache.spark.sql.execution +import java.util.IdentityHashMap + import scala.util.control.NonFatal import org.apache.hadoop.fs.{FileSystem, Path} @@ -79,13 +81,47 @@ class CacheManager extends Logging with AdaptiveSparkPlanHelper { @transient @volatile private var cachedData = IndexedSeq[CachedData]() + /** + * Tracks the sessions that explicitly cached each entry. A cache entry can be shared by + * multiple sessions because cache lookup uses plan semantics rather than session identity. + */ + @transient + private val cacheOwners = new IdentityHashMap[CachedData, Set[String]] + /** Clears all cached tables. */ def clearCache(): Unit = this.synchronized { cachedData.foreach(_.cachedRepresentation.cacheBuilder.clearCache()) cachedData = IndexedSeq[CachedData]() + cacheOwners.clear() CacheManager.logCacheOperation(log"Cleared all Dataframe cache entries") } + /** Clears cached data that is owned only by the given session. */ + private[sql] def clearCache(session: SparkSession): Unit = { + val sessionUUID = session.sessionUUID + val plansToUncache = this.synchronized { + val plans = cachedData.filter { cd => + Option(cacheOwners.get(cd)).contains(Set(sessionUUID)) + } + cachedData = cachedData.filterNot(cd => plans.exists(_ eq cd)) + val ownerEntries = cacheOwners.entrySet().iterator() + while (ownerEntries.hasNext) { + val entry = ownerEntries.next() + val remainingOwners = entry.getValue - sessionUUID + if (remainingOwners.nonEmpty) { + entry.setValue(remainingOwners) + } else { + ownerEntries.remove() + } + } + plans + } + plansToUncache.foreach(_.cachedRepresentation.cacheBuilder.clearCache()) + CacheManager.logCacheOperation( + log"Cleared ${MDC(SIZE, plansToUncache.size)} Dataframe cache entries for session " + + log"${MDC(SESSION_ID, sessionUUID)}") + } + /** Checks if the cache is empty. */ def isEmpty: Boolean = { cachedData.isEmpty @@ -144,7 +180,7 @@ class CacheManager extends Logging with AdaptiveSparkPlanHelper { log"Asked to cache a plan that is inapplicable for caching: " + log"${MDC(LOGICAL_PLAN, unnormalizedPlan)}" ) - } else if (lookupCachedDataInternal(normalizedPlan).nonEmpty) { + } else if (registerCacheOwner(normalizedPlan, spark.sessionUUID)) { logWarning("Asked to cache already cached data.") } else { val sessionWithConfigsOff = getOrCloneSessionWithConfigsOff(spark) @@ -158,12 +194,13 @@ class CacheManager extends Logging with AdaptiveSparkPlanHelper { } this.synchronized { - if (lookupCachedDataInternal(normalizedPlan).nonEmpty) { + if (registerCacheOwner(normalizedPlan, spark.sessionUUID)) { logWarning("Data has already been cached.") } else { // the cache key is the normalized plan val cd = CachedData(normalizedPlan, inMemoryRelation) cachedData = cd +: cachedData + cacheOwners.put(cd, Set(spark.sessionUUID)) CacheManager.logCacheOperation(log"Added Dataframe cache entry:" + log"${MDC(DATAFRAME_CACHE_ENTRY, cd)}") } @@ -171,6 +208,17 @@ class CacheManager extends Logging with AdaptiveSparkPlanHelper { } } + private def registerCacheOwner(plan: LogicalPlan, sessionUUID: String): Boolean = + this.synchronized { + lookupCachedDataInternal(plan) match { + case Some(cd) => + val owners = Option(cacheOwners.get(cd)).getOrElse(Set.empty) + cacheOwners.put(cd, owners + sessionUUID) + true + case None => false + } + } + /** * Un-cache the given plan or all the cache entries that refer to the given plan. * @@ -301,6 +349,7 @@ class CacheManager extends Logging with AdaptiveSparkPlanHelper { val plansToUncache = cachedData.filter(cd => shouldRemove(cd.plan)) this.synchronized { cachedData = cachedData.filterNot(cd => plansToUncache.exists(_ eq cd)) + plansToUncache.foreach(cd => cacheOwners.remove(cd)) } plansToUncache.foreach { _.cachedRepresentation.cacheBuilder.clearCache(blocking) } CacheManager.logCacheOperation(log"Removed ${MDC(SIZE, plansToUncache.size)} Dataframe " + @@ -374,14 +423,22 @@ class CacheManager extends Logging with AdaptiveSparkPlanHelper { } needToRecache.foreach { cd => cd.cachedRepresentation.cacheBuilder.clearCache() - tryRebuildCacheEntry(spark, cd).foreach { entry => - this.synchronized { - if (lookupCachedDataInternal(entry.plan).nonEmpty) { - logWarning("While recaching, data was already added to cache.") - } else { - cachedData = entry +: cachedData - CacheManager.logCacheOperation(log"Re-cached Dataframe cache entry:" + - log"${MDC(DATAFRAME_CACHE_ENTRY, entry)}") + val rebuiltEntry = tryRebuildCacheEntry(spark, cd) + this.synchronized { + val previousOwners = Option(cacheOwners.remove(cd)).getOrElse(Set.empty) + rebuiltEntry.foreach { entry => + if (previousOwners.nonEmpty) { + lookupCachedDataInternal(entry.plan) match { + case Some(existing) => + val existingOwners = Option(cacheOwners.get(existing)).getOrElse(Set.empty) + cacheOwners.put(existing, existingOwners ++ previousOwners) + logWarning("While recaching, data was already added to cache.") + case None => + cachedData = entry +: cachedData + cacheOwners.put(entry, previousOwners) + CacheManager.logCacheOperation(log"Re-cached Dataframe cache entry:" + + log"${MDC(DATAFRAME_CACHE_ENTRY, entry)}") + } } } } From e8b0cf38276e1e0099353727837467951278c24c Mon Sep 17 00:00:00 2001 From: PA Savalle Date: Mon, 31 Aug 2026 08:34:40 +0200 Subject: [PATCH 2/4] add argument to clearCache --- python/pyspark/sql/catalog.py | 17 +++-- python/pyspark/sql/connect/catalog.py | 4 +- python/pyspark/sql/connect/plan.py | 5 +- .../pyspark/sql/connect/proto/catalog_pb2.py | 68 +++++++++---------- .../pyspark/sql/connect/proto/catalog_pb2.pyi | 20 ++++++ python/pyspark/sql/tests/test_catalog.py | 25 +++++++ .../apache/spark/sql/catalog/Catalog.scala | 16 +++++ .../spark/sql/connect/CatalogSuite.scala | 24 +++++++ .../main/protobuf/spark/connect/catalog.proto | 5 +- .../apache/spark/sql/connect/Catalog.scala | 11 ++- .../connect/planner/SparkConnectPlanner.scala | 8 ++- .../apache/spark/sql/classic/Catalog.scala | 16 ++++- .../apache/spark/sql/CachedTableSuite.scala | 16 +++++ 13 files changed, 187 insertions(+), 48 deletions(-) diff --git a/python/pyspark/sql/catalog.py b/python/pyspark/sql/catalog.py index ed7d6ac7482b5..591c261090a30 100644 --- a/python/pyspark/sql/catalog.py +++ b/python/pyspark/sql/catalog.py @@ -1409,15 +1409,24 @@ def uncacheTable(self, tableName: str) -> None: """ self._jcatalog.uncacheTable(tableName) - def clearCache(self) -> None: - """Removes all cached tables from the in-memory cache. + def clearCache(self, allSessions: bool = True) -> None: + """Removes cached tables from the in-memory cache. .. versionadded:: 2.0.0 + Parameters + ---------- + allSessions : bool, optional + Whether to clear cached data across all sessions. If ``False``, only data cached by + the current session is cleared, while data also cached by another session is preserved. + Defaults to ``True``. + + .. versionadded:: 4.4.0 + Notes ----- Cached data is shared across all Spark sessions on the cluster, so clearing - the cache affects all sessions. + the cache affects all sessions by default. Examples -------- @@ -1428,7 +1437,7 @@ def clearCache(self) -> None: False >>> _ = spark.sql("DROP TABLE tbl1") """ - self._jcatalog.clearCache() + self._jcatalog.clearCache(allSessions) def refreshTable(self, tableName: str) -> None: """Invalidates and refreshes all the cached data and metadata of the given table. diff --git a/python/pyspark/sql/connect/catalog.py b/python/pyspark/sql/connect/catalog.py index 9859fef90caaf..5bcc8e5143b50 100644 --- a/python/pyspark/sql/connect/catalog.py +++ b/python/pyspark/sql/connect/catalog.py @@ -302,8 +302,8 @@ def uncacheTable(self, tableName: str) -> None: uncacheTable.__doc__ = PySparkCatalog.uncacheTable.__doc__ - def clearCache(self) -> None: - self._execute_and_fetch(plan.ClearCache()) + def clearCache(self, allSessions: bool = True) -> None: + self._execute_and_fetch(plan.ClearCache(all_sessions=allSessions)) clearCache.__doc__ = PySparkCatalog.clearCache.__doc__ diff --git a/python/pyspark/sql/connect/plan.py b/python/pyspark/sql/connect/plan.py index 297f20e2de054..f7b419f73ff02 100644 --- a/python/pyspark/sql/connect/plan.py +++ b/python/pyspark/sql/connect/plan.py @@ -2622,12 +2622,13 @@ def plan(self, session: "SparkConnectClient") -> proto.Relation: class ClearCache(LogicalPlan): - def __init__(self) -> None: + def __init__(self, all_sessions: bool = True) -> None: super().__init__(None) + self._all_sessions = all_sessions def plan(self, session: "SparkConnectClient") -> proto.Relation: plan = self._create_proto_relation() - plan.catalog.clear_cache.SetInParent() + plan.catalog.clear_cache.all_sessions = self._all_sessions return plan diff --git a/python/pyspark/sql/connect/proto/catalog_pb2.py b/python/pyspark/sql/connect/proto/catalog_pb2.py index c10a4bdf744a4..b55a9658d5a14 100644 --- a/python/pyspark/sql/connect/proto/catalog_pb2.py +++ b/python/pyspark/sql/connect/proto/catalog_pb2.py @@ -40,7 +40,7 @@ DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile( - b'\n\x1bspark/connect/catalog.proto\x12\rspark.connect\x1a\x1aspark/connect/common.proto\x1a\x19spark/connect/types.proto"\x8c\x14\n\x07\x43\x61talog\x12K\n\x10\x63urrent_database\x18\x01 \x01(\x0b\x32\x1e.spark.connect.CurrentDatabaseH\x00R\x0f\x63urrentDatabase\x12U\n\x14set_current_database\x18\x02 \x01(\x0b\x32!.spark.connect.SetCurrentDatabaseH\x00R\x12setCurrentDatabase\x12\x45\n\x0elist_databases\x18\x03 \x01(\x0b\x32\x1c.spark.connect.ListDatabasesH\x00R\rlistDatabases\x12<\n\x0blist_tables\x18\x04 \x01(\x0b\x32\x19.spark.connect.ListTablesH\x00R\nlistTables\x12\x45\n\x0elist_functions\x18\x05 \x01(\x0b\x32\x1c.spark.connect.ListFunctionsH\x00R\rlistFunctions\x12?\n\x0clist_columns\x18\x06 \x01(\x0b\x32\x1a.spark.connect.ListColumnsH\x00R\x0blistColumns\x12?\n\x0cget_database\x18\x07 \x01(\x0b\x32\x1a.spark.connect.GetDatabaseH\x00R\x0bgetDatabase\x12\x36\n\tget_table\x18\x08 \x01(\x0b\x32\x17.spark.connect.GetTableH\x00R\x08getTable\x12?\n\x0cget_function\x18\t \x01(\x0b\x32\x1a.spark.connect.GetFunctionH\x00R\x0bgetFunction\x12H\n\x0f\x64\x61tabase_exists\x18\n \x01(\x0b\x32\x1d.spark.connect.DatabaseExistsH\x00R\x0e\x64\x61tabaseExists\x12?\n\x0ctable_exists\x18\x0b \x01(\x0b\x32\x1a.spark.connect.TableExistsH\x00R\x0btableExists\x12H\n\x0f\x66unction_exists\x18\x0c \x01(\x0b\x32\x1d.spark.connect.FunctionExistsH\x00R\x0e\x66unctionExists\x12X\n\x15\x63reate_external_table\x18\r \x01(\x0b\x32".spark.connect.CreateExternalTableH\x00R\x13\x63reateExternalTable\x12?\n\x0c\x63reate_table\x18\x0e \x01(\x0b\x32\x1a.spark.connect.CreateTableH\x00R\x0b\x63reateTable\x12\x43\n\x0e\x64rop_temp_view\x18\x0f \x01(\x0b\x32\x1b.spark.connect.DropTempViewH\x00R\x0c\x64ropTempView\x12V\n\x15\x64rop_global_temp_view\x18\x10 \x01(\x0b\x32!.spark.connect.DropGlobalTempViewH\x00R\x12\x64ropGlobalTempView\x12Q\n\x12recover_partitions\x18\x11 \x01(\x0b\x32 .spark.connect.RecoverPartitionsH\x00R\x11recoverPartitions\x12\x36\n\tis_cached\x18\x12 \x01(\x0b\x32\x17.spark.connect.IsCachedH\x00R\x08isCached\x12<\n\x0b\x63\x61\x63he_table\x18\x13 \x01(\x0b\x32\x19.spark.connect.CacheTableH\x00R\ncacheTable\x12\x42\n\runcache_table\x18\x14 \x01(\x0b\x32\x1b.spark.connect.UncacheTableH\x00R\x0cuncacheTable\x12<\n\x0b\x63lear_cache\x18\x15 \x01(\x0b\x32\x19.spark.connect.ClearCacheH\x00R\nclearCache\x12\x42\n\rrefresh_table\x18\x16 \x01(\x0b\x32\x1b.spark.connect.RefreshTableH\x00R\x0crefreshTable\x12\x46\n\x0frefresh_by_path\x18\x17 \x01(\x0b\x32\x1c.spark.connect.RefreshByPathH\x00R\rrefreshByPath\x12H\n\x0f\x63urrent_catalog\x18\x18 \x01(\x0b\x32\x1d.spark.connect.CurrentCatalogH\x00R\x0e\x63urrentCatalog\x12R\n\x13set_current_catalog\x18\x19 \x01(\x0b\x32 .spark.connect.SetCurrentCatalogH\x00R\x11setCurrentCatalog\x12\x42\n\rlist_catalogs\x18\x1a \x01(\x0b\x32\x1b.spark.connect.ListCatalogsH\x00R\x0clistCatalogs\x12\x39\n\ndrop_table\x18\x1b \x01(\x0b\x32\x18.spark.connect.DropTableH\x00R\tdropTable\x12\x36\n\tdrop_view\x18\x1c \x01(\x0b\x32\x17.spark.connect.DropViewH\x00R\x08\x64ropView\x12H\n\x0f\x63reate_database\x18\x1d \x01(\x0b\x32\x1d.spark.connect.CreateDatabaseH\x00R\x0e\x63reateDatabase\x12\x42\n\rdrop_database\x18\x1e \x01(\x0b\x32\x1b.spark.connect.DropDatabaseH\x00R\x0c\x64ropDatabase\x12H\n\x0flist_partitions\x18\x1f \x01(\x0b\x32\x1d.spark.connect.ListPartitionsH\x00R\x0elistPartitions\x12\x39\n\nlist_views\x18 \x01(\x0b\x32\x18.spark.connect.ListViewsH\x00R\tlistViews\x12U\n\x14get_table_properties\x18! \x01(\x0b\x32!.spark.connect.GetTablePropertiesH\x00R\x12getTableProperties\x12\\\n\x17get_create_table_string\x18" \x01(\x0b\x32#.spark.connect.GetCreateTableStringH\x00R\x14getCreateTableString\x12\x45\n\x0etruncate_table\x18# \x01(\x0b\x32\x1c.spark.connect.TruncateTableH\x00R\rtruncateTable\x12\x42\n\ranalyze_table\x18$ \x01(\x0b\x32\x1b.spark.connect.AnalyzeTableH\x00R\x0c\x61nalyzeTableB\n\n\x08\x63\x61t_type"\x11\n\x0f\x43urrentDatabase"-\n\x12SetCurrentDatabase\x12\x17\n\x07\x64\x62_name\x18\x01 \x01(\tR\x06\x64\x62Name":\n\rListDatabases\x12\x1d\n\x07pattern\x18\x01 \x01(\tH\x00R\x07pattern\x88\x01\x01\x42\n\n\x08_pattern"a\n\nListTables\x12\x1c\n\x07\x64\x62_name\x18\x01 \x01(\tH\x00R\x06\x64\x62Name\x88\x01\x01\x12\x1d\n\x07pattern\x18\x02 \x01(\tH\x01R\x07pattern\x88\x01\x01\x42\n\n\x08_db_nameB\n\n\x08_pattern"d\n\rListFunctions\x12\x1c\n\x07\x64\x62_name\x18\x01 \x01(\tH\x00R\x06\x64\x62Name\x88\x01\x01\x12\x1d\n\x07pattern\x18\x02 \x01(\tH\x01R\x07pattern\x88\x01\x01\x42\n\n\x08_db_nameB\n\n\x08_pattern"V\n\x0bListColumns\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName\x12\x1c\n\x07\x64\x62_name\x18\x02 \x01(\tH\x00R\x06\x64\x62Name\x88\x01\x01\x42\n\n\x08_db_name"&\n\x0bGetDatabase\x12\x17\n\x07\x64\x62_name\x18\x01 \x01(\tR\x06\x64\x62Name"S\n\x08GetTable\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName\x12\x1c\n\x07\x64\x62_name\x18\x02 \x01(\tH\x00R\x06\x64\x62Name\x88\x01\x01\x42\n\n\x08_db_name"\\\n\x0bGetFunction\x12#\n\rfunction_name\x18\x01 \x01(\tR\x0c\x66unctionName\x12\x1c\n\x07\x64\x62_name\x18\x02 \x01(\tH\x00R\x06\x64\x62Name\x88\x01\x01\x42\n\n\x08_db_name")\n\x0e\x44\x61tabaseExists\x12\x17\n\x07\x64\x62_name\x18\x01 \x01(\tR\x06\x64\x62Name"V\n\x0bTableExists\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName\x12\x1c\n\x07\x64\x62_name\x18\x02 \x01(\tH\x00R\x06\x64\x62Name\x88\x01\x01\x42\n\n\x08_db_name"_\n\x0e\x46unctionExists\x12#\n\rfunction_name\x18\x01 \x01(\tR\x0c\x66unctionName\x12\x1c\n\x07\x64\x62_name\x18\x02 \x01(\tH\x00R\x06\x64\x62Name\x88\x01\x01\x42\n\n\x08_db_name"\xc6\x02\n\x13\x43reateExternalTable\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName\x12\x17\n\x04path\x18\x02 \x01(\tH\x00R\x04path\x88\x01\x01\x12\x1b\n\x06source\x18\x03 \x01(\tH\x01R\x06source\x88\x01\x01\x12\x34\n\x06schema\x18\x04 \x01(\x0b\x32\x17.spark.connect.DataTypeH\x02R\x06schema\x88\x01\x01\x12I\n\x07options\x18\x05 \x03(\x0b\x32/.spark.connect.CreateExternalTable.OptionsEntryR\x07options\x1a:\n\x0cOptionsEntry\x12\x10\n\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n\x05value\x18\x02 \x01(\tR\x05value:\x02\x38\x01\x42\x07\n\x05_pathB\t\n\x07_sourceB\t\n\x07_schema"\xed\x02\n\x0b\x43reateTable\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName\x12\x17\n\x04path\x18\x02 \x01(\tH\x00R\x04path\x88\x01\x01\x12\x1b\n\x06source\x18\x03 \x01(\tH\x01R\x06source\x88\x01\x01\x12%\n\x0b\x64\x65scription\x18\x04 \x01(\tH\x02R\x0b\x64\x65scription\x88\x01\x01\x12\x34\n\x06schema\x18\x05 \x01(\x0b\x32\x17.spark.connect.DataTypeH\x03R\x06schema\x88\x01\x01\x12\x41\n\x07options\x18\x06 \x03(\x0b\x32\'.spark.connect.CreateTable.OptionsEntryR\x07options\x1a:\n\x0cOptionsEntry\x12\x10\n\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n\x05value\x18\x02 \x01(\tR\x05value:\x02\x38\x01\x42\x07\n\x05_pathB\t\n\x07_sourceB\x0e\n\x0c_descriptionB\t\n\x07_schema"+\n\x0c\x44ropTempView\x12\x1b\n\tview_name\x18\x01 \x01(\tR\x08viewName"1\n\x12\x44ropGlobalTempView\x12\x1b\n\tview_name\x18\x01 \x01(\tR\x08viewName"2\n\x11RecoverPartitions\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName")\n\x08IsCached\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName"\x84\x01\n\nCacheTable\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName\x12\x45\n\rstorage_level\x18\x02 \x01(\x0b\x32\x1b.spark.connect.StorageLevelH\x00R\x0cstorageLevel\x88\x01\x01\x42\x10\n\x0e_storage_level"-\n\x0cUncacheTable\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName"\x0c\n\nClearCache"-\n\x0cRefreshTable\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName"#\n\rRefreshByPath\x12\x12\n\x04path\x18\x01 \x01(\tR\x04path"\x10\n\x0e\x43urrentCatalog"6\n\x11SetCurrentCatalog\x12!\n\x0c\x63\x61talog_name\x18\x01 \x01(\tR\x0b\x63\x61talogName"9\n\x0cListCatalogs\x12\x1d\n\x07pattern\x18\x01 \x01(\tH\x00R\x07pattern\x88\x01\x01\x42\n\n\x08_pattern"]\n\tDropTable\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName\x12\x1b\n\tif_exists\x18\x02 \x01(\x08R\x08ifExists\x12\x14\n\x05purge\x18\x03 \x01(\x08R\x05purge"D\n\x08\x44ropView\x12\x1b\n\tview_name\x18\x01 \x01(\tR\x08viewName\x12\x1b\n\tif_exists\x18\x02 \x01(\x08R\x08ifExists"\xdb\x01\n\x0e\x43reateDatabase\x12\x17\n\x07\x64\x62_name\x18\x01 \x01(\tR\x06\x64\x62Name\x12"\n\rif_not_exists\x18\x02 \x01(\x08R\x0bifNotExists\x12M\n\nproperties\x18\x03 \x03(\x0b\x32-.spark.connect.CreateDatabase.PropertiesEntryR\nproperties\x1a=\n\x0fPropertiesEntry\x12\x10\n\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n\x05value\x18\x02 \x01(\tR\x05value:\x02\x38\x01"^\n\x0c\x44ropDatabase\x12\x17\n\x07\x64\x62_name\x18\x01 \x01(\tR\x06\x64\x62Name\x12\x1b\n\tif_exists\x18\x02 \x01(\x08R\x08ifExists\x12\x18\n\x07\x63\x61scade\x18\x03 \x01(\x08R\x07\x63\x61scade"/\n\x0eListPartitions\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName"`\n\tListViews\x12\x1c\n\x07\x64\x62_name\x18\x01 \x01(\tH\x00R\x06\x64\x62Name\x88\x01\x01\x12\x1d\n\x07pattern\x18\x02 \x01(\tH\x01R\x07pattern\x88\x01\x01\x42\n\n\x08_db_nameB\n\n\x08_pattern"3\n\x12GetTableProperties\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName"P\n\x14GetCreateTableString\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName\x12\x19\n\x08\x61s_serde\x18\x02 \x01(\x08R\x07\x61sSerde".\n\rTruncateTable\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName"F\n\x0c\x41nalyzeTable\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName\x12\x17\n\x07no_scan\x18\x02 \x01(\x08R\x06noScanB6\n\x1eorg.apache.spark.connect.protoP\x01Z\x12internal/generatedb\x06proto3' + b'\n\x1bspark/connect/catalog.proto\x12\rspark.connect\x1a\x1aspark/connect/common.proto\x1a\x19spark/connect/types.proto"\x8c\x14\n\x07\x43\x61talog\x12K\n\x10\x63urrent_database\x18\x01 \x01(\x0b\x32\x1e.spark.connect.CurrentDatabaseH\x00R\x0f\x63urrentDatabase\x12U\n\x14set_current_database\x18\x02 \x01(\x0b\x32!.spark.connect.SetCurrentDatabaseH\x00R\x12setCurrentDatabase\x12\x45\n\x0elist_databases\x18\x03 \x01(\x0b\x32\x1c.spark.connect.ListDatabasesH\x00R\rlistDatabases\x12<\n\x0blist_tables\x18\x04 \x01(\x0b\x32\x19.spark.connect.ListTablesH\x00R\nlistTables\x12\x45\n\x0elist_functions\x18\x05 \x01(\x0b\x32\x1c.spark.connect.ListFunctionsH\x00R\rlistFunctions\x12?\n\x0clist_columns\x18\x06 \x01(\x0b\x32\x1a.spark.connect.ListColumnsH\x00R\x0blistColumns\x12?\n\x0cget_database\x18\x07 \x01(\x0b\x32\x1a.spark.connect.GetDatabaseH\x00R\x0bgetDatabase\x12\x36\n\tget_table\x18\x08 \x01(\x0b\x32\x17.spark.connect.GetTableH\x00R\x08getTable\x12?\n\x0cget_function\x18\t \x01(\x0b\x32\x1a.spark.connect.GetFunctionH\x00R\x0bgetFunction\x12H\n\x0f\x64\x61tabase_exists\x18\n \x01(\x0b\x32\x1d.spark.connect.DatabaseExistsH\x00R\x0e\x64\x61tabaseExists\x12?\n\x0ctable_exists\x18\x0b \x01(\x0b\x32\x1a.spark.connect.TableExistsH\x00R\x0btableExists\x12H\n\x0f\x66unction_exists\x18\x0c \x01(\x0b\x32\x1d.spark.connect.FunctionExistsH\x00R\x0e\x66unctionExists\x12X\n\x15\x63reate_external_table\x18\r \x01(\x0b\x32".spark.connect.CreateExternalTableH\x00R\x13\x63reateExternalTable\x12?\n\x0c\x63reate_table\x18\x0e \x01(\x0b\x32\x1a.spark.connect.CreateTableH\x00R\x0b\x63reateTable\x12\x43\n\x0e\x64rop_temp_view\x18\x0f \x01(\x0b\x32\x1b.spark.connect.DropTempViewH\x00R\x0c\x64ropTempView\x12V\n\x15\x64rop_global_temp_view\x18\x10 \x01(\x0b\x32!.spark.connect.DropGlobalTempViewH\x00R\x12\x64ropGlobalTempView\x12Q\n\x12recover_partitions\x18\x11 \x01(\x0b\x32 .spark.connect.RecoverPartitionsH\x00R\x11recoverPartitions\x12\x36\n\tis_cached\x18\x12 \x01(\x0b\x32\x17.spark.connect.IsCachedH\x00R\x08isCached\x12<\n\x0b\x63\x61\x63he_table\x18\x13 \x01(\x0b\x32\x19.spark.connect.CacheTableH\x00R\ncacheTable\x12\x42\n\runcache_table\x18\x14 \x01(\x0b\x32\x1b.spark.connect.UncacheTableH\x00R\x0cuncacheTable\x12<\n\x0b\x63lear_cache\x18\x15 \x01(\x0b\x32\x19.spark.connect.ClearCacheH\x00R\nclearCache\x12\x42\n\rrefresh_table\x18\x16 \x01(\x0b\x32\x1b.spark.connect.RefreshTableH\x00R\x0crefreshTable\x12\x46\n\x0frefresh_by_path\x18\x17 \x01(\x0b\x32\x1c.spark.connect.RefreshByPathH\x00R\rrefreshByPath\x12H\n\x0f\x63urrent_catalog\x18\x18 \x01(\x0b\x32\x1d.spark.connect.CurrentCatalogH\x00R\x0e\x63urrentCatalog\x12R\n\x13set_current_catalog\x18\x19 \x01(\x0b\x32 .spark.connect.SetCurrentCatalogH\x00R\x11setCurrentCatalog\x12\x42\n\rlist_catalogs\x18\x1a \x01(\x0b\x32\x1b.spark.connect.ListCatalogsH\x00R\x0clistCatalogs\x12\x39\n\ndrop_table\x18\x1b \x01(\x0b\x32\x18.spark.connect.DropTableH\x00R\tdropTable\x12\x36\n\tdrop_view\x18\x1c \x01(\x0b\x32\x17.spark.connect.DropViewH\x00R\x08\x64ropView\x12H\n\x0f\x63reate_database\x18\x1d \x01(\x0b\x32\x1d.spark.connect.CreateDatabaseH\x00R\x0e\x63reateDatabase\x12\x42\n\rdrop_database\x18\x1e \x01(\x0b\x32\x1b.spark.connect.DropDatabaseH\x00R\x0c\x64ropDatabase\x12H\n\x0flist_partitions\x18\x1f \x01(\x0b\x32\x1d.spark.connect.ListPartitionsH\x00R\x0elistPartitions\x12\x39\n\nlist_views\x18 \x01(\x0b\x32\x18.spark.connect.ListViewsH\x00R\tlistViews\x12U\n\x14get_table_properties\x18! \x01(\x0b\x32!.spark.connect.GetTablePropertiesH\x00R\x12getTableProperties\x12\\\n\x17get_create_table_string\x18" \x01(\x0b\x32#.spark.connect.GetCreateTableStringH\x00R\x14getCreateTableString\x12\x45\n\x0etruncate_table\x18# \x01(\x0b\x32\x1c.spark.connect.TruncateTableH\x00R\rtruncateTable\x12\x42\n\ranalyze_table\x18$ \x01(\x0b\x32\x1b.spark.connect.AnalyzeTableH\x00R\x0c\x61nalyzeTableB\n\n\x08\x63\x61t_type"\x11\n\x0f\x43urrentDatabase"-\n\x12SetCurrentDatabase\x12\x17\n\x07\x64\x62_name\x18\x01 \x01(\tR\x06\x64\x62Name":\n\rListDatabases\x12\x1d\n\x07pattern\x18\x01 \x01(\tH\x00R\x07pattern\x88\x01\x01\x42\n\n\x08_pattern"a\n\nListTables\x12\x1c\n\x07\x64\x62_name\x18\x01 \x01(\tH\x00R\x06\x64\x62Name\x88\x01\x01\x12\x1d\n\x07pattern\x18\x02 \x01(\tH\x01R\x07pattern\x88\x01\x01\x42\n\n\x08_db_nameB\n\n\x08_pattern"d\n\rListFunctions\x12\x1c\n\x07\x64\x62_name\x18\x01 \x01(\tH\x00R\x06\x64\x62Name\x88\x01\x01\x12\x1d\n\x07pattern\x18\x02 \x01(\tH\x01R\x07pattern\x88\x01\x01\x42\n\n\x08_db_nameB\n\n\x08_pattern"V\n\x0bListColumns\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName\x12\x1c\n\x07\x64\x62_name\x18\x02 \x01(\tH\x00R\x06\x64\x62Name\x88\x01\x01\x42\n\n\x08_db_name"&\n\x0bGetDatabase\x12\x17\n\x07\x64\x62_name\x18\x01 \x01(\tR\x06\x64\x62Name"S\n\x08GetTable\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName\x12\x1c\n\x07\x64\x62_name\x18\x02 \x01(\tH\x00R\x06\x64\x62Name\x88\x01\x01\x42\n\n\x08_db_name"\\\n\x0bGetFunction\x12#\n\rfunction_name\x18\x01 \x01(\tR\x0c\x66unctionName\x12\x1c\n\x07\x64\x62_name\x18\x02 \x01(\tH\x00R\x06\x64\x62Name\x88\x01\x01\x42\n\n\x08_db_name")\n\x0e\x44\x61tabaseExists\x12\x17\n\x07\x64\x62_name\x18\x01 \x01(\tR\x06\x64\x62Name"V\n\x0bTableExists\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName\x12\x1c\n\x07\x64\x62_name\x18\x02 \x01(\tH\x00R\x06\x64\x62Name\x88\x01\x01\x42\n\n\x08_db_name"_\n\x0e\x46unctionExists\x12#\n\rfunction_name\x18\x01 \x01(\tR\x0c\x66unctionName\x12\x1c\n\x07\x64\x62_name\x18\x02 \x01(\tH\x00R\x06\x64\x62Name\x88\x01\x01\x42\n\n\x08_db_name"\xc6\x02\n\x13\x43reateExternalTable\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName\x12\x17\n\x04path\x18\x02 \x01(\tH\x00R\x04path\x88\x01\x01\x12\x1b\n\x06source\x18\x03 \x01(\tH\x01R\x06source\x88\x01\x01\x12\x34\n\x06schema\x18\x04 \x01(\x0b\x32\x17.spark.connect.DataTypeH\x02R\x06schema\x88\x01\x01\x12I\n\x07options\x18\x05 \x03(\x0b\x32/.spark.connect.CreateExternalTable.OptionsEntryR\x07options\x1a:\n\x0cOptionsEntry\x12\x10\n\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n\x05value\x18\x02 \x01(\tR\x05value:\x02\x38\x01\x42\x07\n\x05_pathB\t\n\x07_sourceB\t\n\x07_schema"\xed\x02\n\x0b\x43reateTable\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName\x12\x17\n\x04path\x18\x02 \x01(\tH\x00R\x04path\x88\x01\x01\x12\x1b\n\x06source\x18\x03 \x01(\tH\x01R\x06source\x88\x01\x01\x12%\n\x0b\x64\x65scription\x18\x04 \x01(\tH\x02R\x0b\x64\x65scription\x88\x01\x01\x12\x34\n\x06schema\x18\x05 \x01(\x0b\x32\x17.spark.connect.DataTypeH\x03R\x06schema\x88\x01\x01\x12\x41\n\x07options\x18\x06 \x03(\x0b\x32\'.spark.connect.CreateTable.OptionsEntryR\x07options\x1a:\n\x0cOptionsEntry\x12\x10\n\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n\x05value\x18\x02 \x01(\tR\x05value:\x02\x38\x01\x42\x07\n\x05_pathB\t\n\x07_sourceB\x0e\n\x0c_descriptionB\t\n\x07_schema"+\n\x0c\x44ropTempView\x12\x1b\n\tview_name\x18\x01 \x01(\tR\x08viewName"1\n\x12\x44ropGlobalTempView\x12\x1b\n\tview_name\x18\x01 \x01(\tR\x08viewName"2\n\x11RecoverPartitions\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName")\n\x08IsCached\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName"\x84\x01\n\nCacheTable\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName\x12\x45\n\rstorage_level\x18\x02 \x01(\x0b\x32\x1b.spark.connect.StorageLevelH\x00R\x0cstorageLevel\x88\x01\x01\x42\x10\n\x0e_storage_level"-\n\x0cUncacheTable\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName"E\n\nClearCache\x12&\n\x0c\x61ll_sessions\x18\x01 \x01(\x08H\x00R\x0b\x61llSessions\x88\x01\x01\x42\x0f\n\r_all_sessions"-\n\x0cRefreshTable\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName"#\n\rRefreshByPath\x12\x12\n\x04path\x18\x01 \x01(\tR\x04path"\x10\n\x0e\x43urrentCatalog"6\n\x11SetCurrentCatalog\x12!\n\x0c\x63\x61talog_name\x18\x01 \x01(\tR\x0b\x63\x61talogName"9\n\x0cListCatalogs\x12\x1d\n\x07pattern\x18\x01 \x01(\tH\x00R\x07pattern\x88\x01\x01\x42\n\n\x08_pattern"]\n\tDropTable\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName\x12\x1b\n\tif_exists\x18\x02 \x01(\x08R\x08ifExists\x12\x14\n\x05purge\x18\x03 \x01(\x08R\x05purge"D\n\x08\x44ropView\x12\x1b\n\tview_name\x18\x01 \x01(\tR\x08viewName\x12\x1b\n\tif_exists\x18\x02 \x01(\x08R\x08ifExists"\xdb\x01\n\x0e\x43reateDatabase\x12\x17\n\x07\x64\x62_name\x18\x01 \x01(\tR\x06\x64\x62Name\x12"\n\rif_not_exists\x18\x02 \x01(\x08R\x0bifNotExists\x12M\n\nproperties\x18\x03 \x03(\x0b\x32-.spark.connect.CreateDatabase.PropertiesEntryR\nproperties\x1a=\n\x0fPropertiesEntry\x12\x10\n\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n\x05value\x18\x02 \x01(\tR\x05value:\x02\x38\x01"^\n\x0c\x44ropDatabase\x12\x17\n\x07\x64\x62_name\x18\x01 \x01(\tR\x06\x64\x62Name\x12\x1b\n\tif_exists\x18\x02 \x01(\x08R\x08ifExists\x12\x18\n\x07\x63\x61scade\x18\x03 \x01(\x08R\x07\x63\x61scade"/\n\x0eListPartitions\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName"`\n\tListViews\x12\x1c\n\x07\x64\x62_name\x18\x01 \x01(\tH\x00R\x06\x64\x62Name\x88\x01\x01\x12\x1d\n\x07pattern\x18\x02 \x01(\tH\x01R\x07pattern\x88\x01\x01\x42\n\n\x08_db_nameB\n\n\x08_pattern"3\n\x12GetTableProperties\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName"P\n\x14GetCreateTableString\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName\x12\x19\n\x08\x61s_serde\x18\x02 \x01(\x08R\x07\x61sSerde".\n\rTruncateTable\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName"F\n\x0c\x41nalyzeTable\x12\x1d\n\ntable_name\x18\x01 \x01(\tR\ttableName\x12\x17\n\x07no_scan\x18\x02 \x01(\x08R\x06noScanB6\n\x1eorg.apache.spark.connect.protoP\x01Z\x12internal/generatedb\x06proto3' ) _globals = globals() @@ -106,37 +106,37 @@ _globals["_UNCACHETABLE"]._serialized_start = 4561 _globals["_UNCACHETABLE"]._serialized_end = 4606 _globals["_CLEARCACHE"]._serialized_start = 4608 - _globals["_CLEARCACHE"]._serialized_end = 4620 - _globals["_REFRESHTABLE"]._serialized_start = 4622 - _globals["_REFRESHTABLE"]._serialized_end = 4667 - _globals["_REFRESHBYPATH"]._serialized_start = 4669 - _globals["_REFRESHBYPATH"]._serialized_end = 4704 - _globals["_CURRENTCATALOG"]._serialized_start = 4706 - _globals["_CURRENTCATALOG"]._serialized_end = 4722 - _globals["_SETCURRENTCATALOG"]._serialized_start = 4724 - _globals["_SETCURRENTCATALOG"]._serialized_end = 4778 - _globals["_LISTCATALOGS"]._serialized_start = 4780 - _globals["_LISTCATALOGS"]._serialized_end = 4837 - _globals["_DROPTABLE"]._serialized_start = 4839 - _globals["_DROPTABLE"]._serialized_end = 4932 - _globals["_DROPVIEW"]._serialized_start = 4934 - _globals["_DROPVIEW"]._serialized_end = 5002 - _globals["_CREATEDATABASE"]._serialized_start = 5005 - _globals["_CREATEDATABASE"]._serialized_end = 5224 - _globals["_CREATEDATABASE_PROPERTIESENTRY"]._serialized_start = 5163 - _globals["_CREATEDATABASE_PROPERTIESENTRY"]._serialized_end = 5224 - _globals["_DROPDATABASE"]._serialized_start = 5226 - _globals["_DROPDATABASE"]._serialized_end = 5320 - _globals["_LISTPARTITIONS"]._serialized_start = 5322 - _globals["_LISTPARTITIONS"]._serialized_end = 5369 - _globals["_LISTVIEWS"]._serialized_start = 5371 - _globals["_LISTVIEWS"]._serialized_end = 5467 - _globals["_GETTABLEPROPERTIES"]._serialized_start = 5469 - _globals["_GETTABLEPROPERTIES"]._serialized_end = 5520 - _globals["_GETCREATETABLESTRING"]._serialized_start = 5522 - _globals["_GETCREATETABLESTRING"]._serialized_end = 5602 - _globals["_TRUNCATETABLE"]._serialized_start = 5604 - _globals["_TRUNCATETABLE"]._serialized_end = 5650 - _globals["_ANALYZETABLE"]._serialized_start = 5652 - _globals["_ANALYZETABLE"]._serialized_end = 5722 + _globals["_CLEARCACHE"]._serialized_end = 4677 + _globals["_REFRESHTABLE"]._serialized_start = 4679 + _globals["_REFRESHTABLE"]._serialized_end = 4724 + _globals["_REFRESHBYPATH"]._serialized_start = 4726 + _globals["_REFRESHBYPATH"]._serialized_end = 4761 + _globals["_CURRENTCATALOG"]._serialized_start = 4763 + _globals["_CURRENTCATALOG"]._serialized_end = 4779 + _globals["_SETCURRENTCATALOG"]._serialized_start = 4781 + _globals["_SETCURRENTCATALOG"]._serialized_end = 4835 + _globals["_LISTCATALOGS"]._serialized_start = 4837 + _globals["_LISTCATALOGS"]._serialized_end = 4894 + _globals["_DROPTABLE"]._serialized_start = 4896 + _globals["_DROPTABLE"]._serialized_end = 4989 + _globals["_DROPVIEW"]._serialized_start = 4991 + _globals["_DROPVIEW"]._serialized_end = 5059 + _globals["_CREATEDATABASE"]._serialized_start = 5062 + _globals["_CREATEDATABASE"]._serialized_end = 5281 + _globals["_CREATEDATABASE_PROPERTIESENTRY"]._serialized_start = 5220 + _globals["_CREATEDATABASE_PROPERTIESENTRY"]._serialized_end = 5281 + _globals["_DROPDATABASE"]._serialized_start = 5283 + _globals["_DROPDATABASE"]._serialized_end = 5377 + _globals["_LISTPARTITIONS"]._serialized_start = 5379 + _globals["_LISTPARTITIONS"]._serialized_end = 5426 + _globals["_LISTVIEWS"]._serialized_start = 5428 + _globals["_LISTVIEWS"]._serialized_end = 5524 + _globals["_GETTABLEPROPERTIES"]._serialized_start = 5526 + _globals["_GETTABLEPROPERTIES"]._serialized_end = 5577 + _globals["_GETCREATETABLESTRING"]._serialized_start = 5579 + _globals["_GETCREATETABLESTRING"]._serialized_end = 5659 + _globals["_TRUNCATETABLE"]._serialized_start = 5661 + _globals["_TRUNCATETABLE"]._serialized_end = 5707 + _globals["_ANALYZETABLE"]._serialized_start = 5709 + _globals["_ANALYZETABLE"]._serialized_end = 5779 # @@protoc_insertion_point(module_scope) diff --git a/python/pyspark/sql/connect/proto/catalog_pb2.pyi b/python/pyspark/sql/connect/proto/catalog_pb2.pyi index 72c1c9c6f1c83..1ced6fd3be02c 100644 --- a/python/pyspark/sql/connect/proto/catalog_pb2.pyi +++ b/python/pyspark/sql/connect/proto/catalog_pb2.pyi @@ -1126,9 +1126,29 @@ class ClearCache(google.protobuf.message.Message): DESCRIPTOR: google.protobuf.descriptor.Descriptor + ALL_SESSIONS_FIELD_NUMBER: builtins.int + all_sessions: builtins.bool + """(Optional) Whether to clear cached data across all sessions. Defaults to true when omitted.""" def __init__( self, + *, + all_sessions: builtins.bool | None = ..., + ) -> None: ... + def HasField( + self, + field_name: typing_extensions.Literal[ + "_all_sessions", b"_all_sessions", "all_sessions", b"all_sessions" + ], + ) -> builtins.bool: ... + def ClearField( + self, + field_name: typing_extensions.Literal[ + "_all_sessions", b"_all_sessions", "all_sessions", b"all_sessions" + ], ) -> None: ... + def WhichOneof( + self, oneof_group: typing_extensions.Literal["_all_sessions", b"_all_sessions"] + ) -> typing_extensions.Literal["all_sessions"] | None: ... global___ClearCache = ClearCache diff --git a/python/pyspark/sql/tests/test_catalog.py b/python/pyspark/sql/tests/test_catalog.py index d813384a3a747..e6b41f0897c08 100644 --- a/python/pyspark/sql/tests/test_catalog.py +++ b/python/pyspark/sql/tests/test_catalog.py @@ -449,6 +449,31 @@ def assert_cached(c: bool): spark.catalog.clearCache() assert_cached(False) + def test_clear_cache_current_session(self): + # SPARK-50569: clearCache can preserve data cached by other sessions. + spark = self.spark + other_session = spark.newSession() + first_view = "clear_cache_first_view" + second_view = "clear_cache_second_view" + try: + spark.range(1).createTempView(first_view) + other_session.range(1, 2).createTempView(second_view) + spark.catalog.cacheTable(first_view) + other_session.catalog.cacheTable(second_view) + + spark.catalog.clearCache(allSessions=False) + + self.assertFalse(spark.catalog.isCached(first_view)) + self.assertTrue(other_session.catalog.isCached(second_view)) + other_session.catalog.clearCache(allSessions=False) + self.assertFalse(other_session.catalog.isCached(second_view)) + finally: + spark.catalog.clearCache() + spark.catalog.dropTempView(first_view) + other_session.catalog.dropTempView(second_view) + if hasattr(other_session, "client"): + other_session.client.close() + def test_table_exists(self): # SPARK-36176: testing that table_exists returns correct boolean spark = self.spark diff --git a/sql/api/src/main/scala/org/apache/spark/sql/catalog/Catalog.scala b/sql/api/src/main/scala/org/apache/spark/sql/catalog/Catalog.scala index 1bb88be33d8e8..03c0bdd662daa 100644 --- a/sql/api/src/main/scala/org/apache/spark/sql/catalog/Catalog.scala +++ b/sql/api/src/main/scala/org/apache/spark/sql/catalog/Catalog.scala @@ -645,6 +645,22 @@ abstract class Catalog { */ def clearCache(): Unit + /** + * Removes cached tables from the in-memory cache. + * + * @param allSessions + * when true, removes cached data across all Spark sessions. When false, removes cached data + * owned by the current session while preserving data that is also cached by another session. + * @since 4.4.0 + */ + def clearCache(allSessions: Boolean): Unit = { + if (allSessions) { + clearCache() + } else { + catalogUnsupported("clearCache") + } + } + /** * Invalidates and refreshes all the cached data and metadata of the given table. For * performance reasons, Spark SQL or the external data source library it uses might cache diff --git a/sql/connect/client/jvm/src/test/scala/org/apache/spark/sql/connect/CatalogSuite.scala b/sql/connect/client/jvm/src/test/scala/org/apache/spark/sql/connect/CatalogSuite.scala index 5237554b3625d..518505346c5e9 100644 --- a/sql/connect/client/jvm/src/test/scala/org/apache/spark/sql/connect/CatalogSuite.scala +++ b/sql/connect/client/jvm/src/test/scala/org/apache/spark/sql/connect/CatalogSuite.scala @@ -209,6 +209,30 @@ class CatalogSuite extends ConnectFunSuite with RemoteSparkSession with SQLHelpe } } + test("SPARK-50569: clearCache can be scoped to the current session") { + val otherSession = spark.newSession() + val firstView = "clear_cache_first_view" + val secondView = "clear_cache_second_view" + try { + spark.range(1).createTempView(firstView) + otherSession.range(1, 2).createTempView(secondView) + spark.catalog.cacheTable(firstView) + otherSession.catalog.cacheTable(secondView) + + spark.catalog.clearCache(allSessions = false) + + assert(!spark.catalog.isCached(firstView)) + assert(otherSession.catalog.isCached(secondView)) + otherSession.catalog.clearCache(allSessions = false) + assert(!otherSession.catalog.isCached(secondView)) + } finally { + spark.catalog.clearCache() + spark.catalog.dropTempView(firstView) + otherSession.catalog.dropTempView(secondView) + otherSession.close() + } + } + test("TempView APIs") { val viewName = "view1" val globalViewName = "g_view1" diff --git a/sql/connect/common/src/main/protobuf/spark/connect/catalog.proto b/sql/connect/common/src/main/protobuf/spark/connect/catalog.proto index b5341d16ebee6..422bdcf59c83c 100644 --- a/sql/connect/common/src/main/protobuf/spark/connect/catalog.proto +++ b/sql/connect/common/src/main/protobuf/spark/connect/catalog.proto @@ -223,7 +223,10 @@ message UncacheTable { } // See `spark.catalog.clearCache` -message ClearCache { } +message ClearCache { + // (Optional) Whether to clear cached data across all sessions. Defaults to true when omitted. + optional bool all_sessions = 1; +} // See `spark.catalog.refreshTable` message RefreshTable { diff --git a/sql/connect/common/src/main/scala/org/apache/spark/sql/connect/Catalog.scala b/sql/connect/common/src/main/scala/org/apache/spark/sql/connect/Catalog.scala index ce7a10c4026c2..91ec7ef32479f 100644 --- a/sql/connect/common/src/main/scala/org/apache/spark/sql/connect/Catalog.scala +++ b/sql/connect/common/src/main/scala/org/apache/spark/sql/connect/Catalog.scala @@ -643,8 +643,17 @@ class Catalog(sparkSession: SparkSession) extends catalog.Catalog { * @since 3.5.0 */ override def clearCache(): Unit = { + clearCache(allSessions = true) + } + + /** + * Removes cached tables from the in-memory cache. + * + * @since 4.4.0 + */ + override def clearCache(allSessions: Boolean): Unit = { sparkSession.execute { builder => - builder.getCatalogBuilder.getClearCacheBuilder + builder.getCatalogBuilder.getClearCacheBuilder.setAllSessions(allSessions) } } diff --git a/sql/connect/server/src/main/scala/org/apache/spark/sql/connect/planner/SparkConnectPlanner.scala b/sql/connect/server/src/main/scala/org/apache/spark/sql/connect/planner/SparkConnectPlanner.scala index 4803c5865e145..ee5bf6c00abd1 100644 --- a/sql/connect/server/src/main/scala/org/apache/spark/sql/connect/planner/SparkConnectPlanner.scala +++ b/sql/connect/server/src/main/scala/org/apache/spark/sql/connect/planner/SparkConnectPlanner.scala @@ -306,7 +306,8 @@ class SparkConnectPlanner( case proto.Catalog.CatTypeCase.CACHE_TABLE => transformCacheTable(catalog.getCacheTable) case proto.Catalog.CatTypeCase.UNCACHE_TABLE => transformUncacheTable(catalog.getUncacheTable) - case proto.Catalog.CatTypeCase.CLEAR_CACHE => transformClearCache() + case proto.Catalog.CatTypeCase.CLEAR_CACHE => + transformClearCache(catalog.getClearCache) case proto.Catalog.CatTypeCase.REFRESH_TABLE => transformRefreshTable(catalog.getRefreshTable) case proto.Catalog.CatTypeCase.REFRESH_BY_PATH => @@ -4303,8 +4304,9 @@ class SparkConnectPlanner( emptyLocalRelation } - private def transformClearCache(): LogicalPlan = { - session.catalog.clearCache() + private def transformClearCache(clearCache: proto.ClearCache): LogicalPlan = { + val allSessions = !clearCache.hasAllSessions || clearCache.getAllSessions + session.catalog.clearCache(allSessions) emptyLocalRelation } diff --git a/sql/core/src/main/scala/org/apache/spark/sql/classic/Catalog.scala b/sql/core/src/main/scala/org/apache/spark/sql/classic/Catalog.scala index 40c40f6ea78aa..ff7401284f1ff 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/classic/Catalog.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/classic/Catalog.scala @@ -875,7 +875,21 @@ class Catalog(sparkSession: SparkSession) extends catalog.Catalog with Logging { * @since 2.0.0 */ override def clearCache(): Unit = { - sparkSession.sharedState.cacheManager.clearCache() + clearCache(allSessions = true) + } + + /** + * Removes cached tables or views from the in-memory cache. + * + * @group cachemgmt + * @since 4.4.0 + */ + override def clearCache(allSessions: Boolean): Unit = { + if (allSessions) { + sparkSession.sharedState.cacheManager.clearCache() + } else { + sparkSession.sharedState.cacheManager.clearCache(sparkSession) + } } /** diff --git a/sql/core/src/test/scala/org/apache/spark/sql/CachedTableSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/CachedTableSuite.scala index 66732a9e52b00..abaed72804a06 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/CachedTableSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/CachedTableSuite.scala @@ -83,6 +83,22 @@ class CachedTableSuite extends SharedSparkSession } } + test("SPARK-50569: clearCache can be scoped to the current session") { + val otherSession = spark.newSession() + val firstDataFrame = spark.range(1) + val secondDataFrame = otherSession.range(1, 2) + + firstDataFrame.cache() + secondDataFrame.cache() + spark.catalog.clearCache(allSessions = false) + + assert(spark.sharedState.cacheManager.lookupCachedData(firstDataFrame).isEmpty) + assert(otherSession.sharedState.cacheManager.lookupCachedData(secondDataFrame).isDefined) + + otherSession.catalog.clearCache(allSessions = false) + assert(spark.sharedState.cacheManager.isEmpty) + } + def rddIdOf(tableName: String): Int = { val plan = spark.table(tableName).queryExecution.sparkPlan plan.collect { From f0d07705b06487661d484a83a1b181cf970b14a9 Mon Sep 17 00:00:00 2001 From: PA Savalle Date: Mon, 31 Aug 2026 08:39:52 +0200 Subject: [PATCH 3/4] scalafmt --- .../spark/sql/connect/config/Connect.scala | 5 +++-- .../SparkConnectSessionManagerSuite.scala | 16 +++++++++++----- 2 files changed, 14 insertions(+), 7 deletions(-) diff --git a/sql/connect/server/src/main/scala/org/apache/spark/sql/connect/config/Connect.scala b/sql/connect/server/src/main/scala/org/apache/spark/sql/connect/config/Connect.scala index 62f587c0ab01b..39530d92e2e2b 100644 --- a/sql/connect/server/src/main/scala/org/apache/spark/sql/connect/config/Connect.scala +++ b/sql/connect/server/src/main/scala/org/apache/spark/sql/connect/config/Connect.scala @@ -158,8 +158,9 @@ object Connect { val CONNECT_SESSION_MANAGER_CLEANUP_CACHED_DATA_ENABLED = buildStaticConf("spark.connect.session.manager.cleanupCachedData.enabled") - .doc("When true, cached data persisted by an isolated session is removed when the session " + - "is closed. Cached data that is also persisted by another session is preserved.") + .doc( + "When true, cached data persisted by an isolated session is removed when the session " + + "is closed. Cached data that is also persisted by another session is preserved.") .version("4.4.0") .booleanConf .createWithDefault(false) diff --git a/sql/connect/server/src/test/scala/org/apache/spark/sql/connect/service/SparkConnectSessionManagerSuite.scala b/sql/connect/server/src/test/scala/org/apache/spark/sql/connect/service/SparkConnectSessionManagerSuite.scala index e5bbc683d9084..e3d4f69f44c5e 100644 --- a/sql/connect/server/src/test/scala/org/apache/spark/sql/connect/service/SparkConnectSessionManagerSuite.scala +++ b/sql/connect/server/src/test/scala/org/apache/spark/sql/connect/service/SparkConnectSessionManagerSuite.scala @@ -198,9 +198,11 @@ class SparkConnectSessionManagerSuite extends SharedSparkSession { Connect.CONNECT_SESSION_MANAGER_CLEANUP_CACHED_DATA_ENABLED.key -> cleanupCachedData.toString) { val first = SparkConnectService.sessionManager.getOrCreateIsolatedSession( - SessionKey("user", UUID.randomUUID().toString), None) + SessionKey("user", UUID.randomUUID().toString), + None) val second = SparkConnectService.sessionManager.getOrCreateIsolatedSession( - SessionKey("user", UUID.randomUUID().toString), None) + SessionKey("user", UUID.randomUUID().toString), + None) val firstDataFrame = first.session.range(1) second.session.range(1, 2).createTempView("second_view") second.session.catalog.cacheTable("second_view") @@ -217,7 +219,9 @@ class SparkConnectSessionManagerSuite extends SharedSparkSession { SparkConnectService.sessionManager.closeSession(second.key) assert( - second.session.sharedState.cacheManager.lookupCachedData(secondDataFrame).isDefined === + second.session.sharedState.cacheManager + .lookupCachedData(secondDataFrame) + .isDefined === !cleanupCachedData) spark.catalog.clearCache() } @@ -228,9 +232,11 @@ class SparkConnectSessionManagerSuite extends SharedSparkSession { test("SPARK-50569: cached data cleanup preserves entries persisted by another session") { withSparkConf(Connect.CONNECT_SESSION_MANAGER_CLEANUP_CACHED_DATA_ENABLED.key -> "true") { val first = SparkConnectService.sessionManager.getOrCreateIsolatedSession( - SessionKey("user", UUID.randomUUID().toString), None) + SessionKey("user", UUID.randomUUID().toString), + None) val second = SparkConnectService.sessionManager.getOrCreateIsolatedSession( - SessionKey("user", UUID.randomUUID().toString), None) + SessionKey("user", UUID.randomUUID().toString), + None) val firstDataFrame = first.session.range(1) val secondDataFrame = second.session.range(1) From bd3a644ec362080720aafa07b501e958f059f259 Mon Sep 17 00:00:00 2001 From: PA Savalle Date: Mon, 31 Aug 2026 13:51:44 +0200 Subject: [PATCH 4/4] use is_remote in test --- python/pyspark/sql/tests/test_catalog.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/python/pyspark/sql/tests/test_catalog.py b/python/pyspark/sql/tests/test_catalog.py index e6b41f0897c08..0c778ec819abb 100644 --- a/python/pyspark/sql/tests/test_catalog.py +++ b/python/pyspark/sql/tests/test_catalog.py @@ -18,6 +18,7 @@ from pyspark import StorageLevel from pyspark.errors import AnalysisException, PySparkTypeError +from pyspark.sql import is_remote from pyspark.sql.types import IntegerType, StructField, StructType from pyspark.testing.sqlutils import ReusedSQLTestCase @@ -471,7 +472,7 @@ def test_clear_cache_current_session(self): spark.catalog.clearCache() spark.catalog.dropTempView(first_view) other_session.catalog.dropTempView(second_view) - if hasattr(other_session, "client"): + if is_remote(): other_session.client.close() def test_table_exists(self):