[api][runtime][java][python] Simplify agent output materialization with typed toDataStream/toTable - #1183
Open
yunfengzhou-hub wants to merge 5 commits into
Open
[api][runtime][java][python] Simplify agent output materialization with typed toDataStream/toTable#1183yunfengzhou-hub wants to merge 5 commits into
yunfengzhou-hub wants to merge 5 commits into
Conversation
yunfengzhou-hub
force-pushed
the
issue-1080-typed-output-terminals
branch
2 times, most recently
from
October 1, 2026 10:25
4f64f12 to
791e45e
Compare
Add flink_agents.runtime.output_type_utils with the pure helpers the typed output terminals rely on: infer a RowTypeInfo from a structured Python type (Pydantic model / dataclass / named tuple / TypedDict), convert between a RowTypeInfo and a Table Schema through Flink's Java type conversions so nested ROW and the full column range round-trip, reduce a RowTypeInfo to a picklable row shape, build a positional Row from an emitted value with that shape, and rebuild declared instances. Only the picklable shape is captured in a conversion closure, so it survives serialization to the workers while the py4j-backed RowTypeInfo stays on the driver. Refs apache#1080 Generated-by: Qoder 1.32.1 (Qwen3.8-Max)
Expose the output type on the terminals instead of a separate declaration step. to_datastream gains an optional output_type; to_table takes schema and output_type as optional, requires at least one, and cross-checks the two when both are given (previously both were required positionally). A declared type is applied by a downstream conversion operator, so the agent operator keeps emitting serialized bytes and the unrestricted to_datastream() view still exposes heterogeneous output. The raw stream is now cached at the untyped boundary, which fixes a bug where the first output_type was baked into the shared cache and silently reused by later terminals; several typed views of different types can now coexist on one execution. Refs apache#1080 Generated-by: Qoder 1.32.1 (Qwen3.8-Max)
Add type-carrying overloads to AgentBuilder alongside the existing raw toDataStream() and schema-only toTable(Schema): toDataStream(TypeInformation) and toTable(TypeInformation) materialize a typed stream/table, and toDataStream(Class) / toTable(Class) are conveniences that derive the TypeInformation from the class. The declared type is applied by a downstream conversion operator, so the agent operator keeps emitting Object and the raw stream is left unchanged; each typed terminal layers its own conversion on the shared raw stream, so several typed outputs of different types can coexist with the unrestricted view. toTable(TypeInformation) derives the physical schema from the stream's type. Refs apache#1080 Generated-by: Qoder 1.32.1 (Qwen3.8-Max)
Declare the output type at the terminal in WorkflowMultipleAgentExample via toDataStream(ProductReviewAnalysisRes.class), removing the manual downcast of the raw stream. Add end-to-end coverage in the integration tests with a structured TypedOutput agent exercising toDataStream(TypeInformation) and toTable(TypeInformation), whose schema is derived from the POJO. Refs apache#1080 Generated-by: Qoder 1.32.1 (Qwen3.8-Max)
Describe passing the output type directly to to_datastream/toDataStream and to_table/toTable, and the schema/output_type rules for the Table terminal in both languages. Refs apache#1080 Generated-by: Qoder 1.32.1 (Qwen3.8-Max)
yunfengzhou-hub
force-pushed
the
issue-1080-typed-output-terminals
branch
from
October 1, 2026 13:48
791e45e to
a0b6aae
Compare
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
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 untypedDataStream<Object>, and there is no way to declare the type of aTableoutput. Background and current-state analysis are in #1080.This PR implements the direction agreed there — declare the output type directly on the terminal:
toDataStream(TypeInformation)/toDataStream(Class)andtoTable(TypeInformation)/toTable(Class), alongside the existing untypedtoDataStream()and schema-onlytoTable(Schema).to_datastream(output_type=None), andto_table(schema=None, output_type=None)with at least 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
Object(Java) / serialized bytes (Python); the raw stream's element type never changes and still exposes heterogeneous output throughtoDataStream()/to_datastream().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).to_tablecontract and failure behavior. At least one ofschema/output_typeis required: neither raisesValueError; a non-Schemain theschemaslot raisesTypeError(directing type declarations tooutput_type=); aSchemain theoutput_typeslot raisesTypeError(directing it toschema=); when both are given they are cross-checked and raiseValueErrorif they describe different row types. A givenschemastill drives the physical table, preserving Table-domain information (primary key, computed/metadata columns, watermark); the derived row type reads only the schema's physical columns, since computed/metadata columns are planner-derived.toTable(Schema)andtoTable(TypeInformation)as separate overloads; Python's singleto_table(schema, output_type)accepts either or both.Tests
TypedOutputTerminalTest(raw-stream caching, no type welding onto the shared raw stream,Classoverload ==TypeInformationoverload, multiple typed views coexist, typedtoTablederives columns from the declared type) andAgentBuilderApplyByNameTeststubs updated for the new abstract methods. Focusedmvnbuild green;spotless:checkclean.FlinkIntegrationTestaddstestToDataStreamWithTypeInformationandtestToTableWithTypeInformation(schema derived from the POJO), backed by a structuredTypedOutputagent; examples + integration module test-compile.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_tablerequires at least one and rejects aSchemain either slot) andtest_output_type_utils.py(type inference and schema→row-type derivation, including skipping computed/metadata columns);ruff check/ruff formatclean; runtime unit suite green.API
Public API change within the 0.4 breaking-change window.
AgentBuildergainstoDataStream(TypeInformation),toDataStream(Class),toTable(TypeInformation),toTable(Class). The untypedtoDataStream()and schema-onlytoTable(Schema)remain unchanged. Nothing is removed.to_datastream(output_type=None)is now declared on the abstract API (the rawto_datastream()is unchanged).to_table(schema=None, output_type=None)changes from "both required" to "at least one required, cross-checked when both" — all in-repo callers are migrated in this PR. The first-output-type cache is removed.Documentation
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?
Generated-by: Qoder 1.32.1 (Qwen3.8-Max)