Installation

Workflow Agent

qianmoQqianmoQ· 更新于 2026-10-08· 阅读 37 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

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.

This quickstart introduces two small, progressive streaming examples that demonstrate how to build LLM-powered workflows with Flink Agents:

  • Review Analysis: Processes a stream of product reviews and uses a single agent to extract a rating (1–5) and unsatisfied reasons from each review.
  • Product Improvement Suggestions: Builds on the first example by aggregating per-review analysis in windows to produce product-level summaries (score distribution and common complaints), then applies a second agent to generate concrete improvement suggestions for each product.

Together, these examples show how to build a multi-agent workflow with Flink Agents and run it on a Flink standalone cluster.

Code Walkthrough

Prepare Agents Execution Environment

Create the agents execution environment, and register the available chat model connections, which can be used by the agents, to the environment.

Python

# Set up the Flink streaming environment and the Agents execution environment.
env = StreamExecutionEnvironment.get_execution_environment()
agents_env = AgentsExecutionEnvironment.get_execution_environment(env)

# Add Ollama chat model connection to be used by the ReviewAnalysisAgent
# and ProductSuggestionAgent.
agents_env.add_resource(
    "ollama_server",
    ResourceType.CHAT_MODEL_CONNECTION,
    ollama_server_descriptor,
)

Java

// Set up the Flink streaming environment and the Agents execution environment.
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
AgentsExecutionEnvironment agentsEnv =
        AgentsExecutionEnvironment.getExecutionEnvironment(env);

// Add Ollama chat model connection to be used by the ReviewAnalysisAgent.
agentsEnv.addResource(
        "ollamaChatModelConnection",
        ResourceType.CHAT_MODEL_CONNECTION,
        CustomTypesAndResources.OLLAMA_SERVER_DESCRIPTOR);

Create the Agents

Below is the example code for the ReviewAnalysisAgent, which is used to analyze the product reviews and generate a satisfaction score and potential reasons for dissatisfaction. It demonstrates how to define the prompt, tool, chat model, and action for the agent. Also, it shows how to process the chat response and send the output event. For more details, please refer to the Workflow Agent documentation.

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.EVENT_TYPE)
    @staticmethod
    def process_input(event: Event, ctx: RunnerContext) -> None:
        """Process input event and send chat request for review analysis."""
        input = ProductReview.model_validate(InputEvent.from_event(event).input)
        ctx.short_term_memory.set("id", input.id)

        content = f"""
            "id": {input.id},
            "review": {input.review}
        """
        msg = ChatMessage(role=MessageRole.USER)
        ctx.send_event(
            ChatRequestEvent(
                model="review_analysis_model",
                messages=[msg],
                prompt_args={"input": content},
            )
        )

    @action(ChatResponseEvent.EVENT_TYPE)
    @staticmethod
    def process_chat_response(event: Event, ctx: RunnerContext) -> None:
        """Process chat response event and send output event."""
        chat_response = ChatResponseEvent.from_event(event)
        try:
            json_content = json.loads(chat_response.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 {chat_response.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(listenEventTypes = {InputEvent.EVENT_TYPE})
    public static void processInput(Event event, RunnerContext ctx) throws Exception {
        InputEvent inputEvent = InputEvent.fromEvent(event);
        String input = (String) inputEvent.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, "");

        ctx.sendEvent(
                new ChatRequestEvent(
                        "reviewAnalysisModel", List.of(msg), Map.of("input", content), null));
    }

    @Action(listenEventTypes = {ChatResponseEvent.EVENT_TYPE})
    public static void processChatResponse(Event event, RunnerContext ctx)
            throws Exception {
        ChatResponseEvent chatResponse = ChatResponseEvent.fromEvent(event);
        JsonNode jsonNode = MAPPER.readTree(chatResponse.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)));
    }
}

The code for the ProductSuggestionAgent, which is used to generate product improvement suggestions based on the aggregated analysis results, is similar to the ReviewAnalysisAgent.

Integrate the Agents with Flink

Create the input DataStream by reading the product reviews from a text file as a streaming source, and use the ReviewAnalysisAgent to analyze the product reviews and generate the result DataStream. Finally print the result DataStream.

Python

# Read product reviews from a text file as a streaming source.
# Each line in the file should be a JSON string representing a ProductReview.
product_review_stream = env.from_source(
    # Target the single file, not the resources/ dir: Flink's enumerator
    # recurses, so files like skills/SKILL.md would be parsed as reviews.
    source=FileSource.for_record_stream_format(
        StreamFormat.text_line_format(),
        f"file:///{current_dir}/resources/product_review.txt",
    )
    .monitor_continuously(Duration.of_minutes(1))
    .build(),
    watermark_strategy=WatermarkStrategy.no_watermarks(),
    source_name="streaming_agent_example",
).map(
    lambda x: ProductReview.model_validate_json(
        x
    )  # Deserialize JSON to ProductReview.
)

# Use the ReviewAnalysisAgent to analyze each product review.
review_analysis_res_stream = (
    agents_env.from_datastream(
        input=product_review_stream, key_selector=lambda x: x.id
    )
    .apply(ReviewAnalysisAgent())
    .to_datastream()
)

# Print the analysis results to stdout.
review_analysis_res_stream.print()

# Execute the Flink pipeline with the Flink job name.
agents_env.execute("Workflow Agent Example Job")

Java

// Read product reviews from input_data.txt file as a streaming source.
// Each element represents a ProductReview.
DataStream<String> productReviewStream =
       env.fromSource(
               FileSource.forRecordStreamFormat(
                               new TextLineInputFormat(),
                               new Path(inputDataFile.getAbsolutePath()))
                       .build(),
               WatermarkStrategy.noWatermarks(),
               "streaming-agent-example");

// Use the ReviewAnalysisAgent to analyze each product review.
DataStream<Object> reviewAnalysisResStream =
       agentsEnv
               .fromDataStream(productReviewStream)
               .apply(new ReviewAnalysisAgent())
               .toDataStream();

// Print the analysis results to stdout.
reviewAnalysisResStream.print();

// Execute the Flink pipeline with the Flink job name.
agentsEnv.execute("Workflow Agent Example Job");

Read Input via the Table API

The multiple-agent example reads product reviews through the Flink Table API instead of a DataStream. Because Table rows arrive as Row (Java) / dict (Python) instead of a ProductReview POJO, it uses a dedicated TableReviewAnalysisAgent — identical to ReviewAnalysisAgent except its process_input action reads the id and review columns out of the row, and it ships a key selector that extracts the key from that row.

Python

class TableKeySelector(KeySelector):
    """Extract the partition key from a Table row (dict)."""

    def get_key(self, value: Any) -> str:
        return str(value["id"])

class TableReviewAnalysisAgent(Agent):
    # prompt / tool / chat_model_setup are identical to ReviewAnalysisAgent.

    @action(InputEvent.EVENT_TYPE)
    @staticmethod
    def process_input(event: Event, ctx: RunnerContext) -> None:
        # Table input arrives as a dict keyed by column name, not a POJO.
        input_dict = InputEvent.from_event(event).input
        product_id = str(input_dict["id"])
        review_text = str(input_dict["review"])
        ...

Java

public class TableReviewAnalysisAgent extends Agent {

    /** Extract the partition key from a Table Row. */
    public static class RowKeySelector implements KeySelector<Object, String> {
        @Override
        public String getKey(Object value) {
            return (String) ((Row) value).getField("id");
        }
    }

    // prompt / tool / chat model setup are identical to ReviewAnalysisAgent.

    @Action(listenEventTypes = {InputEvent.EVENT_TYPE})
    public static void processInput(Event event, RunnerContext ctx) {
        // Table input arrives as a Row keyed by column name, not a POJO.
        Row row = (Row) InputEvent.fromEvent(event).getInput();
        String productId = (String) row.getField("id");
        String reviewText = (String) row.getField("review");
        ...
    }
}

Register the input file as a table and feed it to the agent with from_table / fromTable, providing the key selector:

Python

input_table = t_env.from_path("product_reviews")

review_analysis_res_stream = (
    agents_env.from_table(input=input_table, key_selector=TableKeySelector())
    .apply(TableReviewAnalysisAgent())
    .to_datastream()
)

Java

Table inputTable = tableEnv.from("product_reviews");

DataStream<Object> reviewAnalysisResStream =
        agentsEnv
                .fromTable(inputTable, new TableReviewAnalysisAgent.RowKeySelector())
                .apply(new TableReviewAnalysisAgent())
                .toDataStream();

See Integrate with Flink for the full Table API integration reference.

Run the Example

Prerequisites

  • Unix-like environment (we use Linux, Mac OS X, Cygwin, WSL)
  • Git
  • Java 11+
  • Python 3.10, 3.11 or 3.12

Preparation

Prepare Flink and Flink Agents

Follow the installation instructions to setup Flink and the Flink Agents.

Clone the Flink Agents Repository (if not done already)

git clone https://github.com/apache/flink-agents.git
cd flink-agents

For python examples, you can skip this step and submit the python file in installed flink-agents wheel.

Deploy a Standalone Flink Cluster

You can deploy a standalone Flink cluster in your local environment with the following command.

Python

export PYTHONPATH=$(python -c 'import sysconfig; print(sysconfig.get_paths()["purelib"])')
$FLINK_HOME/bin/start-cluster.sh

Java

  1. Build Flink Agents from source to generate example jar. See installation for more details.

  2. Start the Flink cluster

    $FLINK_HOME/bin/start-cluster.sh

To run example on JDK 21+, append jvm option --add-exports=java.base/jdk.internal.vm=ALL-UNNAMED to env.java.opts.all in $FLINK_HOME/conf/config.yaml before start the flink cluster.

You can refer to the local cluster instructions for more detailed step.

If you can’t navigate to the web UI at localhost:8081, you can find the reason in $FLINK_HOME/log. If the reason is port conflict, you can change the port in $FLINK_HOME/conf/config.yaml.

Prepare Ollama

Download and install Ollama from the official website.

Ollama server 0.9.0 or higher is required.

Then pull the qwen3:8b model, which is required by the quickstart examples

ollama pull qwen3:8b

Submit Flink Agents Job to Standalone Flink Cluster

Submit to Flink Cluster

Python

export PYTHONPATH=$(python -c 'import sysconfig; print(sysconfig.get_paths()["purelib"])')

# Run review analysis example
$FLINK_HOME/bin/flink run -py ./flink-agents/python/flink_agents/examples/quickstart/workflow_single_agent_example.py
# or submit the example python file in installed flink-agents wheel
$FLINK_HOME/bin/flink run -py  $PYTHONPATH/flink_agents/examples/quickstart/workflow_single_agent_example.py

# Run product suggestion example
$FLINK_HOME/bin/flink run -py ./flink-agents/python/flink_agents/examples/quickstart/workflow_multiple_agent_example.py
# or submit the example python file in installed flink-agents wheel
$FLINK_HOME/bin/flink run -py  $PYTHONPATH/flink_agents/examples/quickstart/workflow_multiple_agent_example.py

Java

# Run review analysis example
$FLINK_HOME/bin/flink run -c org.apache.flink.agents.examples.WorkflowSingleAgentExample ./flink-agents/examples/target/flink-agents-examples-$VERSION.jar

# Run product suggestion example
$FLINK_HOME/bin/flink run -c org.apache.flink.agents.examples.WorkflowMultipleAgentExample ./flink-agents/examples/target/flink-agents-examples-$VERSION.jar

Now you should see a Flink job submitted to the Flink Cluster in Flink web UI localhost:8081

After a few minutes, you can check for the output in the TaskManager output log.

评论

登录后参与评论

正在加载评论…