wellkilo opened a new pull request, #1172: URL: https://github.com/apache/flink-agents/pull/1172
Linked issue: #939 ### Purpose of change Flink Agents workflows now complete every keyed input when Flink uses its batch keyed-state backend, instead of rejecting explicit batch mode or silently retaining only the last key's output. Bounded `AUTOMATIC` jobs receive the same behavior when Flink selects that backend. #### Runtime flow Both Java and Python entry points still converge on `CompileUtils.connectToAgent`. The previous graph-construction rejection is removed. During `ActionExecutionOperator.open`, the operator records whether Flink created a `BatchExecutionKeyedStateBackend`. `processElement` continues to enqueue action work through the mailbox; before it returns under that backend, it yields until the current input's complete action chain has removed its processing key. General-purpose keyed-state backends skip this drain and preserve cross-key concurrency. #### Key decisions - Detect the backend Flink actually instantiated rather than infer execution mode from graph-construction configuration. This covers explicit `BATCH` and bounded `AUTOMATIC` jobs. - Reuse the existing mailbox continuation and `waitInFlightEventsFinished` path instead of introducing a second synchronous action executor or moving all keyed runtime state into operator-owned maps. - Serialize inputs only under the batch backend. Its state contract already requires a key to be fully processed before moving to another key; streaming backends retain their existing concurrent behavior. ### Behavioral Semantics #### Interaction decisions | Runtime backend | Action form | Behavior | |---|---|---| | Batch keyed-state backend | Synchronous chain | The chain finishes before `processElement` returns | | Batch keyed-state backend | Async continuation | Mailbox yielding continues until the action chain finishes | | General-purpose keyed-state backend | Sync or async | Existing mailbox scheduling and cross-key concurrency are unchanged | #### Behavioral contracts 1. Every input key in a bounded, single-subtask batch job produces its workflow output. 2. Async action continuations finish without losing keyed state when the batch backend is active. 3. `AUTOMATIC` mode behaves correctly when Flink resolves bounded input to batch execution. 4. Streaming/general-purpose state backends retain cross-key concurrent scheduling. #### Failure behavior Mailbox continuation failures still propagate through the operator's existing action-task failure path and fail the task. The batch drain does not retry, fall back, or absorb failures. No new configuration or external-service failure path is introduced. ### Tests | Contract | Tests | |---|---| | 1 | `CompileUtilsTest.processesEveryKeyWithDefaultBatchStateBackend`; `FlinkIntegrationTest.testBatchExecutionWithMultipleKeys` | | 2 | `CompileUtilsTest.processesAsyncActionsWithDefaultBatchStateBackend` | | 3 | `CompileUtilsTest.processesEveryKeyWhenAutomaticModeSelectsBatchExecution` | | 4 | `ActionExecutionOperatorTest.testDifferentKeyDataCanRunConcurrently`; `testExecuteAsyncWithMultipleKeys` | Verification: - Full Java unit-test reactor passed; runtime: 1,067 tests, 0 failures, 1 skipped. - `FlinkIntegrationTest.testBatchExecutionWithMultipleKeys` passed on Flink 1.20.5 and 2.3.0 locally; CI runs the same test on every supported Flink line from 1.20 through 2.3. - Full 36-module Spotless check passed. - Apache RAT license check passed. Not verified locally: the E2E test was not individually executed against Flink 2.0, 2.1, and 2.2; their dist builds succeeded, and the repository's existing Java IT matrix runs the test for each version. <details> <summary>Implementation invariants and supporting evidence</summary> - The unfixed runtime was reproduced after removing the temporary guard: five distinct keys produced only `[10]` instead of `[2, 4, 6, 8, 10]`. - `BatchExecutionKeyedStateBackend` is present under the same fully qualified name in the supported Flink 1.20.5, 2.0.2, 2.1.3, 2.2.1, and 2.3.0 artifacts. - No new dependency, state field, serializer, checkpoint format, event schema, or public method is introduced. </details> ### API No public signature, JSON, event, or IDL changes. Existing callers using explicit batch mode no longer receive the temporary configuration-time `IllegalStateException`; the previously unsupported configuration now executes correctly. Java and Python bridge calls share the same runtime operator behavior. ### Documentation - [ ] `doc-needed` - [ ] `doc-not-needed` - [x] `doc-included` The runtime method and API contracts are recorded in `runtime/Method.md` and `runtime/API.md`. ### Was this patch authored or co-authored using generative AI tooling? - [x] Yes - [ ] No Generated-by: TraeCode (GPT-5) -- 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]
