Skip to content

Execution hangs in RUNNING when a worker fails to initialize in a region's second phase #8667

Description

@sshiv012

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:

  1. The dependee phase starts only the probe, and succeeds.

  2. 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?

  1. Build: two CSV scans → HashJoin (build, probe) → Python UDF.
  2. Use UDF code without from pytexera import *:
    class ProcessTableOperator(UDFTableOperator):
        def process_table(self, table: Table, port: int) -> Iterator[Optional[TableLike]]:
            yield table
  3. 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)

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

No labels
No labels

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions