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]

Reply via email to