yunfengzhou-hub opened a new pull request, #1183:
URL: https://github.com/apache/flink-agents/pull/1183

   Linked issue: close #1080
   
   ### Purpose of change
   
   Flink Agents has no first-class way to materialize agent output as a 
**typed** `DataStream`/`Table`: the output type is not part of the terminal, so 
a typed stream means manually casting the untyped `DataStream<Object>`, and 
there is no way to declare the type of a `Table` output. Background and 
current-state analysis are in #1080.
   
   This PR implements the direction agreed there — declare the output type 
directly on the terminal:
   
   - Java: `toDataStream(TypeInformation)` / `toDataStream(Class)` and 
`toTable(TypeInformation)` / `toTable(Class)`, alongside the existing untyped 
`toDataStream()` and schema-only `toTable(Schema)`.
   - Python: `to_datastream(output_type=None)`, and `to_table(schema=None, 
output_type=None)` with exactly one required (cross-checked when both are 
given).
   
   The declared type is applied by a downstream conversion operator on the 
shared raw stream, so the untyped view is preserved and several typed views of 
different types can coexist on one execution.
   
   ### Behavioral Semantics
   
   - **Untyped view preserved.** The agent operator keeps emitting `Object` 
(Java) / serialized bytes (Python); the raw stream's element type never changes 
and still exposes heterogeneous output through `toDataStream()` / 
`to_datastream()`.
   - **Independent typed views.** The raw stream is built once and cached at 
the untyped boundary; each typed call layers its own conversion operator, so 
multiple typed outputs of different types can coexist.
   - **Table schema derivation.** `toTable(TypeInformation)` / 
`to_table(output_type=...)` with no schema derives the physical columns from 
the declared type (a POJO's / structured type's fields become columns).
   - **Python `to_table` contract and failure behavior.** Exactly one of 
`schema` / `output_type` is required: neither raises `ValueError`; a 
non-`Schema` in the `schema` slot raises `TypeError` (directing type 
declarations to `output_type=`); both are cross-checked and raise `ValueError` 
if they describe different row types. A given `schema` still drives the 
physical table, preserving Table-domain information (primary key, 
computed/metadata columns, watermark).
   - **Table terminals differ by language, intentionally.** Java exposes 
`toTable(Schema)` and `toTable(TypeInformation)` as separate overloads; 
Python's single `to_table(schema, output_type)` accepts either or both.
   
   ### Tests
   
   - **Java unit**: `TypedOutputTerminalTest` (raw-stream caching, no type 
welding onto the shared raw stream, `Class` overload == `TypeInformation` 
overload, multiple typed views coexist), `OutputTypeUtilsTest` 
(schema→RowTypeInfo; adaptToRow for Row/Map/POJO/scalar), and 
`AgentBuilderApplyByNameTest` stubs updated for the new abstract methods. 
Focused `mvn` build green; `spotless:check` clean.
   - **Java e2e**: `FlinkIntegrationTest` adds 
`testToDataStreamWithTypeInformation`, `testToTableWithTypeInformation` (schema 
derived from the POJO), and `testToTableWithMultiColumnSchema` (schema-only 
path), backed by a structured `ReviewOutput` agent; examples + integration 
module test-compile.
   - **Python**: `test_remote_execution_environment.py` (untyped raw view; 
scalar/structured/RowTypeInfo outputs; caching without welding; multiple typed 
views coexist; schema-only, derived-schema, and cross-check match/mismatch 
Table; `to_table` requires exactly one) and `test_output_type_utils.py`; `ruff 
check` / `ruff format` clean; runtime unit suite green.
   - Suites needing external services (Ollama/OpenAI) and a real cluster are 
migrated but not executed locally.
   
   ### API
   
   Public API change within the 0.4 breaking-change window.
   
   - **Java (additive)**: `AgentBuilder` gains `toDataStream(TypeInformation)`, 
`toDataStream(Class)`, `toTable(TypeInformation)`, `toTable(Class)`. The 
untyped `toDataStream()` and schema-only `toTable(Schema)` remain; 
`toTable(Schema)` is also corrected to map elements onto the declared columns 
instead of degrading to a single `f0`. Nothing is removed.
   - **Python**: `to_datastream(output_type=None)` is now declared on the 
abstract API (the raw `to_datastream()` is unchanged). `to_table(schema=None, 
output_type=None)` changes from "both required" to "exactly one required, 
cross-checked when both" — all in-repo callers are migrated in this PR. The 
first-output-type cache is removed.
   
   ### Documentation
   
   - [x] `doc-included` — 
`docs/content/docs/development/integrate_with_flink.md` "Typed outputs" section 
rewritten for the typed terminals in both languages.
   
   ### Was this patch authored or co-authored using generative AI tooling?
   
   - [x] Yes
   
   `Generated-by: Qoder 1.32.1 (Qwen3.8-Max)`
   


-- 
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