Funkcionális munkafolyamat API

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 kiexecutor_invoked / executor_completed / executor_failed a megfigyelhetőség érdekében kerül kibocsátásra. Gyorsítótár-találat esetén ehelyett executor_bypassed kerü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 egy ctx: RunContext paramé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:

Következő lépések

Kapcsolódó témakörök:

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.