What happened?
If a worker fails to initialize in the second phase of a region, the execution hangs in RUNNING forever. That includes, for example, a Python UDF whose code raises at load time. No error reaches the UI, and every operator in the region stays at 0 tuples.
This happens when the region has an operator with a dependee input port, such as a HashJoin probe. The region launches in two phases:
-
The dependee phase starts only the probe, and succeeds.
-
After the dependee port completes, WorkflowExecutionManager.advanceRegionExecutions launches the remaining operators through syncStatusAndTransitionRegionExecutionPhase(). It discards the returned Future:
// WorkflowExecutionManager.scala
unfinishedRegionManagers.foreach(_.syncStatusAndTransitionRegionExecutionPhase())
If initializeExecutor fails for any worker, the chain initExecutors → assignPorts → connectChannels → openOperators → startWorkers stops. None of the region's workers are started. The failure never reaches AdvanceRegionExecutionsHandler's .onFailure → FatalError.
The failure also isn't visible in the logs. AsyncRPCClient.logControlReply calls logger.error(s"received error from $channelId", err). err is a ControlError, not a Throwable, so the logger drops it and only received error from ... is printed.
Expected: the execution fails, and the UI shows the worker's error, e.g. name 'UDFTableOperator' is not defined.
How we found it: a user reported that a workflow "kept hanging" on every new computing unit (internal ticket TEX-81). Monitoring showed:
- the computing unit was idle, with no OOM or restarts,
- every source in the second region had output 0 rows,
- the only clue was a bare
received error from Worker:...-main-0 line.
Re-running while tailing the logs, then reading the coordinator code, led to the dropped Future. The user's UDF was missing from pytexera import *. Adding it made the workflow complete, which confirmed the chain above.
How to reproduce?
- Build: two CSV scans → HashJoin (build, probe) → Python UDF.
- Use UDF code without
from pytexera import *:
class ProcessTableOperator(UDFTableOperator):
def process_table(self, table: Table, port: int) -> Iterator[Optional[TableLike]]:
yield table
- Run. The workflow stays
RUNNING: the HashJoin build side completes, and every other operator stays at 0 with no error shown.
Branch
main
Relevant log output
[ERROR] [CONTROLLER] [AsyncRPCClient] - received error from ChannelIdentity(ActorVirtualIdentity(Worker:WF0-op9-main-0),ActorVirtualIdentity(CONTROLLER),true)
What happened?
If a worker fails to initialize in the second phase of a region, the execution hangs in
RUNNINGforever. That includes, for example, a Python UDF whose code raises at load time. No error reaches the UI, and every operator in the region stays at 0 tuples.This happens when the region has an operator with a dependee input port, such as a HashJoin probe. The region launches in two phases:
The dependee phase starts only the probe, and succeeds.
After the dependee port completes,
WorkflowExecutionManager.advanceRegionExecutionslaunches the remaining operators throughsyncStatusAndTransitionRegionExecutionPhase(). It discards the returnedFuture:// WorkflowExecutionManager.scala unfinishedRegionManagers.foreach(_.syncStatusAndTransitionRegionExecutionPhase())If
initializeExecutorfails for any worker, the chaininitExecutors → assignPorts → connectChannels → openOperators → startWorkersstops. None of the region's workers are started. The failure never reachesAdvanceRegionExecutionsHandler's.onFailure → FatalError.The failure also isn't visible in the logs.
AsyncRPCClient.logControlReplycallslogger.error(s"received error from $channelId", err).erris aControlError, not aThrowable, so the logger drops it and onlyreceived error from ...is printed.Expected: the execution fails, and the UI shows the worker's error, e.g.
name 'UDFTableOperator' is not defined.How we found it: a user reported that a workflow "kept hanging" on every new computing unit (internal ticket TEX-81). Monitoring showed:
received error from Worker:...-main-0line.Re-running while tailing the logs, then reading the coordinator code, led to the dropped
Future. The user's UDF was missingfrom pytexera import *. Adding it made the workflow complete, which confirmed the chain above.How to reproduce?
from pytexera import *:RUNNING: the HashJoin build side completes, and every other operator stays at 0 with no error shown.Branch
main
Relevant log output