Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
72 changes: 68 additions & 4 deletions api/src/main/java/org/apache/flink/agents/api/AgentBuilder.java
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
package org.apache.flink.agents.api;

import org.apache.flink.agents.api.agents.Agent;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.table.api.Schema;
import org.apache.flink.table.api.Table;
Expand Down Expand Up @@ -72,20 +73,83 @@ default AgentBuilder apply(String agentName) {
* Get output DataStream of agent execution.
*
* <p>This method converts the agent's output events into a Flink DataStream that can be further
* processed in the Flink pipeline.
* processed in the Flink pipeline. The returned view is unrestricted: its elements are whatever
* the agent emitted ({@code Object}). To obtain a typed view, declare the output type with
* {@link #toDataStream(TypeInformation)} or {@link #toDataStream(Class)}.
*
* @return DataStream containing outputs from agent execution.
*/
DataStream<Object> toDataStream();

/**
* Get output Table of agent execution.
* Get output DataStream of agent execution, typed as {@code T}.
*
* <p>This method converts the agent's output events into a Flink Table using the provided
* schema for structure definition.
* <p>The agent operator keeps emitting {@code Object} on the shared raw stream and the declared
* type is applied by a downstream conversion operator, so the unrestricted {@link
* #toDataStream()} view is unaffected and heterogeneous output stays available through it. Each
* call is independent, so several typed views of different types can coexist on one execution.
*
* @param typeInformation the declared type of the agent's output elements.
* @param <T> the declared output element type.
* @return DataStream whose elements are typed as {@code T}.
*/
<T> DataStream<T> toDataStream(TypeInformation<T> typeInformation);

/**
* Get output DataStream of agent execution, typed as {@code T}, deriving the {@link
* TypeInformation} from the given output class.
*
* <p>Convenience overload for the common case. Prefer {@link #toDataStream(TypeInformation)}
* for generic or custom types that a {@link Class} cannot express.
*
* @param outputType the declared class of the agent's output elements.
* @param <T> the declared output element type.
* @return DataStream whose elements are typed as {@code T}.
*/
default <T> DataStream<T> toDataStream(Class<T> outputType) {
return toDataStream(TypeInformation.of(outputType));
}

/**
* Get output Table of agent execution, materializing the unrestricted view with an explicit
* physical schema.
*
* <p>The table's row type is derived from the schema's physical columns, and the Table planner
* adapts each agent output element into a matching row. Computed and metadata columns declared
* in the schema are derived by the planner rather than read from the agent output. To declare
* the output with a type instead of a physical schema, use {@link #toTable(TypeInformation)} or
* {@link #toTable(Class)}.
*
* @param schema Schema indicating the structure of the output table.
* @return Table containing outputs from agent execution.
*/
Table toTable(Schema schema);

/**
* Get output Table of agent execution, typed as {@code T}, deriving the physical schema from
* the declared type (for example, a POJO's fields become columns).
*
* <p>The declared type is applied by a downstream conversion operator on the shared raw stream,
* so the unrestricted {@link #toDataStream()} view is unaffected.
*
* @param typeInformation the declared type of the agent's output elements.
* @param <T> the declared output element type.
* @return Table containing the agent's output, typed as {@code T}.
*/
<T> Table toTable(TypeInformation<T> typeInformation);

/**
* Get output Table of agent execution, typed as {@code T}, deriving the {@link TypeInformation}
* from the given output class.
*
* <p>Convenience overload for the common case. Prefer {@link #toTable(TypeInformation)} for
* generic or custom types that a {@link Class} cannot express.
*
* @param outputType the declared class of the agent's output elements.
* @param <T> the declared output element type.
* @return Table containing the agent's output, typed as {@code T}.
*/
default <T> Table toTable(Class<T> outputType) {
return toTable(TypeInformation.of(outputType));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
package org.apache.flink.agents.api;

import org.apache.flink.agents.api.agents.Agent;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.table.api.Schema;
import org.apache.flink.table.api.Table;
Expand Down Expand Up @@ -62,10 +63,20 @@ public DataStream<Object> toDataStream() {
return null;
}

@Override
public <T> DataStream<T> toDataStream(TypeInformation<T> typeInformation) {
return null;
}

@Override
public Table toTable(Schema schema) {
return null;
}

@Override
public <T> Table toTable(TypeInformation<T> typeInformation) {
return null;
}
}

/** Stub builder that overrides apply(String) to resolve from a one-agent registry. */
Expand Down Expand Up @@ -104,9 +115,19 @@ public DataStream<Object> toDataStream() {
return null;
}

@Override
public <T> DataStream<T> toDataStream(TypeInformation<T> typeInformation) {
return null;
}

@Override
public Table toTable(Schema schema) {
return null;
}

@Override
public <T> Table toTable(TypeInformation<T> typeInformation) {
return null;
}
}
}
78 changes: 53 additions & 25 deletions docs/content/docs/development/integrate_with_flink.md
Original file line number Diff line number Diff line change
Expand Up @@ -161,7 +161,6 @@ Integrate the agent with input `Table`, and return the output `Table` can be con

{{< tab "Python" >}}
```python
from pyflink.common.typeinfo import BasicTypeInfo, ExternalTypeInfo, RowTypeInfo
from pyflink.datastream import KeySelector
from pyflink.table import DataTypes, Schema

Expand All @@ -179,19 +178,14 @@ input_table = t_env.from_elements(
["id", "input"],
)

# The output TypeInformation and Schema must be mutually consistent: both
# describe a single "result" INT column here.
output_type = ExternalTypeInfo(RowTypeInfo(
[BasicTypeInfo.INT_TYPE_INFO()],
["result"],
))

schema = (Schema.new_builder().column("result", DataTypes.INT())).build()
# A single output-type declaration: the Schema's physical columns give the
# output row type, and each agent output element is adapted into a matching row.
schema = Schema.new_builder().column("result", DataTypes.INT()).build()

output_table = (
agents_env.from_table(input=input_table, key_selector=MyKeySelector())
.apply(your_agent)
.to_table(schema=schema, output_type=output_type)
.to_table(schema=schema)
)
```
{{< /tab >}}
Expand Down Expand Up @@ -225,11 +219,14 @@ Table inputTable =
Row.of(2, "Bob", 92.0),
Row.of(3, "Charlie", 78.3));

// The agent output is exposed as a single anonymous column named "f0".
// Declare "f0" with the agent's OUTPUT type: a scalar type (e.g. DataTypes.STRING())
// when the agent emits a scalar value, or a nested DataTypes.ROW(...) only when the
// agent emits a composite row.
Schema outputSchema = Schema.newBuilder().column("f0", DataTypes.STRING()).build();
// Declare the output columns. Each agent output element is adapted into a row by
// matching columns by name (from a POJO's getters/public fields, a Row, or a Map);
// a scalar output maps to a single-column schema.
Schema outputSchema =
Schema.newBuilder()
.column("name", DataTypes.STRING())
.column("score", DataTypes.DOUBLE())
.build();

Table outputTable =
agentsEnv
Expand All @@ -244,18 +241,49 @@ Table outputTable =

User should provide `KeySelector` in `from_table()` to tell how to convert the input `Table` to `KeyedStream` internally.

The arguments required by `to_table()` differ by language:
## Typed outputs

`to_datastream()` / `toDataStream()` return an unrestricted stream whose elements are whatever the agent emitted (`Object` in Java). To give the output a concrete type, pass the type straight to the terminal: `to_datastream(output_type)` / `toDataStream(TypeInformation)` materialize a typed stream, and `to_table(output_type=...)` / `toTable(TypeInformation)` materialize a typed table. The unrestricted call stays available and unchanged:

{{< tabs "Typed outputs" >}}

{{< tab "Python" >}}
```python
builder = agents_env.from_datastream(input_stream, key_selector).apply(your_agent)

raw = builder.to_datastream() # unrestricted: whatever the agent emitted

# Pass the output type to the terminal to materialize a typed stream or table.
typed = builder.to_datastream(ReviewOutput) # stream of validated ReviewOutput
table = builder.to_table(output_type=ReviewOutput) # schema derived from the type
# Or pass both: the Schema drives the physical table and is cross-checked
# against the declared type.
both = builder.to_table(review_schema, ReviewOutput)
```
{{< /tab >}}

{{< tab "Java" >}}
```java
AgentBuilder builder =
agentsEnv.fromDataStream(inputStream, keySelector).apply(yourAgent);

DataStream<Object> raw = builder.toDataStream(); // unrestricted

// Pass the output type to the terminal to materialize a typed stream or table.
DataStream<ReviewOutput> typed = builder.toDataStream(ReviewOutput.class);
Table fromType = builder.toTable(ReviewOutput.class); // schema derived from the type
```
{{< /tab >}}

{{< /tabs >}}

- **Python**: provide both `Schema` and `TypeInformation` to define the output `Table` schema.
- **Java**: provide only `Schema` (`toTable(Schema)`); `TypeInformation` is not required.
For the `Table` terminal:

The two languages also name the output columns differently:
- Calling `to_table(output_type=...)` / `toTable(TypeInformation)` with no schema derives the physical columns from the declared type, so a POJO's fields (or a structured Python type's fields) become columns.
- In Python, `to_table(schema, output_type)` accepts both: the `Schema` drives the physical table, preserving Table-domain information such as a primary key, computed columns, metadata columns, or a watermark, and is cross-checked against the declared type, raising if the two describe different row types. At least one of the two is required. In Java the schema-only `toTable(Schema)` and the type-only `toTable(TypeInformation)` are separate overloads.

- **Python**: the `TypeInformation` passed to `to_table()` is a `RowTypeInfo` whose field names become the output columns, so you name them directly (the `"result"` column above matches `RowTypeInfo([...], ["result"])`).
- **Java**: `toTable(Schema)` exposes the agent output as a single anonymous column named `f0` (internally it calls `StreamTableEnvironment.fromDataStream(DataStream<Object>, schema)`), so the `Schema` must reference `f0` — wrap it in a `ROW(...)` only when the agent emits a composite row.
In Java, `toDataStream(Class)` / `toTable(Class)` are conveniences that derive the `TypeInformation` from the class; the `TypeInformation` overloads remain for generic or custom types.

{{< hint info >}}
In Python, `to_table()` currently requires both `Schema` and `TypeInformation`; we plan to support providing only one of them in the future.
{{< /hint >}}
The schema-only `to_table(schema=...)` / `toTable(Schema)` stays available when you have a physical schema but no output type to declare.

For complete, runnable examples, see [`WorkflowMultipleAgentExample.java`](https://github.com/apache/flink-agents/blob/main/examples/src/main/java/org/apache/flink/agents/examples/WorkflowMultipleAgentExample.java) (Java) and [`workflow_multiple_agent_example.py`](https://github.com/apache/flink-agents/blob/main/python/flink_agents/examples/quickstart/workflow_multiple_agent_example.py) (Python).
Each agent output element is adapted to the declared row type by matching columns by name — from a POJO's getters or public fields, a `Row`, or a `Map`. A scalar output maps to a single-column row.
Original file line number Diff line number Diff line change
Expand Up @@ -187,4 +187,38 @@ public static void generateOutput(Event event, RunnerContext ctx) throws Excepti
ctx.sendEvent(new OutputEvent(output));
}
}

/** A structured output type used to exercise the typed output terminals. */
public static class TypedOutput {
public int value;
public String label;

public TypedOutput() {}

public TypedOutput(int value, String label) {
this.value = value;
this.label = label;
}

@Override
public String toString() {
return "TypedOutput{value=" + value + ", label='" + label + "'}";
}
}

/** Agent emitting a structured {@link TypedOutput} derived from each integer input. */
public static class TypedOutputAgent extends Agent {
/**
* Turns each integer input into a deterministic structured output.
*
* @param event the input event carrying an integer id.
* @param ctx the runner context for sending the output event.
*/
@Action(EventType.InputEvent)
public static void processInput(Event event, RunnerContext ctx) throws Exception {
InputEvent inputEvent = InputEvent.fromEvent(event);
int id = (Integer) inputEvent.getInput();
ctx.sendEvent(new OutputEvent(new TypedOutput(id * 10, "item-" + id)));
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
package org.apache.flink.agents.integration.test;

import org.apache.flink.agents.api.AgentsExecutionEnvironment;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.java.functions.KeySelector;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
Expand Down Expand Up @@ -181,6 +182,79 @@ public void testFromDataStreamToTable() throws Exception {
checkResult(results, "test_from_datastream_to_table.txt");
}

@Test
public void testToDataStreamWithTypeInformation() throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);

DataStream<Integer> inputStream = env.fromData(1, 2, 3);
AgentsExecutionEnvironment agentsEnv =
AgentsExecutionEnvironment.getExecutionEnvironment(env);

DataStream<FlinkIntegrationAgent.TypedOutput> outputStream =
agentsEnv
.fromDataStream(inputStream)
.apply(new FlinkIntegrationAgent.TypedOutputAgent())
.toDataStream(TypeInformation.of(FlinkIntegrationAgent.TypedOutput.class));

CloseableIterator<FlinkIntegrationAgent.TypedOutput> results = outputStream.collectAsync();
agentsEnv.execute();

List<String> actual = new ArrayList<>();
while (results.hasNext()) {
actual.add(results.next().toString());
}
actual.sort(Comparator.naturalOrder());

List<String> expected =
new ArrayList<>(
List.of(
"TypedOutput{value=10, label='item-1'}",
"TypedOutput{value=20, label='item-2'}",
"TypedOutput{value=30, label='item-3'}"));
expected.sort(Comparator.naturalOrder());

Assertions.assertEquals(expected, actual);
}

@Test
public void testToTableWithTypeInformation() throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);

DataStream<Integer> inputStream = env.fromData(1, 2, 3);
AgentsExecutionEnvironment agentsEnv =
AgentsExecutionEnvironment.getExecutionEnvironment(env);

Table outputTable =
agentsEnv
.fromDataStream(inputStream)
.apply(new FlinkIntegrationAgent.TypedOutputAgent())
.toTable(TypeInformation.of(FlinkIntegrationAgent.TypedOutput.class));

// The schema is derived from the POJO type; locate columns by name.
List<String> columns = outputTable.getResolvedSchema().getColumnNames();
int valueIdx = columns.indexOf("value");
int labelIdx = columns.indexOf("label");
Assertions.assertTrue(valueIdx >= 0 && labelIdx >= 0, "expected value and label columns");

List<String> actual = collectValueLabel(outputTable, valueIdx, labelIdx);
Assertions.assertEquals(List.of("10:item-1", "20:item-2", "30:item-3"), actual);
}

private List<String> collectValueLabel(Table table, int valueIdx, int labelIdx)
throws Exception {
List<String> actual = new ArrayList<>();
try (CloseableIterator<Row> results = table.execute().collect()) {
while (results.hasNext()) {
Row row = results.next();
actual.add(row.getField(valueIdx) + ":" + row.getField(labelIdx));
}
}
actual.sort(Comparator.naturalOrder());
return actual;
}

private void checkResult(CloseableIterator<?> results, String fileName) throws IOException {
String path =
Objects.requireNonNull(
Expand Down
Loading
Loading