Megjegyzés
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhat bejelentkezni vagy módosítani a címtárat.
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhatja módosítani a címtárat.
Warning
A funkcionális munkafolyamat API kísérleti jellegű, és előzetes értesítés nélkül módosítható vagy eltávolítható a jövőbeli verziókban.
A funkcionális munkafolyamat API-val egyszerű Python aszinkron függvényként írhat munkafolyamatokat. A végrehajtó osztályok és kábelezési élek meghatározása, illetve a WorkflowBuilder használata helyett díszítsen fel egy async függvényt @workflow-val, és használjon natív Python vezérlési folyamatokat – például if/else, for ciklusok, asyncio.gather – a logika kifejezésére.
A graph API-val való párhuzamos összehasonlításért tekintse meg a Munkafolyamatok API-k áttekintését.
@workflow Dekoratőr
Alkalmazza a(z) @workflow elemet egy async függvényre állapot nélküli FunctionalWorkflowDefinition létrehozásához:
from agent_framework import workflow
@workflow
async def text_pipeline(text: str) -> str:
upper = await to_upper_case(text)
return await reverse_text(upper)
A @workflow dekoratőr támogatja a paraméteres űrlapot opcionális argumentumokkal:
from agent_framework import workflow
@workflow(name="my_pipeline", description="Uppercase then reverse")
async def text_pipeline(text: str) -> str:
...
@workflow Paraméterek
| Paraméter | Típus | Leírás |
|---|---|---|
name |
str | None |
A munkafolyamat kijelzőneve. A függvény alapértelmezett értéke __name__. |
description |
str | None |
Nem kötelező, emberi olvasásra alkalmas leírás. |
Munkafolyamat létrehozása
A definíción hívja meg a .build() elemet egy állapottal rendelkező FunctionalWorkflow létrehozásához. Minden beépített munkafolyamat egy logikai hívót vagy munkamenetet jelöl. Hozzon létre külön példányokat az egymástól független hívók számára, hogy elkülönítse futtatási és újrajátszási állapotukat.
Adja át az ellenőrzőpont-tárolást a(z) .build() számára, amikor a munkafolyamatnak tartósan tárolnia kell a lépések eredményeit és az állapotot:
from agent_framework import InMemoryCheckpointStorage
storage = InMemoryCheckpointStorage()
workflow_instance = text_pipeline.build(checkpoint_storage=storage)
| Paraméter | Típus | Leírás |
|---|---|---|
checkpoint_storage |
CheckpointStorage | None |
A lépéseredmények és az állapot futások közötti megőrzésének alapértelmezett tárolója. |
A munkafolyamat függvény szignatúrája
A munkafolyamat-függvény első paramétere megkapja a beépített munkafolyamathoz .run() átadott bemenetet. Csak akkor adjon hozzá paramétert ctx: RunContext , ha HITL- vagy kulcs-/értékállapotra vagy egyéni eseményekre van szüksége – máskülönben nem kötelező:
# No ctx needed — just a plain pipeline
@workflow
async def simple_pipeline(data: str) -> str:
result = await process(data)
return result
# ctx needed for HITL, state, or custom events
@workflow
async def hitl_pipeline(data: str, ctx: RunContext) -> str:
feedback = await ctx.request_info({"draft": data}, response_type=str)
return feedback
RunContext először típusannotációval észlelve, majd a paraméternév ctx alapján, így mind a ctx: RunContext, mind a csupasz ctx paraméter működik.
Munkafolyamat futtatása
Hívás .run() az FunctionalWorkflow objektumon, amelyet a .build() visszaadott:
workflow_instance = text_pipeline.build()
# .run() wraps the result in a WorkflowRunResult with events and state
result = await workflow_instance.run("hello world")
print(result.text) # first output as a string
print(result.get_outputs()) # list of terminal outputs
print(result.get_intermediate_outputs()) # list of intermediate outputs
print(result.get_final_state()) # WorkflowRunState.IDLE
run() Paraméterek
| Paraméter | Típus | Leírás |
|---|---|---|
message |
Any | None |
A munkafolyamat-függvényhez első argumentumként átadott bemenet. |
stream |
bool |
Ha True, visszaad egy ResponseStream-et, ami WorkflowEvent objektumokat eredményez. Alapértelmezett érték: False. |
responses |
dict[str, Any] | None |
HITL-válaszok request_id-hoz társítva. Felfüggesztett munkafolyamat folytatására szolgál. |
checkpoint_id |
str | None |
Ellenőrzőpont, ahonnan vissza szeretne állítani. A checkpoint_storage értékét a build során vagy ehhez a futtatáshoz be kell állítani. |
checkpoint_storage |
CheckpointStorage | None |
Felülbírálja a futtatáshoz szükséges buildelési idő ellenőrzőpont-tárolót. |
include_status_events |
bool |
Az állapotváltozási események belefoglalása a nem streamelt eredménybe. |
Adjon meg egy bemeneti módot hívásonként: message, responsesvagy checkpoint_id. A kivétel a külső bemenet használatával történő ellenőrzőpontból való folytatás, ahol a checkpoint_id és a responses együtt átadható.
WorkflowRunResult
run() (nem streamelés) egy WorkflowRunResult. Főbb módszerek:
| Metódus/tulajdonság | Returns | Leírás |
|---|---|---|
.text |
str |
Első kimenet karakterláncként. Üres sztring, ha nincs sztringkimenet. |
.get_outputs() |
list[Any] |
A munkafolyamat által generált összes terminálkimenet (type == "output" elemmel rendelkező események). |
.get_intermediate_outputs() |
list[Any] |
A munkafolyamat által létrehozott összes köztes kimenet (a type == "intermediate" elemet tartalmazó események). |
.get_final_state() |
WorkflowRunState |
Utolsó futtatási állapot (IDLE, IDLE_WITH_PENDING_REQUESTS, FAILED, ...). |
.get_request_info_events() |
list[WorkflowEvent] |
Függőben lévő HITL-kérelmek, amikor az állapot IDLE_WITH_PENDING_REQUESTS. |
Streaming
Adjon meg stream=True az események megkapásához, amint előállnak:
from agent_framework import workflow
@workflow
async def data_pipeline(url: str) -> str:
raw = await fetch_data(url)
return await transform_data(raw)
workflow_instance = data_pipeline.build()
# stream=True returns a ResponseStream you iterate with async for
stream = workflow_instance.run("https://example.com/api/data", stream=True)
async for event in stream:
if event.type == "output":
print(f"Output: {event.data}")
# After iteration, get_final_response() returns the WorkflowRunResult
result = await stream.get_final_response()
print(f"Final state: {result.get_final_state()}")
Tekintse meg python/samples/03-workflows/functional/basic_streaming_pipeline.py a teljes példát.
@step Dekoratőr
@step egy választható dekorátor, amely eredmény gyorsítótárazást, eseménykibocsátást és lépésenkénti ellenőrzőpontozást ad hozzá az egyes aszinkron függvényekhez.
from agent_framework import step, workflow
@step
async def fetch_data(url: str) -> dict:
# expensive — hits a real API
return await http_get(url)
@workflow
async def pipeline(url: str) -> str:
raw = await fetch_data(url)
return process(raw)
Mit @step végez a munkafolyamaton belül?
-
Gyorsítótárazza az eredményeket – az eredményt a rendszer tárolja
(step_name, call_index). A HITL-folytatás vagy az ellenőrzőpont visszaállításakor a befejezett lépés az ismételt végrehajtás helyett azonnal visszaadja a mentett eredményt. -
Eseményeket bocsát ki –
executor_invoked/executor_completed/executor_faileda megfigyelhetőség érdekében kerül kibocsátásra. Gyorsítótár-találat esetén ehelyettexecutor_bypassedkerül kibocsátásra. -
Ellenőrzőpontokat ment – ha a beépített munkafolyamat rendelkezik
checkpoint_storage–, a rendszer minden lépés befejezése után ment egy ellenőrzőpontot. -
Injektál
RunContext– ha a lépésfüggvény deklarál egyctx: RunContextparamétert, az aktív környezet automatikusan injektálásra kerül.
A futó munkafolyamaton @step kívül transzparens – a függvény ugyanúgy viselkedik, mint a leválasztatlan verziója, így teljesen tesztelhető elszigetelten.
Mikor érdemes használni a @step
Függvényeken, amelyeknek az @step, például az ügynökhívásokon, külső API-kéréseken, vagy bármely olyan műveleten használható, ahol az újrafuttatás költséges lenne vagy mellékhatásokat okozhatna. Az egyszerű függvények (anélkül @step) továbbra is működnek belül @workflow; egyszerűen újrafuttatják őket, amikor a munkafolyamat folytatódik.
from agent_framework import InMemoryCheckpointStorage, step, workflow
@step # cached — won't re-run on resume
async def call_llm(prompt: str) -> str:
return (await agent.run(prompt)).text
# No @step — cheap, fine to re-run
async def validate(text: str) -> bool:
return len(text) > 0
@workflow
async def pipeline(topic: str) -> str:
draft = await call_llm(f"Write about: {topic}")
ok = await validate(draft)
return draft if ok else ""
storage = InMemoryCheckpointStorage()
workflow_instance = pipeline.build(checkpoint_storage=storage)
@step egy paramétert name is elfogad:
@step(name="transform")
async def transform_data(raw: dict) -> str:
...
Tekintse meg python/samples/03-workflows/functional/steps_and_checkpointing.py a teljes példát.
RunContext
RunContext a munkafolyamatba és a lépésfüggvényekbe injektált végrehajtási környezet. Csak HITL, kulcs/érték állapot vagy egyéni események használatakor van rá szüksége.
Importálja a következőből agent_framework:
from agent_framework import RunContext, workflow
ctx.request_info() — Ember a rendszerben
ctx.request_info() felfüggeszti a munkafolyamatot, hogy megvárja a külső bemenetet:
@workflow
async def review_pipeline(topic: str, ctx: RunContext) -> str:
draft = await write_draft(topic)
feedback = await ctx.request_info(
{"draft": draft, "instructions": "Please review this draft"},
response_type=str,
request_id="review_request",
)
return await revise_draft(draft, feedback)
Paraméterek:
| Paraméter | Típus | Leírás |
|---|---|---|
request_data |
Any |
Adatcsomag, amely leírja, hogy milyen bemenetre van szükség (szótár, Pydantic-modell, karakterlánc, ...). |
response_type |
type |
A válasz várható Python típusa. |
request_id |
str | None |
A kérés stabil azonosítója. Ha nincs megadva, a rendszer egy determinisztikus auto::<index> azonosítót hoz létre a hívásrendelésből. |
Visszajátszás szemantikája: Az első végrehajtáskor egy belső jelzést ad, request_info() amely felfüggeszti a munkafolyamatot. A jel soha nem látható a kódban. A hívó megkap egy WorkflowRunResult-t get_final_state() == WorkflowRunState.IDLE_WITH_PENDING_REQUESTS-vel.
A folytatáshoz hívja meg a(z) .run(responses={request_id: value}) elemet ugyanazon a létrehozott munkafolyamaton. A munkafolyamat felülről újrafut, és request_info() azonnal visszaadja a megadott értéket.
@step-dekorált függvények, amelyek a felfüggesztés előtt futottak, újraindításkor gyorsítótárazott eredményeiket adják vissza újrafuttatás helyett.
A válasz kezelése:
workflow_instance = review_pipeline.build()
# Phase 1 — run until the workflow pauses
result1 = await workflow_instance.run("AI Safety")
assert result1.get_final_state() == WorkflowRunState.IDLE_WITH_PENDING_REQUESTS
requests = result1.get_request_info_events()
print(requests[0].request_id) # "review_request"
# Phase 2 — resume with the human's answer
result2 = await workflow_instance.run(
responses={"review_request": "Add more details about alignment research"}
)
print(result2.text)
Tekintse meg python/samples/03-workflows/functional/hitl_review.py a teljes példát.
ctx.request_info() függvényen belül @step is támogatott.
ctx.add_event() — Egyéni események
Alkalmazásspecifikus ctx.add_event() eseményeket bocsáthat ki a keretrendszer életciklus-eseményei mellett. További részletekért és példákért lásd: Egyéni események kibocsátása.
ctx.get_state()
/
ctx.set_state() — Kulcs/érték állapota
Használja a ctx.get_state() és ctx.set_state() elemeket az értékek tárolására, amelyek megmaradnak a HITL-megszakítások során, és szerepelnek az ellenőrzőpontokban. További részletekért lásd: Munkafolyamat állapota.
Az állapotértékek JSON-szerializálhatóaknak kell lenniük az Ellenőrzőpont-tároló konfigurálásakor.
ctx.is_streaming()
Amikor az aktuális futtatást True-vel indították, stream=True-t ad vissza. Hasznos belső lépésfüggvények, amelyek a streamelési mód alapján szeretnék módosítani a működésüket.
get_run_context()
Egy futó munkafolyamat bármely pontjáról lekéri az aktívt RunContext – olyan segédfüggvényekben hasznos, amelyek nem deklarálnak paramétert ctx :
from agent_framework import get_run_context
async def helper():
ctx = get_run_context()
if ctx is not None:
ctx.set_state("helper_ran", True)
Akkor ad None vissza, ha egy futó munkafolyamaton kívül hívják meg.
Párhuzamosság asyncio.gather -tel
Használja a standard Python párhuzamosságot a fan-out/fan-in feladatokhoz – nincs szükség keretrendszerbeli primitívekre.
import asyncio
from agent_framework import workflow
@workflow
async def research_pipeline(topic: str) -> str:
web, papers, news = await asyncio.gather(
research_web(topic),
research_papers(topic),
research_news(topic),
)
return await synthesize([web, papers, news])
asyncio.gather akkor is működik, ha a függvényeket dekorátorral látják el @step.
Tekintse meg python/samples/03-workflows/functional/parallel_pipeline.py a teljes példát.
Ügynökök meghívása munkafolyamatokon belül
Az ügynökhívások egyszerű függvényhívásként működnek belül @workflow:
from agent_framework import Agent, workflow
writer = Agent(name="WriterAgent", instructions="Write a short poem.", client=client)
reviewer = Agent(name="ReviewerAgent", instructions="Review the poem.", client=client)
@workflow
async def poem_workflow(topic: str) -> str:
poem = (await writer.run(f"Write a poem about: {topic}")).text
review = (await reviewer.run(f"Review this poem: {poem}")).text
return f"Poem:\n{poem}\n\nReview: {review}"
Adja hozzá @step az ügynököt hívó függvényekhez, ha azt szeretné, hogy az eredmények gyorsítótárazva legyenek a HITL-folytatások vagy az ellenőrzőpont-visszaállítások során.
from agent_framework import step
@step
async def write_poem(topic: str) -> str:
return (await writer.run(f"Write a poem about: {topic}")).text
Tekintse meg python/samples/03-workflows/functional/agent_integration.py a teljes példát.
.as_agent() – Munkafolyamat használata ügynökként
Hozzon létre egy FunctionalWorkflow, majd csomagolja be ügynökkompatibilis objektumként a következővel .as_agent():
from agent_framework import workflow
@workflow
async def poem_workflow(topic: str) -> str:
...
# Build the workflow and wrap it as an agent
agent = poem_workflow.build().as_agent(name="PoemAgent")
# Use with the standard agent interface
response = await agent.run("Write a poem about the ocean")
print(response.text)
# Or use in a larger workflow or orchestration
.as_agent() Olyan FunctionalWorkflowAgent felületet ad vissza, amely ugyanazt run() az interfészt teszi elérhetővé, mint más ügynökobjektumok, így a funkcionális munkafolyamatok bármely olyan rendszerrel összeállíthatók, amelyek ügynököket fogadnak el.
| Paraméter | Típus | Leírás |
|---|---|---|
name |
str | None |
Az ügynök megjelenítendő neve. A munkafolyamat nevének alapértelmezett értéke. |
description |
str | None |
Nem kötelező leírás felülbírálása. A munkafolyamat leírásának alapértelmezett értéke. |
context_providers |
Sequence[Any] | None |
Az ügynökhöz társítható opcionális kontextusszolgáltatók. |
Samples
Futtatható példák a következő mintamappákban találhatók:
-
python/samples/01-get-started/– bevezető@workflowpéldák -
python/samples/03-workflows/functional/– teljes funkcionalitású funkcionális munkafolyamat-minták
Következő lépések
Kapcsolódó témakörök:
- Végrehajtók – feldolgozási egységek a gráfalapú API-ban
- Human-in-the-loop – HITL gráfalapú munkafolyamatokban
- Ellenőrzőpontok – ellenőrzőpontok tárolása és folytatása
- Események – munkafolyamat-eseménytípusok
- Munkafolyamatok használata ügynökökként
A funkcionális munkafolyamat API jelenleg nem érhető el a C# számára.
A funkcionális munkafolyamat-API jelenleg nem érhető el Go nyelven. Használja a gráfalapú munkafolyamat API-kat a Workflow Builder > Execution alkalmazásban.