Graph-Agenten
Erstellen Sie komplexe, zustandsbehaftete Workflows mit LangGraph-Γ€hnlicher Orchestrierung und nativer ADK-Rust-Integration.
Γbersicht
GraphAgent ermΓΆglicht es Ihnen, Workflows als gerichtete Graphen mit Knoten und Kanten zu definieren und unterstΓΌtzt dabei:
- AgentNode: LLM-Agenten als Graph-Knoten mit benutzerdefinierten Ein-/Ausgabe-Mappern einbinden
- Zyklische Workflows: Native UnterstΓΌtzung fΓΌr Schleifen und iterative Schlussfolgerung (ReAct-Muster)
- Bedingtes Routing: Dynamisches Kanten-Routing basierend auf dem Zustand
- Zustandsverwaltung: Typisierter Zustand mit Reducern (ΓΌberschreiben, anhΓ€ngen, summieren, benutzerdefiniert)
- Checkpointing: Persistenter Zustand fΓΌr Fehlertoleranz und Human-in-the-Loop
- Streaming: Mehrere Stream-Modi (Werte, Aktualisierungen, Nachrichten, Debug)
Das adk-graph-Crate bietet LangGraph-artige Workflow-Orchestrierung zum Erstellen komplexer, zustandsbehafteter Agent-Workflows. Es bringt graphbasierte Workflow-Funktionen in das ADK-Rust-Γkosystem und bleibt dabei vollstΓ€ndig kompatibel mit dem Agentensystem von ADK.
Wichtige Vorteile:
- Visuelles Workflow-Design: Definieren Sie komplexe Logik als intuitive Knoten-und-Kanten-Graphen
- Parallele AusfΓΌhrung: Mehrere Knoten kΓΆnnen gleichzeitig ausgefΓΌhrt werden, fΓΌr bessere Leistung
- Zustandspersistenz: Integriertes Checkpointing fΓΌr Fehlertoleranz und Human-in-the-Loop
- LLM-Integration: Native UnterstΓΌtzung zum Einbinden von ADK-Agenten als Graph-Knoten
- Flexibles Routing: Statische Kanten, bedingtes Routing und dynamische Entscheidungsfindung
Was Sie erstellen werden
In diesem Leitfaden erstellen Sie eine Textverarbeitungspipeline, die Γbersetzung und Zusammenfassung parallel ausfΓΌhrt:
βββββββββββββββββββββββ
User Input β β
βββββββββββββββββΆ β START β
β β
ββββββββββββ¬βββββββββββ
β
βββββββββββββββββ΄ββββββββββββββββ
β β
βΌ βΌ
ββββββββββββββββββββ ββββββββββββββββββββ
β TRANSLATOR β β SUMMARIZER β
β β β β
β π«π· French β β π One sentence β
β Translation β β Summary β
βββββββββββ¬βββββββββ βββββββββββ¬βββββββββ
β β
βββββββββββββββββ¬ββββββββββββββββ
β
βΌ
βββββββββββββββββββββββ
β COMBINE β
β β
β π Merge Results β
ββββββββββββ¬βββββββββββ
β
βΌ
βββββββββββββββββββββββ
β END β
β β
β β
Complete β
βββββββββββββββββββββββ
Wichtige Konzepte:
- Knoten - Verarbeitungseinheiten, die Arbeit ausfΓΌhren (LLM-Agenten, Funktionen oder benutzerdefinierte Logik)
- Kanten - Kontrollfluss zwischen Knoten (statische Verbindungen oder bedingtes Routing)
- Zustand - Gemeinsame Daten, die durch den Graphen flieΓen und zwischen Knoten erhalten bleiben
- Parallele AusfΓΌhrung - Mehrere Knoten kΓΆnnen gleichzeitig ausgefΓΌhrt werden, fΓΌr bessere Leistung
Die Kernkomponenten verstehen
π§ Knoten: Die Arbeiter Knoten sind der Ort, an dem die eigentliche Arbeit stattfindet. Jeder Knoten kann:
- AgentNode: Einen LLM-Agenten einbinden, um natΓΌrliche Sprache zu verarbeiten
- Funktionsknoten: Benutzerdefinierten Rust-Code zur Datenverarbeitung ausfΓΌhren
- Integrierte Knoten: Vorgegebene Logik wie ZΓ€hler oder Validierer verwenden
Stellen Sie sich Knoten als spezialisierte Arbeiter in einer Montagelinie vor - jeder hat eine bestimmte Aufgabe und Expertise.
π Kanten: Die Flusssteuerung Kanten bestimmen, wie die AusfΓΌhrung durch Ihren Graphen verlΓ€uft:
- Statische Kanten: Direkte Verbindungen (
A β B β C) - Bedingte Kanten: Dynamisches Routing basierend auf dem Zustand (
if sentiment == "positive" β positive_handler) - Parallele Kanten: Mehrere Pfade von einem Knoten (
START β [translator, summarizer])
Kanten sind wie Ampeln und StraΓenschilder, die den Arbeitsfluss lenken.
πΎ Zustand: Der gemeinsame Speicher Zustand ist ein Key-Value-Speicher, aus dem alle Knoten lesen und in den sie schreiben kΓΆnnen:
- Eingabedaten: Anfangsinformationen, die in den Graphen eingespeist werden
- Zwischenergebnisse: Die Ausgabe eines Knotens wird zur Eingabe fΓΌr einen anderen
- Endausgabe: Das fertige Ergebnis nach der gesamten Verarbeitung
Zustand wirkt wie ein gemeinsames Whiteboard, auf dem Knoten Informationen fΓΌr andere hinterlassen kΓΆnnen.
β‘ Parallele AusfΓΌhrung: Der Geschwindigkeitsschub Wenn mehrere Kanten einen Knoten verlassen, werden diese Zielknoten gleichzeitig ausgefΓΌhrt:
- Schnellere Verarbeitung: UnabhΓ€ngige Aufgaben laufen gleichzeitig
- Ressourceneffizienz: Bessere Ausnutzung von CPU und I/O
- Skalierbarkeit: BewΓ€ltigen Sie komplexere Workflows ohne lineare Verlangsamung
Das ist, als wΓΌrden mehrere Arbeiter verschiedene Teile eines Auftrags gleichzeitig bearbeiten, anstatt in einer Schlange zu warten.
Schnellstart
1. Erstellen Sie Ihr Projekt
cargo new graph_demo
cd graph_demo
FΓΌgen Sie AbhΓ€ngigkeiten zu Cargo.toml hinzu:
[dependencies]
adk-graph = { version = "2.0.0", features = ["sqlite"] }
adk-agent = "2.0.0"
adk-model = "2.0.0"
adk-core = "2.0.0"
tokio = { version = "1", features = ["full"] }
dotenvy = "0.15"
serde_json = "1.0"
Erstellen Sie .env mit Ihrem API-SchlΓΌssel:
echo 'GOOGLE_API_KEY=your-api-key' > .env
2. Beispiel fΓΌr parallele Verarbeitung
Hier ist ein vollstΓ€ndiges, lauffΓ€higes Beispiel, das Text parallel verarbeitet:
use adk_agent::LlmAgentBuilder;
use adk_graph::{
agent::GraphAgent,
edge::{END, START},
node::{AgentNode, ExecutionConfig, NodeOutput},
state::State,
};
use adk_model::GeminiModel;
use serde_json::json;
use std::sync::Arc;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
dotenvy::dotenv().ok();
let api_key = std::env::var("GOOGLE_API_KEY")?;
let model = Arc::new(GeminiModel::new(&api_key, "gemini-2.5-flash")?);
// Create specialized LLM agents
let translator_agent = Arc::new(
LlmAgentBuilder::new("translator")
.description("Translates text to French")
.model(model.clone())
.instruction("Translate the input text to French. Only output the translation.")
.build()?,
);
let summarizer_agent = Arc::new(
LlmAgentBuilder::new("summarizer")
.description("Summarizes text")
.model(model.clone())
.instruction("Summarize the input text in one sentence.")
.build()?,
);
// Wrap agents as graph nodes with input/output mappers
let translator_node = AgentNode::new(translator_agent)
.with_input_mapper(|state| {
let text = state.get("input").and_then(|v| v.as_str()).unwrap_or("");
adk_core::Content::new("user").with_text(text)
})
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
for event in events {
if let Some(content) = event.content() {
let text: String = content.parts.iter()
.filter_map(|p| p.text())
.collect::<Vec<_>>()
.join("");
if !text.is_empty() {
updates.insert("translation".to_string(), json!(text));
}
}
}
updates
});
let summarizer_node = AgentNode::new(summarizer_agent)
.with_input_mapper(|state| {
let text = state.get("input").and_then(|v| v.as_str()).unwrap_or("");
adk_core::Content::new("user").with_text(text)
})
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
for event in events {
if let Some(content) = event.content() {
let text: String = content.parts.iter()
.filter_map(|p| p.text())
.collect::<Vec<_>>()
.join("");
if !text.is_empty() {
updates.insert("summary".to_string(), json!(text));
}
}
}
updates
});
// Build the graph with parallel execution
let agent = GraphAgent::builder("text_processor")
.description("Processes text with translation and summarization in parallel")
.channels(&["input", "translation", "summary", "result"])
.node(translator_node)
.node(summarizer_node)
.node_fn("combine", |ctx| async move {
let translation = ctx.get("translation").and_then(|v| v.as_str()).unwrap_or("N/A");
let summary = ctx.get("summary").and_then(|v| v.as_str()).unwrap_or("N/A");
let result = format!(
"=== Processing Complete ===\n\n\
French Translation:\n{}\n\n\
Summary:\n{}",
translation, summary
);
Ok(NodeOutput::new().with_update("result", json!(result)))
})
// Parallel execution: both nodes start simultaneously
.edge(START, "translator")
.edge(START, "summarizer")
.edge("translator", "combine")
.edge("summarizer", "combine")
.edge("combine", END)
.build()?;
// Execute the graph
let mut input = State::new();
input.insert("input".to_string(), json!("AI is transforming how we work and live."));
let result = agent.invoke(input, ExecutionConfig::new("thread-1")).await?;
println!("{}", result.get("result").and_then(|v| v.as_str()).unwrap_or(""));
Ok(())
}
Beispielausgabe:
=== Processing Complete ===
French Translation:
L'IA transforme notre faΓ§on de travailler et de vivre.
Summary:
AI is revolutionizing work and daily life through technological transformation.
Wie die GraphausfΓΌhrung funktioniert
Das Gesamtbild
Graph-Agenten fΓΌhren AusfΓΌhrungen in Super-Schritten aus - alle bereiten Knoten laufen parallel, dann wartet der Graph, bis alle abgeschlossen sind, bevor der nΓ€chste Schritt beginnt:
Step 1: START βββ¬βββΆ translator (running)
ββββΆ summarizer (running)
β³ Wait for both to complete...
Step 2: translator βββ¬βββΆ combine (running)
summarizer βββ
β³ Wait for combine to complete...
Step 3: combine βββΆ END β
Datenfluss durch Knoten
Jeder Knoten kann aus dem gemeinsamen Zustand lesen und in ihn schreiben:
βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
β STEP 1: Initial state β
βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ€
β β
β State: { "input": "AI is transforming how we work" } β
β β
β β β
β β
β ββββββββββββββββββββ ββββββββββββββββββββ β
β β translator β β summarizer β β
β β reads "input" β β reads "input" β β
β ββββββββββββββββββββ ββββββββββββββββββββ β
β β
βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
β
βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
β STEP 2: After parallel execution β
βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ€
β β
β State: { β
β "input": "AI is transforming how we work", β
β "translation": "L'IA transforme notre faΓ§on de travailler", β
β "summary": "AI is revolutionizing work through technology" β
β } β
β β
β β β
β β
β ββββββββββββββββββββββββββββββββββββββββ β
β β combine β β
β β reads "translation" + "summary" β β
β β writes "result" β β
β ββββββββββββββββββββββββββββββββββββββββ β
β β
βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
β
βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
β STEP 3: Final state β
βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ€
β β
β State: { β
β "input": "AI is transforming how we work", β
β "translation": "L'IA transforme notre faΓ§on de travailler", β
β "summary": "AI is revolutionizing work through technology", β
β "result": "=== Processing Complete ===\n\nFrench..." β
β } β
β β
βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
Was es mΓΆglich macht
| Komponente | Rolle |
|---|---|
AgentNode | UmhΓΌllt LLM Agenten mit Ein-/Ausgabe-Mappern |
input_mapper | Transformiert Zustand β Agenteneingabe Content |
output_mapper | Transformiert Agent-Ereignisse β Zustandsaktualisierungen |
channels | Deklariert Zustandsfelder, die der Graph verwenden wird |
edge() | Definiert den AusfΓΌhrungsfluss zwischen Knoten |
ExecutionConfig | Stellt die Thread-ID fΓΌr das Checkpointing bereit |
Bedingtes Routing mit LLM-Klassifizierung
Erstellen Sie intelligente Routing-Systeme, bei denen LLMs den AusfΓΌhrungspfad bestimmen:
Visualisierung: Routing basierend auf Stimmung
βββββββββββββββββββββββ
User Feedback β β
βββββββββββββββββΆ β CLASSIFIER β
β π§ Analyze tone β
ββββββββββββ¬βββββββββββ
β
βββββββββββββββββΌββββββββββββββββ
β β β
βΌ βΌ βΌ
ββββββββββββββββββββ ββββββββββββββββββββ ββββββββββββββββββββ
β POSITIVE β β NEGATIVE β β NEUTRAL β
β β β β β β
β π Thank you! β β π Apologize β β π Ask more β
β Celebrate β β Help fix β β questions β
ββββββββββββββββββββ ββββββββββββββββββββ ββββββββββββββββββββ
VollstΓ€ndiger Beispielcode
use adk_agent::LlmAgentBuilder;
use adk_graph::{
edge::{END, Router, START},
graph::StateGraph,
node::{AgentNode, ExecutionConfig},
state::State,
};
use adk_model::GeminiModel;
use serde_json::json;
use std::sync::Arc;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
dotenvy::dotenv().ok();
let api_key = std::env::var("GOOGLE_API_KEY")?;
let model = Arc::new(GeminiModel::new(&api_key, "gemini-2.5-flash")?);
// Create classifier agent
let classifier_agent = Arc::new(
LlmAgentBuilder::new("classifier")
.description("Classifies text sentiment")
.model(model.clone())
.instruction(
"You are a sentiment classifier. Analyze the input text and respond with \
ONLY one word: 'positive', 'negative', or 'neutral'. Nothing else.",
)
.build()?,
);
// Create response agents for each sentiment
let positive_agent = Arc::new(
LlmAgentBuilder::new("positive")
.description("Handles positive feedback")
.model(model.clone())
.instruction(
"You are a customer success specialist. The customer has positive feedback. \
Express gratitude, reinforce the positive experience, and suggest ways to \
share their experience. Be warm and appreciative. Keep response under 3 sentences.",
)
.build()?,
);
let negative_agent = Arc::new(
LlmAgentBuilder::new("negative")
.description("Handles negative feedback")
.model(model.clone())
.instruction(
"You are a customer support specialist. The customer has a complaint. \
Acknowledge their frustration, apologize sincerely, and offer help. \
Be empathetic. Keep response under 3 sentences.",
)
.build()?,
);
let neutral_agent = Arc::new(
LlmAgentBuilder::new("neutral")
.description("Handles neutral feedback")
.model(model.clone())
.instruction(
"You are a customer service representative. The customer has neutral feedback. \
Ask clarifying questions to better understand their needs. Be helpful and curious. \
Keep response under 3 sentences.",
)
.build()?,
);
// Create AgentNodes with mappers
let classifier_node = AgentNode::new(classifier_agent)
.with_input_mapper(|state| {
let text = state.get("feedback").and_then(|v| v.as_str()).unwrap_or("");
adk_core::Content::new("user").with_text(text)
})
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
for event in events {
if let Some(content) = event.content() {
let text: String = content.parts.iter()
.filter_map(|p| p.text())
.collect::<Vec<_>>()
.join("")
.to_lowercase()
.trim()
.to_string();
let sentiment = if text.contains("positive") { "positive" }
else if text.contains("negative") { "negative" }
else { "neutral" };
updates.insert("sentiment".to_string(), json!(sentiment));
}
}
updates
});
// Response nodes (similar pattern for each)
let positive_node = AgentNode::new(positive_agent)
.with_input_mapper(|state| {
let text = state.get("feedback").and_then(|v| v.as_str()).unwrap_or("");
adk_core::Content::new("user").with_text(text)
})
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
for event in events {
if let Some(content) = event.content() {
let text: String = content.parts.iter()
.filter_map(|p| p.text())
.collect::<Vec<_>>()
.join("");
updates.insert("response".to_string(), json!(text));
}
}
updates
});
// Build graph with conditional routing
let graph = StateGraph::with_channels(&["feedback", "sentiment", "response"])
.add_node(classifier_node)
.add_node(positive_node)
// ... add negative_node and neutral_node similarly
.add_edge(START, "classifier")
.add_conditional_edges(
"classifier",
Router::by_field("sentiment"), // Route based on sentiment field
[
("positive", "positive"),
("negative", "negative"),
("neutral", "neutral"),
],
)
.add_edge("positive", END)
.add_edge("negative", END)
.add_edge("neutral", END)
.compile()?;
// Test with different feedback
let mut input = State::new();
input.insert("feedback".to_string(), json!("Your product is amazing! I love it!"));
let result = graph.invoke(input, ExecutionConfig::new("feedback-1")).await?;
println!("Sentiment: {}", result.get("sentiment").and_then(|v| v.as_str()).unwrap_or(""));
println!("Response: {}", result.get("response").and_then(|v| v.as_str()).unwrap_or(""));
Ok(())
}
Beispielfluss:
Input: "Your product is amazing! I love it!"
β
Classifier: "positive"
β
Positive Agent: "Thank you so much for the wonderful feedback!
We're thrilled you love our product.
Would you consider leaving a review to help others?"
ReAct-Muster: Denken + Handeln
Erstellen Sie Agenten, die Tools iterativ verwenden kΓΆnnen, um komplexe Probleme zu lΓΆsen:
Visualisierung: ReAct-Zyklus
βββββββββββββββββββββββ
User Question β β
βββββββββββββββββΆ β REASONER β
β π§ Think + Act β
ββββββββββββ¬βββββββββββ
β
βΌ
βββββββββββββββββββββββ
β Has tool calls? β
β β
ββββββββββββ¬βββββββββββ
β
βββββββββββββββββ΄ββββββββββββββββ
β β
βΌ βΌ
ββββββββββββββββββββ ββββββββββββββββββββ
β YES β β NO β
β β β β
β π Loop back β β β
Final answer β
β to reasoner β β END β
βββββββββββ¬βββββββββ ββββββββββββββββββββ
β
βββββββββββββββββββ
β
βΌ
βββββββββββββββββββββββ
β REASONER β
β π§ Think + Act β
β (next iteration) β
βββββββββββββββββββββββ
VollstΓ€ndiges ReAct-Beispiel
use adk_agent::LlmAgentBuilder;
use adk_core::{Part, Tool};
use adk_graph::{
edge::{END, START},
graph::StateGraph,
node::{AgentNode, ExecutionConfig, NodeOutput},
state::State,
};
use adk_model::GeminiModel;
use adk_tool::FunctionTool;
use serde_json::json;
use std::sync::Arc;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
dotenvy::dotenv().ok();
let api_key = std::env::var("GOOGLE_API_KEY")?;
let model = Arc::new(GeminiModel::new(&api_key, "gemini-2.5-flash")?);
// Create tools
let weather_tool = Arc::new(FunctionTool::new(
"get_weather",
"Get the current weather for a location. Takes a 'location' parameter (city name).",
|_ctx, args| async move {
let location = args.get("location").and_then(|v| v.as_str()).unwrap_or("unknown");
Ok(json!({
"location": location,
"temperature": "72Β°F",
"condition": "Sunny",
"humidity": "45%"
}))
},
)) as Arc<dyn Tool>;
let calculator_tool = Arc::new(FunctionTool::new(
"calculator",
"Perform mathematical calculations. Takes an 'expression' parameter (string).",
|_ctx, args| async move {
let expr = args.get("expression").and_then(|v| v.as_str()).unwrap_or("0");
let result = match expr {
"2 + 2" => "4",
"10 * 5" => "50",
"100 / 4" => "25",
"15 - 7" => "8",
_ => "Unable to evaluate",
};
Ok(json!({ "result": result, "expression": expr }))
},
)) as Arc<dyn Tool>;
// Create reasoner agent with tools
let reasoner_agent = Arc::new(
LlmAgentBuilder::new("reasoner")
.description("Reasoning agent with tools")
.model(model.clone())
.instruction(
"You are a helpful assistant with access to tools. Use tools when needed to answer questions. \
When you have enough information, provide a final answer without using more tools.",
)
.tool(weather_tool)
.tool(calculator_tool)
.build()?,
);
// Create reasoner node that detects tool usage
let reasoner_node = AgentNode::new(reasoner_agent)
.with_input_mapper(|state| {
let question = state.get("question").and_then(|v| v.as_str()).unwrap_or("");
adk_core::Content::new("user").with_text(question)
})
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
let mut has_tool_calls = false;
let mut response = String::new();
for event in events {
if let Some(content) = event.content() {
for part in &content.parts {
match part {
Part::FunctionCall { .. } => {
has_tool_calls = true;
}
Part::Text { text } => {
response.push_str(text);
}
_ => {}
}
}
}
}
updates.insert("has_tool_calls".to_string(), json!(has_tool_calls));
updates.insert("response".to_string(), json!(response));
updates
});
// Build ReAct graph with cycle
let graph = StateGraph::with_channels(&["question", "has_tool_calls", "response", "iteration"])
.add_node(reasoner_node)
.add_node_fn("counter", |ctx| async move {
let i = ctx.get("iteration").and_then(|v| v.as_i64()).unwrap_or(0);
Ok(NodeOutput::new().with_update("iteration", json!(i + 1)))
})
.add_edge(START, "counter")
.add_edge("counter", "reasoner")
.add_conditional_edges(
"reasoner",
|state| {
let has_tools = state.get("has_tool_calls").and_then(|v| v.as_bool()).unwrap_or(false);
let iteration = state.get("iteration").and_then(|v| v.as_i64()).unwrap_or(0);
// Safety limit
if iteration >= 5 { return END.to_string(); }
if has_tools {
"counter".to_string() // Loop back for more reasoning
} else {
END.to_string() // Done - final answer
}
},
[("counter", "counter"), (END, END)],
)
.compile()?
.with_recursion_limit(10);
// Test the ReAct agent
let mut input = State::new();
input.insert("question".to_string(), json!("What's the weather in Paris and what's 15 + 25?"));
let result = graph.invoke(input, ExecutionConfig::new("react-1")).await?;
println!("Final answer: {}", result.get("response").and_then(|v| v.as_str()).unwrap_or(""));
println!("Iterations: {}", result.get("iteration").and_then(|v| v.as_i64()).unwrap_or(0));
Ok(())
}
Beispielfluss:
Question: "What's the weather in Paris and what's 15 + 25?"
Iteration 1:
Reasoner: "I need to get weather info and do math"
β Calls get_weather(location="Paris") and calculator(expression="15 + 25")
β has_tool_calls = true β Loop back
Iteration 2:
Reasoner: "Based on the results: Paris is 72Β°F and sunny, 15 + 25 = 40"
β No tool calls β has_tool_calls = false β END
Final Answer: "The weather in Paris is 72Β°F and sunny with 45% humidity.
And 15 + 25 equals 40."
AgentNode
Kapselt jeden ADK Agent (typischerweise LlmAgent) als Graphknoten:
Was der Agent sieht
Ein Agent innerhalb eines Graphen lΓ€uft unter einem Kontext, der von der Invocation abgeleitet ist, die
den Graphen gestartet hat, sodass er sich genauso verhΓ€lt wie auΓerhalb eines solchen. Wenn ein Runner
einen GraphAgent aufruft, werden die IdentitΓ€t und die Dienste des Aufrufers automatisch mitΓΌbertragen:
| Γbernommen durch | Hinweis |
|---|---|
app_name, user_id, session_id | Die des Aufrufers, nicht eine synthetische |
| Scopes und Request-Metadaten | Damit Scope-PrΓΌfungen die Grants des Aufrufers sehen |
| Geheimdienst, Speicher, Artefakte, geteilter Zustand | Genau wie auΓerhalb des Graphen verfΓΌgbar |
| Abbruch | Runner::interrupt erreicht einen Agenten, der als Knoten ausgefΓΌhrt wird |
RunConfig | Vom Aufrufer geerbt |
branch | Abgeleitet, als {caller_branch}.{agent_name}, sodass die Ereignisse eines Knotens zugeordnet werden kΓΆnnen |
Ein Graph, der direkt aufgerufen wird β graph.invoke(state, ExecutionConfig::new("thread")) β
hat keinen Aufruf, von dem er erben kann. Das ist der Standalone-Modus: Der Knoten erhΓ€lt
user_id = "graph_user", app_name = "graph_app", Branch main, keine Geheimnisse und kein
Speicher. Es ist ein bewusster Modus, um einen Graphen auΓerhalb von Runner auszufΓΌhren, nicht ein
Fallback, auf den man in der Produktion zurΓΌckgreifen sollte.
Um die Verbindung manuell herzustellen β zum Beispiel, wenn Sie einen Graphen ΓΌber Ihren eigenen Executor steuern β ΓΌbergeben Sie den Aufruf explizit:
let config = ExecutionConfig::new(ctx.session_id()).with_parent_context(ctx.clone());
Hinweis: Der Knoten lΓ€uft weiterhin in seiner eigenen In-Memory-Graph-Session, sodass der GesprΓ€chsverlauf des Agenten innerhalb eines Knotens auf den Knoten beschrΓ€nkt ist und nicht an die Session des Aufrufers angehΓ€ngt wird.
let node = AgentNode::new(llm_agent)
.with_input_mapper(|state| {
// Transform graph state to agent input Content
let text = state.get("input").and_then(|v| v.as_str()).unwrap_or("");
adk_core::Content::new("user").with_text(text)
})
.with_output_mapper(|events| {
// Transform agent events to state updates
let mut updates = std::collections::HashMap::new();
for event in events {
if let Some(content) = event.content() {
let text: String = content.parts.iter()
.filter_map(|p| p.text())
.collect::<Vec<_>>()
.join("");
updates.insert("output".to_string(), json!(text));
}
}
updates
});
Funktionsknoten
Einfache asynchrone Funktionen, die den Zustand verarbeiten:
.node_fn("process", |ctx| async move {
let input = ctx.state.get("input").unwrap();
let output = process_data(input).await?;
Ok(NodeOutput::new().with_update("output", output))
})
Kantenarten
Statische Kanten
Direkte Verbindungen zwischen Knoten:
.edge(START, "first_node")
.edge("first_node", "second_node")
.edge("second_node", END)
Bedingte Kanten
Dynamisches Routing basierend auf dem Zustand:
.conditional_edge(
"router",
|state| {
match state.get("next").and_then(|v| v.as_str()) {
Some("research") => "research_node".to_string(),
Some("write") => "write_node".to_string(),
_ => END.to_string(),
}
},
[
("research_node", "research_node"),
("write_node", "write_node"),
(END, END),
],
)
Router-Hilfen
Verwenden Sie integrierte Router fΓΌr gΓ€ngige Muster:
use adk_graph::edge::Router;
// Route based on a state field value
.conditional_edge("classifier", Router::by_field("sentiment"), [
("positive", "positive_handler"),
("negative", "negative_handler"),
("neutral", "neutral_handler"),
])
// Route based on boolean field
.conditional_edge("check", Router::by_bool("approved"), [
("true", "execute"),
("false", "reject"),
])
// Limit iterations
.conditional_edge("loop", Router::max_iterations("count", 5), [
("continue", "process"),
("done", END),
])
Parallele AusfΓΌhrung
Mehrere Kanten von einem einzelnen Knoten werden parallel ausgefΓΌhrt:
let agent = GraphAgent::builder("parallel_processor")
.channels(&["input", "translation", "summary", "analysis"])
.node(translator_node)
.node(summarizer_node)
.node(analyzer_node)
.node(combiner_node)
// All three start simultaneously
.edge(START, "translator")
.edge(START, "summarizer")
.edge(START, "analyzer")
// Wait for all to complete before combining
.edge("translator", "combiner")
.edge("summarizer", "combiner")
.edge("analyzer", "combiner")
.edge("combiner", END)
.build()?;
Zyklische Graphen (ReAct-Muster)
Erstellen Sie iterative Reasoning-Agenten mit Zyklen:
use adk_core::Part;
// Create agent with tools
let reasoner = Arc::new(
LlmAgentBuilder::new("reasoner")
.model(model)
.instruction("Use tools to answer questions. Provide final answer when done.")
.tool(search_tool)
.tool(calculator_tool)
.build()?
);
let reasoner_node = AgentNode::new(reasoner)
.with_input_mapper(|state| {
let question = state.get("question").and_then(|v| v.as_str()).unwrap_or("");
adk_core::Content::new("user").with_text(question)
})
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
let mut has_tool_calls = false;
let mut response = String::new();
for event in events {
if let Some(content) = event.content() {
for part in &content.parts {
match part {
Part::FunctionCall { name, .. } => {
has_tool_calls = true;
}
Part::Text { text } => {
response.push_str(text);
}
_ => {}
}
}
}
}
updates.insert("has_tool_calls".to_string(), json!(has_tool_calls));
updates.insert("response".to_string(), json!(response));
updates
});
// Build graph with cycle
let react_agent = StateGraph::with_channels(&["question", "has_tool_calls", "response", "iteration"])
.add_node(reasoner_node)
.add_node_fn("counter", |ctx| async move {
let i = ctx.get("iteration").and_then(|v| v.as_i64()).unwrap_or(0);
Ok(NodeOutput::new().with_update("iteration", json!(i + 1)))
})
.add_edge(START, "counter")
.add_edge("counter", "reasoner")
.add_conditional_edges(
"reasoner",
|state| {
let has_tools = state.get("has_tool_calls").and_then(|v| v.as_bool()).unwrap_or(false);
let iteration = state.get("iteration").and_then(|v| v.as_i64()).unwrap_or(0);
// Safety limit
if iteration >= 5 { return END.to_string(); }
if has_tools {
"counter".to_string() // Loop back
} else {
END.to_string() // Done
}
},
[("counter", "counter"), (END, END)],
)
.compile()?
.with_recursion_limit(10);
Multi-Agent-Supervisor
Leiten Sie Aufgaben an spezialisierte Agenten weiter:
// Create supervisor agent
let supervisor = Arc::new(
LlmAgentBuilder::new("supervisor")
.model(model.clone())
.instruction("Route tasks to: researcher, writer, or coder. Reply with agent name only.")
.build()?
);
let supervisor_node = AgentNode::new(supervisor)
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
for event in events {
if let Some(content) = event.content() {
let text: String = content.parts.iter()
.filter_map(|p| p.text())
.collect::<Vec<_>>()
.join("")
.to_lowercase();
let next = if text.contains("researcher") { "researcher" }
else if text.contains("writer") { "writer" }
else if text.contains("coder") { "coder" }
else { "done" };
updates.insert("next_agent".to_string(), json!(next));
}
}
updates
});
// Build supervisor graph
let graph = StateGraph::with_channels(&["task", "next_agent", "research", "content", "code"])
.add_node(supervisor_node)
.add_node(researcher_node)
.add_node(writer_node)
.add_node(coder_node)
.add_edge(START, "supervisor")
.add_conditional_edges(
"supervisor",
Router::by_field("next_agent"),
[
("researcher", "researcher"),
("writer", "writer"),
("coder", "coder"),
("done", END),
],
)
// Agents report back to supervisor
.add_edge("researcher", "supervisor")
.add_edge("writer", "supervisor")
.add_edge("coder", "supervisor")
.compile()?;
Zustandsverwaltung
Zustands-Schema mit Reducern
Steuern Sie, wie Zustandsaktualisierungen zusammengefΓΌhrt werden:
let schema = StateSchema::builder()
.channel("current_step") // Overwrite (default)
.list_channel("messages") // Append to list
.channel_with_reducer("count", Reducer::Sum) // Sum values
.channel_with_reducer("data", Reducer::Custom(Arc::new(|old, new| {
// Custom merge logic
merge_json(old, new)
})))
.build();
let agent = GraphAgent::builder("stateful")
.state_schema(schema)
// ... nodes and edges
.build()?;
Reducer-Typen
| Reducer | Verhalten |
|---|---|
Overwrite | Alten Wert durch neuen ersetzen (Standard) |
Append | An Liste anhΓ€ngen |
Sum | Numerische Werte hinzufΓΌgen |
Custom | Benutzerdefinierte ZusammenfΓΌhrungsfunktion |
Checkpointing
Aktiviere persistenten Zustand fΓΌr Fehlertoleranz und Human-in-the-loop:
Im Speicher (Entwicklung)
use adk_graph::checkpoint::MemoryCheckpointer;
let checkpointer = Arc::new(MemoryCheckpointer::new());
let graph = StateGraph::with_channels(&["task", "result"])
// ... nodes and edges
.compile()?
.with_checkpointer_arc(checkpointer.clone());
SQLite (Produktion)
use adk_graph::checkpoint::SqliteCheckpointer;
let checkpointer = SqliteCheckpointer::new("checkpoints.db").await?;
let graph = StateGraph::with_channels(&["task", "result"])
// ... nodes and edges
.compile()?
.with_checkpointer(checkpointer);
Was ein Checkpoint erfasst
Ein Checkpoint speichert den akkumulierten Zustand, die Schrittzahl und die Frontier β die Knoten, die noch ausgefΓΌhrt werden mΓΌssen. Er wird geschrieben, nachdem sich die Frontier weiterbewegt hat, sodass das Fortsetzen nie einen Knoten erneut ausfΓΌhrt, der bereits abgeschlossen wurde, und seine Aktualisierungen nie doppelt anwendet. Ein Lauf, der abgeschlossen wird, checkpointet eine leere Frontier, sodass das Fortsetzen eines abgeschlossenen Threads den endgΓΌltigen Zustand zurΓΌckgibt, anstatt den Graphen neu zu starten.
Zwei FΓ€lle checkpointen absichtlich die Frontier, die gerade ausgefΓΌhrt wurde, statt der nΓ€chsten, weil der unterbrochene Knoten seine Aktualisierungen noch nicht erzeugt hat und beim Fortsetzen erneut ausgefΓΌhrt werden muss:
| Situation | Frontier gespeichert |
|---|---|
| Super-Schritt abgeschlossen | Die nΓ€chsten auszufΓΌhrenden Knoten |
| AusfΓΌhrung abgeschlossen | Leer |
| Interrupt ausgelΓΆst (blockierend oder streaming) | Die Knoten, die ausgefΓΌhrt wurden |
Streamed Runs checkpointen nach demselben Zeitplan wie blockierende Runs, auch wenn ein Interrupt den Stream beendet, sodass eine Pause mit Human-in-the-Loop in beiden AusfΓΌhrungsmodi fortgesetzt werden kann.
Checkpoint-Verlauf (Zeitreise)
Nur Lesezugriff.
TimeTravelHandle::state_history(from, to)gibt den Zustand zurΓΌck, der an jedem checkpointierten Schritt gespeichert wurde. Es fΓΌhrt nichts aus β kein Node lΓ€uft, kein Event wird neu erzeugt, und kein Seiteneffekt wiederholt sich. Um ab einem Punkt in der Historie erneut auszufΓΌhren, verwendefork_at, um diesen Checkpoint zu verzweigen, und rufe den Graphen auf dem abgezweigten Thread auf. Die Methode hieΓ frΓΌherreplayund wurde als erneute AusfΓΌhrung des Graphen dokumentiert, was sie nie getan hat.
Checkpoints ermΓΆglichen auΓerdem dauerhaftes Fortsetzen β wenn eine GraphausfΓΌhrung abstΓΌrzt oder der Prozess neu startet, wird die AusfΓΌhrung vom zuletzt persistent gespeicherten Checkpoint fortgesetzt, statt von vorn zu beginnen. Verwende SqliteCheckpointer oder PostgresCheckpointer fΓΌr ausfallsichere Persistenz.
// List all checkpoints for a thread
let checkpoints = checkpointer.list("thread-id").await?;
for cp in checkpoints {
println!("Step {}: {:?}", cp.step, cp.state.get("status"));
}
// Load a specific checkpoint
if let Some(checkpoint) = checkpointer.load_by_id(&checkpoint_id).await? {
println!("State at step {}: {:?}", checkpoint.step, checkpoint.state);
}
Human-in-the-Loop
Pausiere die AusfΓΌhrung fΓΌr eine menschliche Freigabe mithilfe dynamischer Interrupts:
use adk_graph::{error::GraphError, node::NodeOutput};
// Planner agent assesses risk
let planner_node = AgentNode::new(planner_agent)
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
for event in events {
if let Some(content) = event.content() {
let text: String = content.parts.iter()
.filter_map(|p| p.text())
.collect::<Vec<_>>()
.join("");
// Extract risk level from LLM response
let risk = if text.to_lowercase().contains("risk: high") { "high" }
else if text.to_lowercase().contains("risk: medium") { "medium" }
else { "low" };
updates.insert("plan".to_string(), json!(text));
updates.insert("risk_level".to_string(), json!(risk));
}
}
updates
});
// Review node with dynamic interrupt
let graph = StateGraph::with_channels(&["task", "plan", "risk_level", "approved", "result"])
.add_node(planner_node)
.add_node(executor_node)
.add_node_fn("review", |ctx| async move {
let risk = ctx.get("risk_level").and_then(|v| v.as_str()).unwrap_or("low");
let approved = ctx.get("approved").and_then(|v| v.as_bool());
// Already approved - continue
if approved == Some(true) {
return Ok(NodeOutput::new());
}
// High/medium risk - interrupt for approval
if risk == "high" || risk == "medium" {
return Ok(NodeOutput::interrupt_with_data(
&format!("{} RISK: Human approval required", risk.to_uppercase()),
json!({
"plan": ctx.get("plan"),
"risk_level": risk,
"action": "Set 'approved' to true to continue"
})
));
}
// Low risk - auto-approve
Ok(NodeOutput::new().with_update("approved", json!(true)))
})
.add_edge(START, "planner")
.add_edge("planner", "review")
.add_edge("review", "executor")
.add_edge("executor", END)
.compile()?
.with_checkpointer_arc(checkpointer.clone());
// Execute - may pause for approval
let thread_id = "task-001";
let result = graph.invoke(input, ExecutionConfig::new(thread_id)).await;
match result {
Err(GraphError::Interrupted(interrupt)) => {
println!("*** EXECUTION PAUSED ***");
println!("Reason: {}", interrupt.interrupt);
println!("Plan awaiting approval: {:?}", interrupt.state.get("plan"));
// Human reviews and approves...
// Update state with approval
graph.update_state(thread_id, [("approved".to_string(), json!(true))]).await?;
// Resume execution
let final_result = graph.invoke(State::new(), ExecutionConfig::new(thread_id)).await?;
println!("Final result: {:?}", final_result.get("result"));
}
Ok(result) => {
println!("Completed without interrupt: {:?}", result);
}
Err(e) => {
println!("Error: {}", e);
}
}
Statische Interrupts
Verwende interrupt_before oder interrupt_after fΓΌr obligatorische Pausenpunkte:
let graph = StateGraph::with_channels(&["task", "plan", "result"])
.add_node(planner_node)
.add_node(executor_node)
.add_edge(START, "planner")
.add_edge("planner", "executor")
.add_edge("executor", END)
.compile()?
.with_interrupt_before(&["executor"]); // Always pause before execution
Streaming-AusfΓΌhrung
Streame Events, wΓ€hrend der Graph ausgefΓΌhrt wird:
use futures::StreamExt;
use adk_graph::stream::StreamMode;
let stream = agent.stream(input, config, StreamMode::Updates);
while let Some(event) = stream.next().await {
match event? {
StreamEvent::NodeStart(name) => println!("Starting: {}", name),
StreamEvent::Updates { node, updates } => {
println!("{} updated state: {:?}", node, updates);
}
StreamEvent::NodeEnd(name) => println!("Completed: {}", name),
StreamEvent::Done(state) => println!("Final state: {:?}", state),
_ => {}
}
}
Stream-Modi
| Modus | Beschreibung |
|---|---|
Values | VollstΓ€ndigen Status nach jedem Knoten streamen |
Updates | Nur StatusΓ€nderungen streamen |
Messages | Nachrichten-Typ-Aktualisierungen streamen |
Debug | Alle internen Ereignisse streamen |
Messages-Modus liest Tokens aus Node::execute_stream, sobald sie erzeugt werden.
Jeder Knoten lΓ€uft in diesem Modus pro Super-Schritt einmal: Der Knoten meldet seine Zustandsaktualisierungen im Stream als ein StreamEvent::Updates-Ereignis, und der Executor ΓΌbernimmt diese, statt den Knoten ein zweites Mal auszufΓΌhren, um sie zu sammeln. Das ist besonders wichtig fΓΌr AgentNode, wo eine zweite AusfΓΌhrung einen zweiten berechneten Model-Aufruf pro Knoten bedeuten wΓΌrde.
Wichtig: Ein benutzerdefinierter
Node, derexecute_streamΓΌberschreibt, muss einStreamEvent::Updates-Ereignis liefern, das seine Zustandsaktualisierungen trΓ€gt. Ohne dieses streamt der Knoten Ereignisse, trΓ€gt aber imMessages-Modus keinen Zustand bei. Der standardmΓ€Γigeexecute_stream, derexecuteumschlieΓt, erledigt das fΓΌr dich.
Timeout-Richtlinien gelten fΓΌr die gestreamte AusfΓΌhrung selbst. FΓΌr einen Stream bedeutet idle_timeout, dass innerhalb des Limits kein Ereignis erzeugt wurde.
ADK Integration
GraphAgent implementiert das ADK Agent-Trait und funktioniert daher mit:
- Runner: Mit
adk-runnerfΓΌr die StandardausfΓΌhrung verwenden - Callbacks: VollstΓ€ndige UnterstΓΌtzung fΓΌr Callbacks vor/nach
- Sessions: Funktioniert mit
adk-sessionfΓΌr den GesprΓ€chsverlauf - Streaming: Gibt ADK
EventStreamzurΓΌck
use adk_runner::Runner;
let graph_agent = GraphAgent::builder("workflow")
.before_agent_callback(|ctx| async {
println!("Starting graph execution for session: {}", ctx.session_id());
Ok(())
})
.after_agent_callback(|ctx, event| async {
if let Some(content) = event.content() {
println!("Graph completed with content");
}
Ok(())
})
// ... graph definition
.build()?;
// GraphAgent implements Agent trait - use with Launcher or Runner
// See adk-runner README for Runner configuration
Beispiele
Validierte Graph-Beispiele in diesem Repository:
cargo run --manifest-path examples/tier_examples/standard/Cargo.toml --bin 11-standard-graph
cargo run --manifest-path examples/tier_examples/standard/Cargo.toml --bin 12-standard-sequential
cargo run --manifest-path examples/competitive_graph_resume/Cargo.toml
Der in diese Website eingebettete ADK-Rust Playground enthΓ€lt die vollstΓ€ndige Graph-Galerie mit echter LLM-Integration.
Vergleich mit LangGraph
| Funktion | LangGraph | adk-graph |
|---|---|---|
| Zustandsverwaltung | TypedDict + Reducers | StateSchema + Reducers |
| AusfΓΌhrungsmodell | Pregel-Supersteps | Pregel-Supersteps |
| Checkpointing | Memory, SQLite, Postgres | Memory, SQLite |
| Mensch-im-Loop | interrupt_before/after | interrupt_before/after + dynamisch |
| Streaming | 5 Modi | 5 Modi |
| Zyklen | Native UnterstΓΌtzung | Native UnterstΓΌtzung |
| Typsicherheit | Python-Typisierung | Rust-Typsystem |
| LLM Integration | LangChain | AgentNode + ADK Agenten |
Vorherige: β Multi-Agent-Systeme | NΓ€chste: Echtzeit-Agenten β