joeyutong commented on code in PR #924:
URL: https://github.com/apache/flink-agents/pull/924#discussion_r3702067216
##########
runtime/src/main/java/org/apache/flink/agents/runtime/operator/ActionExecutionOperator.java:
##########
@@ -369,36 +392,54 @@ private void processActionTaskForKey(Object key) throws
Exception {
key, sequenceNumber, actionTask.action,
actionTask.event);
}
- // Set up durable execution context for fine-grained recovery
- durableExecManager.setupDurableExecutionContext(
- actionTask, actionState, sequenceNumber);
-
- ActionTask.ActionTaskResult actionTaskResult =
- actionTask.invoke(
- getRuntimeContext().getUserCodeClassLoader(),
- this.pythonBridge.getPythonActionExecutor());
-
- // We remove the contexts from the map after the task is
processed. They will be added
- // back later if the action task has a generated action task,
meaning it is not
- // finished.
- contextManager.removeMemoryContext(actionTask);
- durableExecManager.removeDurableContext(actionTask);
- contextManager.removeContinuationContext(actionTask);
- contextManager.removePythonAwaitableRef(actionTask);
- durableExecManager.maybePersistTaskResult(
- key,
- sequenceNumber,
- actionTask.action,
- actionTask.event,
- actionTask.getRunnerContext(),
- actionTaskResult);
- isFinished = actionTaskResult.isFinished();
- outputEvents = actionTaskResult.getOutputEvents();
- generatedActionTaskOpt = actionTaskResult.getGeneratedActionTask();
+ notifyActionStarted(actionTask);
+ try {
+ // Set up durable execution context for fine-grained recovery
+ durableExecManager.setupDurableExecutionContext(
+ actionTask, actionState, sequenceNumber);
+
+ ActionTask.ActionTaskResult actionTaskResult =
+ actionTask.invoke(
+ getRuntimeContext().getUserCodeClassLoader(),
+ this.pythonBridge.getPythonActionExecutor());
+
+ // Drop task-local contexts after each step; continuations
transfer them back.
+ contextManager.removeMemoryContext(actionTask);
+ durableExecManager.removeDurableContext(actionTask);
+ contextManager.removeContinuationContext(actionTask);
+ contextManager.removePythonAwaitableRef(actionTask);
+ durableExecManager.maybePersistTaskResult(
+ key,
+ sequenceNumber,
+ actionTask.action,
+ actionTask.event,
+ actionTask.getRunnerContext(),
+ actionTaskResult);
+ isFinished = actionTaskResult.isFinished();
+ outputEvents = actionTaskResult.getOutputEvents();
+ generatedActionTaskOpt =
actionTaskResult.getGeneratedActionTask();
+ notifyFinished = isFinished;
+ } catch (Exception e) {
Review Comment:
Good catch. I now catch `Throwable` around Action invocation so raw errors
emit `failed` before propagating, and emit `finished` immediately after a
completed result is persisted, before processing output Events. This keeps the
Action lifecycle paired without attributing listener or routing failures to the
Action. I added regression tests for both paths.
--
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]