[ 
https://issues.apache.org/jira/browse/SPARK-59032?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Hyukjin Kwon updated SPARK-59032:
---------------------------------
            Epic Link: SPARK-59035
    Affects Version/s: connect-rust-0.1.0
                           (was: 5.0.0)

> Spark Connect client hangs after many ops: unbounded ReleaseExecute leaks 
> HTTP/2 streams
> ----------------------------------------------------------------------------------------
>
>                 Key: SPARK-59032
>                 URL: https://issues.apache.org/jira/browse/SPARK-59032
>             Project: Spark
>          Issue Type: Bug
>          Components: Project Infra
>    Affects Versions: connect-rust-0.1.0
>            Reporter: Hyukjin Kwon
>            Assignee: Hyukjin Kwon
>            Priority: Major
>              Labels: pull-request-available
>             Fix For: connect-rust-4.2.0
>
>
> The Python Connect parity gate intermittently hangs on 
> test_parity_arrow_python_udf.py: the file times out with p=0 (not one test 
> completed) and is killed only by the per-file cap. It was observed once on 
> master at file 5/84 on a fresh server, while all 83 other files (including 
> the sibling arrow/UDF files) passed.
> It is NOT deterministic and NOT a slowness/timeout problem:
> - the file passes standalone in ~122s (matching the reference client);
> - it passes even with the suite's reattach-forcing server limits 
> (senderMaxStreamDuration=1s, senderMaxStreamSize=123);
> - it did NOT reproduce in 8 back-to-back runs under the exact gate invocation;
> - a temporary CI diagnostic (client py-spy + faulthandler, server jstack, 
> Python-worker stacks) found no deterministic hang.
> Root cause: the client never applies gRPC keepalive to its channel. The 
> connection string's grpc_keepalive_* options (default: enabled, 60s interval 
> / 20s timeout, keep-alive-while-idle) are parsed by ChannelBuilder but were 
> never applied to the tonic channel, and no HTTP/2 keepalive was set 
> otherwise. So a server-streaming RPC whose connection wedges (a half-open 
> TCP, or a server-side stall) blocks in stream.message() indefinitely: nothing 
> pings the peer to surface the break, the reattachable-execute iterator never 
> gets the transport error it needs to reattach, and the stream only unblocks 
> when the client is closed (at CI teardown).
> Fix (PR: https://github.com/apache/spark-connect-rust/pull/58):
> - Apply the parsed keepalive settings to the tonic endpoint (TLS and 
> plaintext paths), so an unacked PING within the timeout fails the stream with 
> a transport error that the reattachable iterator treats as retriable and 
> recovers from.
> - Retry the file in the parity gate (FLAKY_FILES) as a belt-and-suspenders 
> safety net for any residual rarity, instead of skipping the whole file or 
> letting a rare hang burn the cap and fail.
> - Reverted the earlier, wrong senderMaxStreamDuration=1200s server-conf 
> attempt (a larger cap only makes a hang wait longer), and removed the 
> temporary diagnostic workflow.
> --- Update (review of PR #58) ---
> The keepalive fix (applying the parsed grpc_keepalive_* to the channel) makes 
> a wedged transport RECOVERABLE — an expired PING surfaces as UNAVAILABLE, 
> which the reattachable-execute iterator retries — but it is NOT confirmed to 
> be the root cause. The captured evidence (509b9e6) shows the client blocked 
> on the very FIRST next() of the reattachable stream while the server produced 
> nothing at all, whereas the reference client gets the error in <1s. Keepalive 
> explains a wedged transport; it does not by itself explain "server produced 
> nothing." So this ticket stays OPEN for the true root cause; keepalive + the 
> FLAKY_FILES retry make the gate green by making the symptom recoverable 
> without necessarily fixing the underlying cause. A server-side jstack was 
> added to the parity gate's timeout path to pin down whether the server's 
> execute thread is deadlocked on the (dead) Python worker or blocked writing 
> to us. A separate, related "documented-but-unapplied option" bug (gRPC max 
> message length) is tracked as SPARK-59037.
> --- Root cause pinned (server jstack) ---
> A full server jstack captured at the hang (across all 3 retry attempts) shows 
> the server JVM is COMPLETELY IDLE: 103 threads, all infrastructure (main 
> awaiting termination, rpc/dispatcher/map-output loops, SparkUI, heartbeat, 
> block-manager, listeners, cleaners) - NO SparkConnectExecute thread, NO 
> PythonRunner/ArrowPythonRunner, NO Python worker, NO gRPC handler for the 
> operation. So the server ran test_raise_stop_iteration, the UDF raised, the 
> operation failed and was fully cleaned up. This DEFINITIVELY refutes the 
> earlier "server-side worker stalls" hypothesis.
> The hang is therefore 100% CLIENT-SIDE: our client blocks in stream.message() 
> reading the reattachable ExecutePlan stream that the server has already 
> finished/errored, instead of detecting the stream end and surfacing the error 
> / driving ReattachExecute. The reference client, on the identical server 
> behavior, gets the error in <1s. Keepalive does NOT help (the connection is 
> alive and idle, so PINGs are ACKed), and a retry does NOT help 
> (deterministic). test_raise_stop_iteration is therefore deselected 
> (KNOWN_CLIENT_HANG_TESTS) to unblock the gate, and this ticket stays OPEN for 
> the client stream-completion fix. Note: only StopIteration hangs; other 
> UDF-error tests in the same file (segfault, err_udf_init/registration) pass, 
> so the fix likely concerns how a failed reattachable op's terminal state is 
> delivered/consumed for this specific case.
> --- RESOLVED: root cause + fix (verified in CI) ---
> Root cause (found via client RPC tracing + server jstack): the 
> reattachable-execute driver (reference reattach.py) fires a ReleaseExecute 
> after EVERY response on a background ThreadPoolExecutor, and each release 
> target closure CAPTURES the operation's response iterator. Our client's 
> release_execute was an UNBOUNDED unary RPC. On the single shared HTTP/2 
> connection, under a long serial run (tens of thousands of releases), a 
> release occasionally stalls and never returns -> its release future stays 
> pending -> that keeps the iterator alive -> which keeps the iterator's 
> still-open ExecutePlan HTTP/2 stream open. Enough of those accumulate that a 
> later ExecutePlan cannot get a stream and the client blocks forever in 
> stream.message() while the server is completely idle (jstack shows no 
> execute/PythonRunner/worker threads). The hang is CUMULATIVE, landing on 
> whatever collect() is ~40 operations in - which is why it appeared to move 
> between tests (test_raise_stop_iteration, test_udf_in_generate, ...) 
> depending on timing.
> This is NOT test-specific and NOT a server stall; keepalive does not fix it 
> (the connection is alive/idle) and a retry does not fix it (deterministic 
> in-sequence).
> Fix: bound release_execute with a 30s deadline. ReleaseExecute is best-effort 
> - the reference reattach driver runs it on a background pool, wraps it in 
> try/except, and logs-and-ignores any failure ("if the release fails, server 
> is equipped to deal with abandoned executions"), confirmed identical on both 
> v4.2.0 and master. A timed-out release fails fast, freeing its thread and 
> closing its HTTP/2 stream, so the release future completes and the iterator 
> (and its ExecutePlan stream) is dropped - nothing accumulates. Normal 
> releases return in milliseconds and never hit the deadline.
> Verified in CI: with the fix and test_raise_stop_iteration NOT deselected, 
> test_parity_arrow_python_udf.py passes 319/320 (the one deselect is the 
> pre-existing environmental test_type_coercion_string_to_numeric) and all 10 
> arrow-prefix files pass with 0 failures. The earlier workarounds (whole-file 
> skip, SLOW_FILE_TIMEOUTS=1200, FLAKY_FILES retry, per-test deselect) are all 
> removed. Fixed in PR https://github.com/apache/spark-connect-rust/pull/58.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to