eemario commented on code in PR #1147:
URL: https://github.com/apache/flink-agents/pull/1147#discussion_r4121626562
##########
runtime/src/main/java/org/apache/flink/agents/runtime/operator/ActionExecutionOperator.java:
##########
@@ -236,10 +347,23 @@ public void open() throws Exception {
getContainingTask().getEnvironment().getTaskManagerInfo().getTmpDirectories(),
getRuntimeContext().getJobInfo().getJobId(),
metricGroup,
- this::checkMailboxThread,
+ parallelExecutionWithoutCoroutineEnabled
+ ? parallelExecutionLock::checkReentrant
+ : this::checkMailboxThread,
jobIdentifier,
getRuntimeContext().getUserCodeClassLoader());
+ if (parallelExecutionWithoutCoroutineEnabled) {
+ executionCoordinator =
+ new ParallelExecutionCoordinator(
+ parallelExecutionLock,
+ mailboxExecutor::execute,
+ Work::new,
+ pythonBridge::releaseCurrentThreadInterpreter,
Review Comment:
Thanks for the detailed reproduction!
I fixed this by creating cached resource `PyObject` handles through the
owner interpreter, even when initialization is triggered by a managed worker.
Worker access is serialized by the parallel execution lock, while resource
opening remains on the worker to preserve reentrant cache lookup.
This decouples cached resource lifetime from worker lifetime: idle worker
retirement can safely close its interpreter, and cached handles remain valid
until `ResourceCache` closes before the owner interpreter.
I also added a native Pemja E2E test covering resource creation, worker
retirement, reuse by another worker, and final cleanup. Reverting to
worker-interpreter creation reproduces the original native crash.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]