[
https://issues.apache.org/jira/browse/SPARK-58447?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Jiwon Park updated SPARK-58447:
-------------------------------
Affects Version/s: 4.3.0
(was: 5.0.0)
> Preserve operation tags in Spark Connect Python foreachBatch callbacks
> ----------------------------------------------------------------------
>
> Key: SPARK-58447
> URL: https://issues.apache.org/jira/browse/SPARK-58447
> Project: Spark
> Issue Type: Bug
> Components: Connect
> Affects Versions: 4.3.0
> Reporter: Jiwon Park
> Priority: Major
>
> Spark Connect does not preserve operation tags from the `writeStream.start()`
> request in Python `foreachBatch` callbacks.
> The callback runs in a separate Python worker that creates a local Connect
> client. Because operation tags are client thread-local, the callback session
> does not contain the tags from the originating request. Consequently, nested
> actions such as `count`, writes, and SQL commands are executed without those
> tags.
> A minimal reproduction is:
> spark.addTag("foreach-batch")
> def process_batch(df, _):
> assert "foreach-batch" in df.sparkSession.getTags()
> df.count()
> query = stream.writeStream.foreachBatch(process_batch).start()
> query.processAllAvailable()
> The assertion fails with Spark Connect. The callback session should inherit
> the operation tags captured when the streaming query starts.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]