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

ASF GitHub Bot updated SPARK-58292:
-----------------------------------
    Labels: pull-request-available  (was: )

> Recreate the netty worker EventLoopGroup when a worker event loop thread dies
> -----------------------------------------------------------------------------
>
>                 Key: SPARK-58292
>                 URL: https://issues.apache.org/jira/browse/SPARK-58292
>             Project: Spark
>          Issue Type: Bug
>          Components: Spark Core
>    Affects Versions: 5.0.0
>            Reporter: ZequnLin
>            Priority: Major
>              Labels: pull-request-available
>
> When a netty worker event-loop thread dies (a Throwable escaping run() at the 
> runIo()/select level; per-task exceptions are swallowed by safeExecute), it 
> is driven to ST_TERMINATED by SingleThreadEventExecutor.doStartThread()'s 
> finally block. That thread is never replaced in the fixed-size 
> MultithreadEventExecutorGroup (children is final, no repopulation), cannot be 
> restarted (startThread() only starts from ST_NOT_STARTED/ST_SUSPENDED), and 
> keeps being handed out by the round-robin EventExecutorChooser (no liveness 
> check). This poisons the TransportClientFactory worker group and produces two 
> failure modes:
> 1. New connections (~1/N fail): a fresh channel bound to the dead loop fails 
> registration with RejectedExecutionException("event executor terminated") 
> (caught by AbstractChannel.AbstractUnsafe.register), so createClient throws 
> IOException. Each connect attempt round-robins across N worker threads, so 
> roughly 1 in N attempts binds to the dead loop and fails.
> 2. Reused cached client (worse - silent hang): a pooled TransportClient 
> pinned to the dead loop still has an open socket, so isActive() is true and 
> createClient keeps returning it. writeAndFlush().addListener() then submits 
> to the dead loop; netty's safeExecute swallows the RejectedExecutionException 
> (only logs "Failed to submit a listener notification task. Event loop shut 
> down?"), the callback is orphaned, and the fetch (broadcast/RDD/RPC) hangs 
> forever.
> Fix (entirely in common/network-common):
> (a) TransportClient.isActive() returns false when 
> channel.eventLoop().isShuttingDown() is true, so a poisoned client is no 
> longer reused.
> (b) In TransportClientFactory.createClient, when the connect fails with a 
> RejectedExecutionException whose message is exactly "event executor 
> terminated" (the terminated-loop rejection only), recreate the worker 
> EventLoopGroup - synchronized and identity-guarded so N concurrent callers 
> swap it exactly once - then rethrow so the existing RetryingBlockTransferor 
> IOException retry path binds onto a fresh, all-live group. The superseded 
> group is not shut down eagerly (its live threads still serve already-open 
> channels) and is shut down best-effort on close().



--
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