Skip to content

Commit b47c144

Browse files
LuciferYangreviewer
andcommitted
[SPARK-44462][SS][CONNECT][FOLLOWUP] Assert foreachBatch close() unregisters the cloned session
### What changes were proposed in this pull request? Followup to #55410. The two `ForeachBatchSessionManager` cleanup tests in `StreamingForeachBatchHelperSuite` checked that `close()` is terminal and that the registered cleaner runs, but neither asserted the side effect `close()` exists to guarantee: that the cloned `SessionHolder` is removed from `SparkConnectService.sessionManager`. Because `close()` sets its `closed` flag independently of the `closeSession` call, a future change that dropped or broke the unregister step would still pass both tests while re-introducing the never-expiring-holder leak. This adds assertions that the cloned `SessionKey` is present in the session manager before `close()` and absent after. ### Why are the changes needed? To make the cloned-session leak-prevention contract regression-sensitive. Spotted in review of #59200, the branch-4.3 backport. ### Does this PR introduce _any_ user-facing change? No. ### How was this patch tested? Pass GitHub Actions. ### Was this patch authored or co-authored using generative AI tooling? Generated-by: Claude Code Closes #59202 from LuciferYang/SPARK-44462-assert-cloned-session-unregister. Lead-authored-by: YangJie <yangjie01@baidu.com> Co-authored-by: reviewer <reviewer@local> Signed-off-by: yangjie01 <yangjie01@baidu.com> (cherry picked from commit 38422ae) Signed-off-by: yangjie01 <yangjie01@baidu.com>
1 parent 9cadc90 commit b47c144

1 file changed

Lines changed: 7 additions & 1 deletion

File tree

‎sql/connect/server/src/test/scala/org/apache/spark/sql/connect/planner/StreamingForeachBatchHelperSuite.scala‎

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,7 @@ import org.scalatestplus.mockito.MockitoSugar
2727

2828
import org.apache.spark.SparkIllegalStateException
2929
import org.apache.spark.sql.connect.SparkConnectTestUtils
30-
import org.apache.spark.sql.connect.service.SessionHolder
30+
import org.apache.spark.sql.connect.service.{SessionHolder, SparkConnectService}
3131
import org.apache.spark.sql.streaming.StreamingQuery
3232
import org.apache.spark.sql.streaming.StreamingQueryListener
3333
import org.apache.spark.sql.test.SharedSparkSession
@@ -128,8 +128,14 @@ class StreamingForeachBatchHelperSuite extends SharedSparkSession with MockitoSu
128128
val clonedHolder = manager.getOrCreateClonedSessionHolder(batchDf)
129129
// The id is pinned for the whole query: every batch resolves to the same holder.
130130
assert(manager.getOrCreateClonedSessionHolder(batchDf) eq clonedHolder)
131+
// The holder is registered in the global session manager and never expires by inactivity, so
132+
// close() must unregister it; a broken unregister step would leak it until server shutdown.
133+
assert(
134+
SparkConnectService.sessionManager.getIsolatedSessionIfPresent(clonedHolder.key).isDefined)
131135

132136
manager.close()
137+
assert(
138+
SparkConnectService.sessionManager.getIsolatedSessionIfPresent(clonedHolder.key).isEmpty)
133139
checkError(
134140
exception =
135141
intercept[SparkIllegalStateException](manager.getOrCreateClonedSessionHolder(batchDf)),

0 commit comments

Comments
 (0)