This is an automated email from the ASF dual-hosted git repository.
Sxnan pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/flink-agents-demos.git
The following commit(s) were added to refs/heads/main by this push:
new b0cda72 Update Flink Agents to 0.3.1 (#8)
b0cda72 is described below
commit b0cda72ee378db0314beef882d8d00ee4c82d752
Author: yunfengzhou-hub <[email protected]>
AuthorDate: Thu Aug 6 17:46:40 2026 +0800
Update Flink Agents to 0.3.1 (#8)
AI-Contributed/Feature: 0/133
AI-Contributed/UT: 0/0
---
flink-operations-agent-demo/README.md | 2 +-
.../bin/internal/setup_flink.sh | 26 +++++++-
.../custom_types_and_prompts.py | 25 ++++++-
.../operations-agent-job/operations_agent.py | 78 +++++++++++++---------
.../operations-agent-job/pyproject.toml | 2 +-
5 files changed, 98 insertions(+), 35 deletions(-)
diff --git a/flink-operations-agent-demo/README.md
b/flink-operations-agent-demo/README.md
index 7f0ecac..dc75119 100644
--- a/flink-operations-agent-demo/README.md
+++ b/flink-operations-agent-demo/README.md
@@ -62,7 +62,7 @@ Finally, the system outputs a complete **operations record**
and feeds it back t
- **DashScope API Key**: For AI-powered diagnosis
- Sign up at https://dashscope.aliyun.com
- Set environment variable: `export DASHSCOPE_API_KEY=your_api_key_here`
- - If you prefer not to use Tongyi, you can modify the chat model section
in `operations-agent-job/operations_agent.py`. We didn't use Ollama by default
because smaller models have poor demo performance, while larger models run
slowly on personal computers. For more chat model options, see the [Flink
Agents
Documentation](https://nightlies.apache.org/flink/flink-agents-docs-release-0.2/docs/development/chat_models/).
+ - If you prefer not to use Tongyi, you can modify the chat model section
in `operations-agent-job/operations_agent.py`. We didn't use Ollama by default
because smaller models have poor demo performance, while larger models run
slowly on personal computers. For more chat model options, see the [Flink
Agents
Documentation](https://nightlies.apache.org/flink/flink-agents-docs-release-0.3/docs/development/chat_models/).
## Quick Start
diff --git a/flink-operations-agent-demo/bin/internal/setup_flink.sh
b/flink-operations-agent-demo/bin/internal/setup_flink.sh
index ff70ca7..0adf933 100755
--- a/flink-operations-agent-demo/bin/internal/setup_flink.sh
+++ b/flink-operations-agent-demo/bin/internal/setup_flink.sh
@@ -71,15 +71,37 @@ else
fi
# Create a virtual environment in a directory named 'venv'
+# flink-agents 0.3.x requires Python >=3.10,<3.13; apache-flink 1.20 supports
up to 3.11
if [ ! -d "venv" ]; then
- python3 -m venv venv
+ PYTHON_BIN=""
+ for candidate in python3.11 python3.10 python3; do
+ if command -v "$candidate" >/dev/null 2>&1; then
+ if "$candidate" -c 'import sys; sys.exit(0 if (3,10) <=
sys.version_info[:2] <= (3,11) else 1)'; then
+ PYTHON_BIN="$candidate"
+ break
+ fi
+ fi
+ done
+ if [ -z "$PYTHON_BIN" ]; then
+ echo "Error: No suitable Python interpreter found (requires Python
3.10 or 3.11)"
+ exit 1
+ fi
+ echo "Creating virtual environment with $PYTHON_BIN ($($PYTHON_BIN
--version))"
+ "$PYTHON_BIN" -m venv venv
fi
# Activate the virtual environment
# On Linux/macOS:
source venv/bin/activate
-pip install "flink-agents==0.2.1" "apache-flink==1.20.3" "elasticsearch~=8.19"
"setuptools>=75.3,<82"
+# Constrain setuptools in PEP 517 build environments (apache-beam sdist build
+# imports pkg_resources, which was removed in setuptools>=81)
+BUILD_CONSTRAINT_FILE="$(mktemp)"
+echo "setuptools<81" > "$BUILD_CONSTRAINT_FILE"
+
+PIP_CONSTRAINT="$BUILD_CONSTRAINT_FILE" pip install "flink-agents==0.3.1"
"apache-flink==1.20.3" "elasticsearch~=8.19" "setuptools>=75.3,<82"
+
+rm -f "$BUILD_CONSTRAINT_FILE"
# Set PYTHONPATH to your Python site-packages directory
export PYTHONPATH=$(python -c 'import sysconfig;
print(sysconfig.get_paths()["purelib"])')
diff --git
a/flink-operations-agent-demo/operations-agent-job/custom_types_and_prompts.py
b/flink-operations-agent-demo/operations-agent-job/custom_types_and_prompts.py
index c7c6240..fee81cf 100644
---
a/flink-operations-agent-demo/operations-agent-job/custom_types_and_prompts.py
+++
b/flink-operations-agent-demo/operations-agent-job/custom_types_and_prompts.py
@@ -16,6 +16,8 @@
# limitations under the License.
#################################################################################
+from typing import ClassVar
+
from flink_agents.api.chat_message import ChatMessage, MessageRole
from flink_agents.api.events.event import Event
from flink_agents.api.prompts.prompt import Prompt
@@ -257,4 +259,25 @@ class ProblemRemedyResult(BaseModel):
class ProblemRemedyRequestEvent(Event):
"""Job Adjust Request Event."""
- diagnoses_result: str
+
+ EVENT_TYPE: ClassVar[str] = "problem_remedy_request_event"
+
+ def __init__(self, diagnoses_result: str) -> None:
+ """Create a ProblemRemedyRequestEvent with the diagnosis result."""
+ super().__init__(
+ type=ProblemRemedyRequestEvent.EVENT_TYPE,
+ attributes={"diagnoses_result": diagnoses_result},
+ )
+
+ @classmethod
+ def from_event(cls, event: Event) -> "ProblemRemedyRequestEvent":
+ """Reconstruct a typed ProblemRemedyRequestEvent from a base Event."""
+ assert "diagnoses_result" in event.attributes
+ result =
ProblemRemedyRequestEvent(diagnoses_result=event.attributes["diagnoses_result"])
+ result.id = event.id
+ return result
+
+ @property
+ def diagnoses_result(self) -> str:
+ """Return the diagnosis result."""
+ return self.get_attr("diagnoses_result")
diff --git
a/flink-operations-agent-demo/operations-agent-job/operations_agent.py
b/flink-operations-agent-demo/operations-agent-job/operations_agent.py
index ffa3aff..554dfb0 100644
--- a/flink-operations-agent-demo/operations-agent-job/operations_agent.py
+++ b/flink-operations-agent-demo/operations-agent-job/operations_agent.py
@@ -132,7 +132,7 @@ class FlinkJobOperationsAgent(Agent):
"""EmbeddingModelConnection responsible for ollama model service
connection."""
return ResourceDescriptor(
clazz=ResourceName.EmbeddingModel.OLLAMA_CONNECTION,
- host="http://localhost:11434",
+ base_url="http://localhost:11434",
)
@embedding_model_setup
@@ -158,7 +158,7 @@ class FlinkJobOperationsAgent(Agent):
dims=768,
)
- @action(InputEvent, ChatResponseEvent)
+ @action(InputEvent.EVENT_TYPE, ChatResponseEvent.EVENT_TYPE)
@staticmethod
def simple_problem_identification(event: Event, ctx: RunnerContext) ->
None:
"""Perform simple problem identification for job diagnosis.
@@ -168,11 +168,13 @@ class FlinkJobOperationsAgent(Agent):
ctx: Runner context for sending events
"""
# Extract job information from input
- if isinstance(event, InputEvent):
- job_info: JobInfo = event.input
- if job_info is None:
+ if event.get_type() == InputEvent.EVENT_TYPE:
+ raw_input = InputEvent.from_event(event).input
+ if raw_input is None:
logging.error("job_info is None, skipping")
return
+ # Events are JSON-serialized in transit, so the input may arrive
as a plain dict
+ job_info: JobInfo = JobInfo.model_validate(raw_input) if
isinstance(raw_input, dict) else raw_input
logging.info(f"🔍 Starting diagnosis: {job_info}")
# Store job information in short-term memory
@@ -184,16 +186,23 @@ class FlinkJobOperationsAgent(Agent):
# Send chat request for AI analysis - let LLM decide which tools
to use
logging.info("🤖 Requesting AI problem identification with tool
selection...")
- msg = ChatMessage(role=MessageRole.USER, extra_args={"job_info":
job_info.model_dump_json()})
-
ctx.send_event(ChatRequestEvent(model="problem_identification_chat_model",
messages=[msg]))
- elif isinstance(event, ChatResponseEvent):
+ msg = ChatMessage(role=MessageRole.USER)
+ ctx.send_event(
+ ChatRequestEvent(
+ model="problem_identification_chat_model",
+ messages=[msg],
+ prompt_args={"job_info": job_info.model_dump_json()},
+ )
+ )
+ elif event.get_type() == ChatResponseEvent.EVENT_TYPE:
+ response = ChatResponseEvent.from_event(event).response
sop = ctx.sensory_memory.get("sop")
problem_diagnosis_res =
ctx.sensory_memory.get("problem_diagnosis_res")
if sop is None and problem_diagnosis_res is None:
- ctx.sensory_memory.set("problem_identification_res",
event.response.content)
- logging.info(f"SOP is not set, retrieving SOP with query:
{event.response.content}")
+ ctx.sensory_memory.set("problem_identification_res",
response.content)
+ logging.info(f"SOP is not set, retrieving SOP with query:
{response.content}")
try:
- result =
ProblemIdentificationResult.model_validate_json(event.response.content)
+ result =
ProblemIdentificationResult.model_validate_json(response.content)
if result.has_issue:
ctx.send_event(
ContextRetrievalRequestEvent(
@@ -224,11 +233,11 @@ class FlinkJobOperationsAgent(Agent):
# In production, use high-performance models and improve
error handling logic.
ctx.send_event(
ContextRetrievalRequestEvent(
- query=event.response.content,
vector_store="vector_store", max_results=1
+ query=response.content,
vector_store="vector_store", max_results=1
)
)
- @action(ContextRetrievalResponseEvent, ChatResponseEvent)
+ @action(ContextRetrievalResponseEvent.EVENT_TYPE,
ChatResponseEvent.EVENT_TYPE)
@staticmethod
def deep_problem_analysis(event: Event, ctx: RunnerContext) -> None:
"""Perform deep problem diagnosis using retrieved SOP context.
@@ -237,8 +246,8 @@ class FlinkJobOperationsAgent(Agent):
event: ContextRetrievalResponseEvent containing retrieved SOP, or
ChatResponseEvent containing AI diagnosis result
ctx: Runner context for sending output events
"""
- if isinstance(event, ContextRetrievalResponseEvent):
- sop = event.documents[0].content
+ if event.get_type() == ContextRetrievalResponseEvent.EVENT_TYPE:
+ sop =
ContextRetrievalResponseEvent.from_event(event).documents[0].content
logging.info(f"📚 SOP retrieved: {sop}")
ctx.sensory_memory.set("sop", sop)
@@ -253,11 +262,14 @@ class FlinkJobOperationsAgent(Agent):
# Send chat request for AI analysis - let LLM decide which
tools to use
logging.info("🤖 Requesting AI analysis with tool selection...")
- msg = ChatMessage(
- role=MessageRole.USER,
- extra_args={"job_info": job_info, "sop": sop,
"problem_identification_res": problem_identification_res},
+ msg = ChatMessage(role=MessageRole.USER)
+ ctx.send_event(
+ ChatRequestEvent(
+ model="diagnosis_chat_model",
+ messages=[msg],
+ prompt_args={"job_info": job_info, "sop": sop,
"problem_identification_res": problem_identification_res},
+ )
)
- ctx.send_event(ChatRequestEvent(model="diagnosis_chat_model",
messages=[msg]))
except Exception as e:
error_msg = f"Diagnosis failed for job {job_info.job_id}:
{e!s}"
@@ -276,7 +288,8 @@ class FlinkJobOperationsAgent(Agent):
)
)
- elif isinstance(event, ChatResponseEvent):
+ elif event.get_type() == ChatResponseEvent.EVENT_TYPE:
+ response = ChatResponseEvent.from_event(event).response
sop = ctx.sensory_memory.get("sop")
problem_diagnosis_res =
ctx.sensory_memory.get("problem_diagnosis_res")
if sop is not None and problem_diagnosis_res is None:
@@ -288,7 +301,7 @@ class FlinkJobOperationsAgent(Agent):
logging.info(f"🤖 Received AI analysis for job: {job_id}")
- result =
ProblemDiagnosisResult.model_validate_json(event.response.content)
+ result =
ProblemDiagnosisResult.model_validate_json(response.content)
if result.need_adjustment:
ctx.send_event(
@@ -328,7 +341,7 @@ class FlinkJobOperationsAgent(Agent):
)
)
- @action(ProblemRemedyRequestEvent, ChatResponseEvent)
+ @action(ProblemRemedyRequestEvent.EVENT_TYPE, ChatResponseEvent.EVENT_TYPE)
@staticmethod
def try_problem_remedy(event: Event, ctx: RunnerContext) -> None:
"""Try to address the problem based on the diagnosis results.
@@ -337,8 +350,9 @@ class FlinkJobOperationsAgent(Agent):
event: ProblemRemedyRequestEvent containing job adjustment request
ctx: Runner context for sending output events
"""
- if isinstance(event, ProblemRemedyRequestEvent):
- ctx.sensory_memory.set("problem_diagnosis_res",
event.diagnoses_result)
+ if event.get_type() == ProblemRemedyRequestEvent.EVENT_TYPE:
+ diagnoses_result =
ProblemRemedyRequestEvent.from_event(event).diagnoses_result
+ ctx.sensory_memory.set("problem_diagnosis_res", diagnoses_result)
# Send chat request for AI analysis - let LLM decide which tools
to use
logging.info("🤖 Requesting AI job adjustment...")
@@ -347,12 +361,16 @@ class FlinkJobOperationsAgent(Agent):
job_name = ctx.sensory_memory.get("job_name")
job_info = JobInfo(base_url=base_url, job_id=job_id,
job_name=job_name)
- msg = ChatMessage(
- role=MessageRole.USER,
- extra_args={"job_info": job_info, "problem_diagnosis_res":
event.diagnoses_result},
+ msg = ChatMessage(role=MessageRole.USER)
+ ctx.send_event(
+ ChatRequestEvent(
+ model="remedy_chat_model",
+ messages=[msg],
+ prompt_args={"job_info": job_info,
"problem_diagnosis_res": diagnoses_result},
+ )
)
- ctx.send_event(ChatRequestEvent(model="remedy_chat_model",
messages=[msg]))
- elif isinstance(event, ChatResponseEvent):
+ elif event.get_type() == ChatResponseEvent.EVENT_TYPE:
+ response = ChatResponseEvent.from_event(event).response
problem_diagnosis_res =
ctx.sensory_memory.get("problem_diagnosis_res")
if problem_diagnosis_res is not None:
try:
@@ -363,7 +381,7 @@ class FlinkJobOperationsAgent(Agent):
logging.info(f"🤖 Received AI job adjustment result for
job: {job_id}")
- result =
ProblemRemedyResult.model_validate_json(event.response.content)
+ result =
ProblemRemedyResult.model_validate_json(response.content)
# Create output event
output_data = {
diff --git a/flink-operations-agent-demo/operations-agent-job/pyproject.toml
b/flink-operations-agent-demo/operations-agent-job/pyproject.toml
index 8e58c1f..c3e379d 100644
--- a/flink-operations-agent-demo/operations-agent-job/pyproject.toml
+++ b/flink-operations-agent-demo/operations-agent-job/pyproject.toml
@@ -24,7 +24,7 @@ requires-python = ">=3.10"
dependencies = [
"apache-flink==1.20.3",
- "flink-agents",
+ "flink-agents>=0.3.1",
"pydantic>=2.0.0",
"pyyaml>=6.0.0",
]