Workflow Agent
Overview
A workflow style agent in Flink-Agents is an agent whose reasoning and behavior are organized as a directed workflow of modular steps, called actions, connected by events. This design is inspired by the need to orchestrate complex, multi-stage tasks in a transparent, extensible, and data-centric way, leveraging Apache Flink’s streaming architecture.
In Flink-Agents, a workflow agent is defined as a class that inherits from the Agent base class. The agent’s logic is expressed as a set of actions, each of which is a function decorated with @action(EventType) in python (or a method annotated with @action(listenEvents = {}) in java). Actions consume events, perform reasoning or tool calls, and emit new events, which may trigger downstream actions. This event-driven workflow forms a directed cyclic graph of computation, where each node is an action and each edge is an event type.
A workflow agent is well-suited for scenarios where the solution requires explicit orchestration, branching, or multi-step reasoning, such as data enrichment, multi-tool pipelines, or complex business logic.
For guidance on choosing Java or Python, see Should I choose Java or Python?.
Workflow Agent Example
Python
class ReviewAnalysisAgent(Agent):
"""An agent that uses a large language model (LLM) to analyze product reviews
and generate a satisfaction score and potential reasons for dissatisfaction.
This agent receives a product review and produces a satisfaction score and a list
of reasons for dissatisfaction. It handles prompt construction, LLM interaction,
and output parsing.
"""
@prompt
@staticmethod
def review_analysis_prompt() -> Prompt:
"""Prompt for review analysis."""
return review_analysis_prompt
@tool
@staticmethod
def notify_shipping_manager(id: str, review: str) -> None:
"""Notify the shipping manager when product received a negative review due to
shipping damage.
Parameters
----------
id : str
The id of the product that received a negative review due to shipping damage
review: str
The negative review content
"""
# reuse the declared function, but for parsing the tool metadata, we write doc
# string here again.
notify_shipping_manager(id=id, review=review)
@chat_model_setup
@staticmethod
def review_analysis_model() -> ResourceDescriptor:
"""ChatModel which focus on review analysis."""
return ResourceDescriptor(
clazz=ResourceName.ChatModel.OLLAMA_SETUP,
connection="ollama_server",
model="qwen3:8b",
prompt="review_analysis_prompt",
tools=["notify_shipping_manager"],
extract_reasoning=True,
)
@action(InputEvent)
@staticmethod
def process_input(event: InputEvent, ctx: RunnerContext) -> None:
"""Process input event and send chat request for review analysis."""
input: ProductReview = event.input
ctx.short_term_memory.set("id", input.id)
content = f"""
"id": {input.id},
"review": {input.review}
"""
msg = ChatMessage(role=MessageRole.USER, extra_args={"input": content})
ctx.send_event(ChatRequestEvent(model="review_analysis_model", messages=[msg]))
@action(ChatResponseEvent)
@staticmethod
def process_chat_response(event: ChatResponseEvent, ctx: RunnerContext) -> None:
"""Process chat response event and send output event."""
try:
json_content = json.loads(event.response.content)
ctx.send_event(
OutputEvent(
output=ProductReviewAnalysisRes(
id=ctx.short_term_memory.get("id"),
score=json_content["score"],
reasons=json_content["reasons"],
)
)
)
except Exception:
logging.exception(
f"Error processing chat response {event.response.content}"
)
# To fail the agent, you can raise an exception here.Java
/**
* An agent that uses a large language model (LLM) to analyze product reviews and generate a
* satisfaction score and potential reasons for dissatisfaction.
*
* <p>This agent receives a product review and produces a satisfaction score and a list of reasons
* for dissatisfaction. It handles prompt construction, LLM interaction, and output parsing.
*/
public class ReviewAnalysisAgent extends Agent {
private static final ObjectMapper MAPPER = new ObjectMapper();
@Prompt
public static org.apache.flink.agents.api.prompt.Prompt reviewAnalysisPrompt() {
return REVIEW_ANALYSIS_PROMPT;
}
@ChatModelSetup
public static ResourceDescriptor reviewAnalysisModel() {
return ResourceDescriptor.Builder.newBuilder(ResourceName.ChatModel.OLLAMA_SETUP)
.addInitialArgument("connection", "ollamaChatModelConnection")
.addInitialArgument("model", "qwen3:8b")
.addInitialArgument("prompt", "reviewAnalysisPrompt")
.addInitialArgument("tools", Collections.singletonList("notifyShippingManager"))
.addInitialArgument("extract_reasoning", true)
.build();
}
/**
* Tool for notifying the shipping manager when product received a negative review due to
* shipping damage.
*
* @param id The id of the product that received a negative review due to shipping damage
* @param review The negative review content
*/
@Tool(
description =
"Notify the shipping manager when product received a negative review due to shipping damage.")
public static void notifyShippingManager(
@ToolParam(name = "id") String id, @ToolParam(name = "review") String review) {
CustomTypesAndResources.notifyShippingManager(id, review);
}
/** Process input event and send chat request for review analysis. */
@Action(listenEvents = {InputEvent.class})
public static void processInput(InputEvent event, RunnerContext ctx) throws Exception {
String input = (String) event.getInput();
MAPPER.configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false);
CustomTypesAndResources.ProductReview inputObj =
MAPPER.readValue(input, CustomTypesAndResources.ProductReview.class);
ctx.getShortTermMemory().set("id", inputObj.getId());
String content =
String.format(
"{\n" + "\"id\": %s,\n" + "\"review\": \"%s\"\n" + "}",
inputObj.getId(), inputObj.getReview());
ChatMessage msg = new ChatMessage(MessageRole.USER, "", Map.of("input", content));
ctx.sendEvent(new ChatRequestEvent("reviewAnalysisModel", List.of(msg)));
}
@Action(listenEvents = ChatResponseEvent.class)
public static void processChatResponse(ChatResponseEvent event, RunnerContext ctx)
throws Exception {
JsonNode jsonNode = MAPPER.readTree(event.getResponse().getContent());
JsonNode scoreNode = jsonNode.findValue("score");
JsonNode reasonsNode = jsonNode.findValue("reasons");
if (scoreNode == null || reasonsNode == null) {
throw new IllegalStateException(
"Invalid response from LLM: missing 'score' or 'reasons' field.");
}
List<String> result = new ArrayList<>();
if (reasonsNode.isArray()) {
for (JsonNode node : reasonsNode) {
result.add(node.asText());
}
}
ctx.sendEvent(
new OutputEvent(
new CustomTypesAndResources.ProductReviewAnalysisRes(
ctx.getShortTermMemory().get("id").getValue().toString(),
scoreNode.asInt(),
result)));
}
}Action
An action is a piece of code that can be executed. Each action listens to at least one type of event. When an event of the listening type occurs, the action will be triggered. An action can also generate new events, to trigger other actions.
To declare an action in Agent, user can use @action to decorate a function of Agent class in python (or annotate a method of Agent class in java), and declare the listened event types as decorator/annotation parameters.
The decorated/annotated function signature should be (Event, RunnerContext) -> None. In Python, actions can also be defined as async def when using async execution (see Async Execution).
Python
class ReviewAnalysisAgent(Agent):
@action(InputEvent)
@staticmethod
def process_input(event: InputEvent, ctx: RunnerContext) -> None:
# the action logicJava
public class ReviewAnalysisAgent extends Agent {
/** Process input event and send chat request for review analysis. */
@Action(listenEvents = {InputEvent.class})
public static void processInput(InputEvent event, RunnerContext ctx) throws Exception {
// the action logic
}
}In the function, user can also send new events, to trigger other actions, or output the data.
Python
@action(InputEvent)
@staticmethod
def process_input(event: InputEvent, ctx: RunnerContext) -> None:
# send ChatRequestEvent
ctx.send_event(ChatRequestEvent(model=xxx, messages=xxx))
# output data to downstream
ctx.send_event(OutputEvent(output=xxx))Java
@Action(listenEvents = {InputEvent.class})
public static void processInput(InputEvent event, RunnerContext ctx) throws Exception {
// send ChatRequestEvent
ctx.sendEvent(new ChatRequestEvent("my_model", messages));
// output data to downstream
ctx.sendEvent(new OutputEvent(xxx));
}Durable Execution
Use durable execution when you wrap a time-consuming or side-effecting operation. The framework persists the result and replays it on recovery when the same call is encountered, so the function will not be called again and side effects are avoided. Action code outside durable_execute / durable_execute_async is always re-executed during recovery.
Constraints:
- The function must be deterministic and called in the same order on recovery.
- Access to Memory and
send_eventis prohibited inside the function/callable. - Arguments and results must be serializable.
Durable execution requires an external action state store. See Exactly-Once Action Consistency on how to setup and configure the external action state store.
Best-effort replay:
- Results may not be reused if call order or arguments change (non-deterministic actions), which clears subsequent cached results and re-executes.
- If a failure happens after a function completes but before its result is persisted, the call will be re-executed.
- In Python async actions, if
ctx.durable_execute_async(...)is not awaited, the result is not recorded and cannot be replayed.
Python
Python actions can call ctx.durable_execute(...) to run a synchronous durable code block.
@action(InputEvent)
@staticmethod
def process_input(event: InputEvent, ctx: RunnerContext) -> None:
def slow_external_call(data: str) -> str:
time.sleep(2)
return f"Processed: {data}"
# Synchronous durable execution
result = ctx.durable_execute(slow_external_call, event.input)
ctx.send_event(OutputEvent(output=result))Java
Java actions use DurableCallable<T> with ctx.durableExecute(...), where getId() must be stable and getResultClass() supports recovery deserialization.
@Action(listenEvents = {InputEvent.class})
public static void processInput(InputEvent event, RunnerContext ctx) throws Exception {
DurableCallable<String> call = new DurableCallable<>() {
@Override
public String getId() {
// Stable, deterministic ID for this call
return "slow_external_call";
}
@Override
public Class<String> getResultClass() {
return String.class;
}
@Override
public String call() throws Exception {
Thread.sleep(2000);
return "Processed: " + event.getInput();
}
};
String result = ctx.durableExecute(call);
ctx.sendEvent(new OutputEvent(result));
}Async Execution
Async execution uses the same durable semantics but yields while waiting for a thread-pool task. This is useful for high-latency I/O.
Python
Define an async def action and await ctx.durable_execute_async(...).
@action(InputEvent)
@staticmethod
async def process_with_async(event: InputEvent, ctx: RunnerContext) -> None:
def slow_external_call(data: str) -> str:
time.sleep(2)
return f"Processed: {data}"
result = await ctx.durable_execute_async(slow_external_call, event.input)
ctx.send_event(OutputEvent(output=result))Python async actions only support
await ctx.durable_execute_async(...). Standard asyncio functions likeasyncio.gather,asyncio.wait,asyncio.create_task, andasyncio.sleepare NOT supported because there is no asyncio event loop.
Java
Use ctx.durableExecuteAsync(DurableCallable); on JDK 21+ it yields using Continuation, and on JDK < 21 it falls back to synchronous execution.
@Action(listenEvents = {InputEvent.class})
public static void processInput(InputEvent event, RunnerContext ctx) throws Exception {
DurableCallable<String> call = new DurableCallable<>() {
@Override
public String getId() {
return "slow_external_call";
}
@Override
public Class<String> getResultClass() {
return String.class;
}
@Override
public String call() throws Exception {
Thread.sleep(2000);
return "Processed: " + event.getInput();
}
};
String result = ctx.durableExecuteAsync(call);
ctx.sendEvent(new OutputEvent(result));
}To use async execution on JDK 21+, user should append jvm option
--add-exports=java.base/jdk.internal.vm=ALL-UNNAMEDto env.java.opts.all before start the flink cluster.
Event
Events are messages passed between actions. Events may carry payloads. A single event may trigger multiple actions if they are all listening to its type.
There are 2 special types of event.
InputEvent: Generated by the framework, carrying an input data record that arrives at the agent ininputfield . Actions listening to theInputEventwill be the entry points of agent.OutputEvent: The framework will listen toOutputEvent, and convert its payload inoutputfield into outputs of the agent. By generatingOutputEvent, actions can emit output data.
User can define own event by extends Event.
Python
class MyEvent(Event):
value: AnyJava
public class MyEvent extends Event {
private Object value;
}Then, user can define actions listen to or send MyEvent.
The payload of python
Eventshould beBaseModelserializable, of javaEventshould be json serializable.
Built-in Events and Actions
There are several built-in Event and Action in Flink-Agents:
- See Chat Models for how to chat with a LLM leveraging built-in action and events.
- See Tool Use for how to programmatically use a tool leveraging built-in action and events.
评论
登录后参与评论
KnowForge