Hinweis
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, sich anzumelden oder das Verzeichnis zu wechseln.
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, das Verzeichnis zu wechseln.
Diese Seite bietet eine Übersicht über Prüfpunkte im Microsoft Agent Framework-Workflowsystem.
Overview
Prüfpunkte ermöglichen es Ihnen, den Status eines Workflows an bestimmten Punkten während der Ausführung zu speichern und von diesen Punkten später fortzusetzen. Dieses Feature eignet sich besonders für die folgenden Szenarien:
- Lang andauernde Workflows, bei denen Sie den Fortschrittsverlust im Falle von Fehlern vermeiden möchten.
- Lang andauernde Workflows, bei denen Sie die Ausführung zu einem späteren Zeitpunkt anhalten und fortsetzen möchten.
- Workflows, die regelmäßige Zustandsspeicherung für Überwachungs- oder Compliancezwecke erfordern.
- Workflows, die in verschiedenen Umgebungen oder Instanzen migriert werden müssen.
Wann werden Prüfpunkte erstellt?
Denken Sie daran, dass Workflows in Supersteps ausgeführt werden, wie im Workflowausführungsmodell dokumentiert. Prüfpunkte werden am Ende jedes Supersteps erstellt, nachdem alle Ausführenden in diesem Superstep ihre Ausführung abgeschlossen haben. Ein Prüfpunkt erfasst den gesamten Status des Workflows, einschließlich:
- Der aktuelle Status aller Executoren
- Alle ausstehenden Nachrichten im Workflow für den nächsten Superstep
- Ausstehende Anforderungen und Antworten
- Freigegebene Zustände
Note
Ab Python Version 1.13.0 erstellen Workflows auch einen Eintragsprüfpunkt vor dem ersten Superstep zum Aufzeichnen der Workfloweingabe und einen weiteren Eintragsprüfpunkt, wenn Antworten auf Anforderungsereignisse übermittelt werden. Diese Prüfpunkte machen die Ausführung des vollständigen Workflows erneut abspielbar. Diese Version enthält geringfügige Änderungen für Anwendungen, die von Iterationsanzahlen, Nachrichtenquell-IDs oder Prüfpunkt-Sortierung abhängen. Vorhandene Prüfpunkte werden weiterhin unterstützt. Details zur Migration finden Sie unter Upgrade Python Workflowprüfpunkte auf 1.13.0.
Erfassen von Prüfpunkten
Zum Aktivieren der Prüfpunkte muss beim Ausführen des Workflows eine CheckpointManager Bereitstellung erfolgen. Auf einen Prüfpunkt kann dann über ein SuperStepCompletedEvent oder über die Checkpoints Eigenschaft während des Laufs zugegriffen werden.
using Microsoft.Agents.AI.Workflows;
// Create a checkpoint manager to manage checkpoints
CheckpointManager checkpointManager = CheckpointManager.CreateInMemory();
// Run the workflow with checkpointing enabled
StreamingRun run = await InProcessExecution
.RunStreamingAsync(workflow, input, checkpointManager)
.ConfigureAwait(false);
await foreach (WorkflowEvent evt in run.WatchStreamAsync().ConfigureAwait(false))
{
if (evt is SuperStepCompletedEvent superStepCompletedEvt)
{
// Access the checkpoint
CheckpointInfo? checkpoint = superStepCompletedEvt.CompletionInfo?.Checkpoint;
}
}
// Checkpoints can also be accessed from the run directly
IReadOnlyList<CheckpointInfo> checkpoints = run.Checkpoints;
Um das Checkpointing zu aktivieren, muss beim Erstellen eines Workflows ein CheckpointStorage bereitgestellt werden. Über den Speicher kann dann auf einen Prüfpunkt zugegriffen werden. Agent Framework enthält drei integrierte Implementierungen – wählen Sie die Implementierung aus, die Ihren Anforderungen an Haltbarkeit und Bereitstellung entspricht:
| Provider | Paket | Durability | Am besten geeignet für: |
|---|---|---|---|
InMemoryCheckpointStorage |
agent-framework |
Nur in Bearbeitung | Tests, Demos, kurzlebige Workflows |
FileCheckpointStorage |
agent-framework |
Lokaler Datenträger | Workflows mit einem Computer, lokale Entwicklung |
CosmosCheckpointStorage |
agent-framework-azure-cosmos |
Azure Cosmos DB (ein Microsoft-Datenbankdienst) | Produktion, verteilte, prozessübergreifende Workflows |
Alle drei implementieren dasselbe CheckpointStorage Protokoll, sodass Sie Anbieter austauschen können, ohne Workflow- oder Executorcode zu ändern.
InMemoryCheckpointStorage hält Prüfpunkte im Prozessspeicher. Am besten geeignet für Tests, Demos und kurzlebige Workflows, bei denen Sie keine Haltbarkeit über Neustarts hinweg benötigen.
from agent_framework import (
InMemoryCheckpointStorage,
WorkflowBuilder,
)
# Create a checkpoint storage to manage checkpoints
checkpoint_storage = InMemoryCheckpointStorage()
# Build a workflow with checkpointing enabled
builder = WorkflowBuilder(start_executor=start_executor, checkpoint_storage=checkpoint_storage)
builder.add_edge(start_executor, executor_b)
builder.add_edge(executor_b, executor_c)
builder.add_edge(executor_b, end_executor)
workflow = builder.build()
# Run the workflow
async for event in workflow.run(input, stream=True):
...
# Access checkpoints from the storage
checkpoints = await checkpoint_storage.list_checkpoints(workflow_name=workflow.name)
Um Checkpointing zu aktivieren, konfigurieren Sie die Ausführungsumgebung mit einem Checkpoint-Manager. Ein Checkpoint kann dann über workflow.SuperStepCompletedEvent oder über die Checkpoint-Liste des Durchlaufs aufgerufen werden.
checkpointManager := checkpoint.NewInMemoryManager()
run, err := inproc.Default.
WithCheckpointing(checkpointManager).
RunStreaming(ctx, wf, input)
if err != nil {
return err
}
defer run.Close(ctx)
var checkpoints []workflow.CheckpointInfo
for evt, err := range run.WatchStream(ctx) {
if err != nil {
return err
}
if completed, ok := evt.(workflow.SuperStepCompletedEvent); ok && completed.CompletionInfo != nil {
if completed.CompletionInfo.CheckpointInfo != nil {
checkpoints = append(checkpoints, *completed.CompletionInfo.CheckpointInfo)
}
}
}
// Checkpoints can also be accessed from the run directly.
checkpoints = run.Checkpoints()
Fortsetzen ab Prüfpunkten
Sie können einen Workflow von einem bestimmten Checkpoint direkt im selben Ausführungslauf fortsetzen.
// Assume we want to resume from the 6th checkpoint
CheckpointInfo savedCheckpoint = run.Checkpoints[5];
// Restore the state directly on the same run instance.
await run.RestoreCheckpointAsync(savedCheckpoint).ConfigureAwait(false);
await foreach (WorkflowEvent evt in run.WatchStreamAsync().ConfigureAwait(false))
{
if (evt is WorkflowOutputEvent workflowOutputEvt)
{
Console.WriteLine($"Workflow completed with result: {workflowOutputEvt.Data}");
}
}
Sie können einen Workflow aus einem bestimmten Prüfpunkt direkt in derselben Workflowinstanz fortsetzen.
# Assume we want to resume from the 6th checkpoint
saved_checkpoint = checkpoints[5]
async for event in workflow.run(checkpoint_id=saved_checkpoint.checkpoint_id, stream=True):
...
Sie können einen Streaming-Lauf direkt im selben Lauf auf einen bestimmten Checkpoint zurücksetzen.
// Assume we want to resume from the 6th checkpoint.
savedCheckpoint := checkpoints[5]
if err := run.RestoreCheckpoint(ctx, savedCheckpoint); err != nil {
return err
}
for evt, err := range run.WatchStream(ctx) {
if err != nil {
return err
}
if outputEvent, ok := evt.(workflow.OutputEvent); ok {
fmt.Printf("Workflow completed with result: %v\n", outputEvent.Output)
}
}
Rehydratieren von Prüfpunkten
Ein rehydrierter Workflow muss die Topologie und die Executor-Identitäten des Workflows beibehalten, der den Checkpoint erstellt hat. Wie die Identität des Executors aufgelöst wird, hängt vom SDK und vom Typ des Executors ab.
Sie können einen Workflow auch aus einem Prüfpunkt in eine neue Ausführungsinstanz rehydratisieren.
// A rehydrated workflow must preserve the topology and executor identities of the workflow that
// created the checkpoint. This executor-only workflow rebuilds identically because its executors
// use fixed ids. Agent-based workflows must recreate each local agent with the same
// ChatClientAgentOptions.Id (and, if set, the same Name), otherwise the executor ids no longer
// match the checkpoint and resume fails.
var newWorkflow = WorkflowFactory.BuildWorkflow();
const int CheckpointIndex = 5;
Console.WriteLine($"\n\nHydrating a new workflow instance from the {CheckpointIndex + 1}th checkpoint.");
CheckpointInfo savedCheckpoint = checkpoints[CheckpointIndex];
await using StreamingRun newCheckpointedRun =
await InProcessExecution.ResumeStreamingAsync(newWorkflow, savedCheckpoint, checkpointManager);
Important
Der Workflow, der an ResumeStreamingAsync übergeben wird, muss dieselbe Struktur und dieselben Executor-Identitäten aufweisen wie der Workflow, der den Checkpoint erstellt hat. Wenn der Workflow lokale ChatClientAgent Instanzen enthält, die über Anforderungen, Abhängigkeitseinfügungsbereiche, Prozesse oder Bereitstellungen hinweg rekonstruiert werden, weisen Sie jedem Agent einen stabilen ChatClientAgentOptions.IdZu. Wenn ein Agent auch einen Name festlegt, lassen Sie diesen Name ebenfalls unverändert.
Weisen Sie beispielsweise eine ID zu, die die logische Rolle des Agents darstellt:
// Give each agent a stable, unique Id so its workflow executor identity stays the same when the
// workflow is reconstructed (for example per request or dependency-injection scope), which keeps
// checkpoints resumable. If an agent also has a Name, keep that stable too, since the executor
// identity includes it. Use a fixed logical role here, not a conversation, request, or user id.
internal const string IntakeAgentName = "Assistant";
public AIAgent IntakeAgent { get; } = chatClient.AsAIAgent(new ChatClientAgentOptions
{
Id = "intake-agent",
Name = IntakeAgentName,
ChatOptions = new()
{
Instructions =
"""
You receive a user request and are responsible for routing to the correct initial expert agent.
""",
},
});
Wenden Sie dieses Muster auf jeden Agent an, der am Workflow teilnimmt. Agent-IDs müssen innerhalb des Workflows eindeutig sein und beim Rekonstruieren desselben logischen Agents wiederverwendet werden. Verwenden Sie keine Unterhaltungs-IDs, Anforderungs-IDs, Benutzer-IDs, persönlich identifizierbare Informationen oder geheime Schlüssel als Agent-IDs.
Wenn ein Agent Name festgelegt ist, wird die aktuelle .NET-Workflow-Executor-Identität sowohl aus Name als auch aus Id abgeleitet, sodass die Änderung eines der beiden Werte den neu erstellten Workflow mit dem Checkpoint inkompatibel macht. Das Zuweisen stabiler Werte repariert keine Prüfpunkte, die mit unterschiedlichen oder zufällig generierten IDs erstellt wurden; starten Sie stattdessen eine neue Sitzungs- und Prüfpunktlinie.
Verwandte Szenarien finden Sie unter Workflows als Agents und Handoff-Orchestrierung.
Alternativ können Sie eine neue Workflowinstanz aus einem Prüfpunkt rehydratisieren.
from agent_framework import WorkflowBuilder
builder = WorkflowBuilder(start_executor=start_executor)
builder.add_edge(start_executor, executor_b)
builder.add_edge(executor_b, executor_c)
builder.add_edge(executor_b, end_executor)
# This workflow instance doesn't require checkpointing enabled.
workflow = builder.build()
# Assume we want to resume from the 6th checkpoint
saved_checkpoint = checkpoints[5]
async for event in workflow.run(
checkpoint_id=saved_checkpoint.checkpoint_id,
checkpoint_storage=checkpoint_storage,
stream=True,
):
...
Alternativ können Sie eine neue Workflowinstanz aus einem Prüfpunkt rehydratisieren.
// Assume we want to resume from the 6th checkpoint
savedCheckpoint := checkpoints[5]
newWorkflow := buildWorkflow()
newRun, err := inproc.Default.
WithCheckpointing(checkpointManager).
ResumeStreaming(ctx, newWorkflow, savedCheckpoint)
if err != nil {
return err
}
defer newRun.Close(ctx)
for evt, err := range newRun.WatchStream(ctx) {
if err != nil {
return err
}
if outputEvent, ok := evt.(workflow.OutputEvent); ok {
fmt.Printf("Workflow completed with result: %v\n", outputEvent.Output)
}
}
Vollstreckungsstatus speichern
Um sicherzustellen, dass der Status eines Executors in einem Prüfpunkt erfasst wird, muss der Executor die OnCheckpointingAsync Methode überschreiben und den Status im Workflowkontext speichern.
using Microsoft.Agents.AI.Workflows;
internal sealed partial class CustomExecutor() : Executor("CustomExecutor")
{
private const string StateKey = "CustomExecutorState";
private List<string> messages = new();
[MessageHandler]
private async ValueTask HandleAsync(string message, IWorkflowContext context)
{
this.messages.Add(message);
// Executor logic...
}
protected override ValueTask OnCheckpointingAsync(IWorkflowContext context, CancellationToken cancellation = default)
{
return context.QueueStateUpdateAsync(StateKey, this.messages);
}
}
Um sicherzustellen, dass der Zustand beim Fortsetzen eines Checkpoints ordnungsgemäß wiederhergestellt wird, muss der Ausführende die OnCheckpointRestoredAsync-Methode überschreiben und seinen Zustand aus dem Workflow-Kontext laden.
protected override async ValueTask OnCheckpointRestoredAsync(IWorkflowContext context, CancellationToken cancellation = default)
{
this.messages = await context.ReadStateAsync<List<string>>(StateKey).ConfigureAwait(false);
}
Um sicherzustellen, dass der Status eines Executors in einem Prüfpunkt erfasst wird, muss der Executor die Methode on_checkpoint_save überschreiben und seinen Status als Dictionary zurückgeben.
class CustomExecutor(Executor):
def __init__(self, id: str) -> None:
super().__init__(id=id)
self._messages: list[str] = []
@handler
async def handle(self, message: str, ctx: WorkflowContext):
self._messages.append(message)
# Executor logic...
async def on_checkpoint_save(self) -> dict[str, Any]:
return {"messages": self._messages}
Um sicherzustellen, dass der Zustand beim Fortsetzen eines Prüfpunkts ordnungsgemäß wiederhergestellt wird, muss der Executor die on_checkpoint_restore Methode überschreiben und den Zustand aus dem bereitgestellten Statusverzeichnis wiederherstellen.
async def on_checkpoint_restore(self, state: dict[str, Any]) -> None:
self._messages = state.get("messages", [])
Um sicherzustellen, dass der Status des Executors in einem Checkpoint erfasst wird, fügen Sie dem Executor Checkpoint-Hooks hinzu und speichern Sie den Status über den Workflow-Kontext.
type customExecutor struct {
messages []string
}
func (e *customExecutor) Handle(message string) {
e.messages = append(e.messages, message)
}
func (e *customExecutor) OnCheckpoint(ctx *workflow.Context) error {
return ctx.QueueStateUpdate("CustomExecutorState", "", slices.Clone(e.messages))
}
Wiederherstellen des Zustands in OnCheckpointRestoredFunc:
func (e *customExecutor) OnCheckpointRestored(ctx *workflow.Context) error {
value, err := ctx.ReadState("CustomExecutorState", "")
if err != nil {
return err
}
if value == nil {
e.messages = nil
return nil
}
messages, ok := value.([]string)
if !ok {
return fmt.Errorf("unexpected custom executor state type %T", value)
}
e.messages = slices.Clone(messages)
return nil
}
executorState := &customExecutor{}
custom := workflow.NewExecutor("CustomExecutor", executorState).Extend(&workflow.Executor{
OnCheckpointFunc: executorState.OnCheckpoint,
OnCheckpointRestoredFunc: executorState.OnCheckpointRestored,
}).Bind()
Sicherheitsüberlegungen
Important
Prüfpunktspeicher ist eine Vertrauensgrenze. Unabhängig davon, ob Sie die integrierten Speicherimplementierungen oder eine benutzerdefinierte Implementierung verwenden, muss das Speicher-Back-End als vertrauenswürdige private Infrastruktur behandelt werden. Laden Sie niemals Prüfpunkte aus nicht vertrauenswürdigen oder potenziell manipulierten Quellen.
Stellen Sie sicher, dass der für Prüfpunkte verwendete Speicherort entsprechend gesichert ist. Nur autorisierte Dienste und Benutzer sollten Lese- oder Schreibzugriff auf Prüfpunktdaten haben.
Pickle-Serialisierung
Sowohl FileCheckpointStorage als auch CosmosCheckpointStorage verwenden das modul pickle Python, um nicht-JSON-nativen Zustand wie Datenklassen, Datetimes und benutzerdefinierte Objekte zu serialisieren. Um die Risiken einer willkürlichen Codeausführung während der Deserialisierung zu minimieren, verwenden beide Anbieter standardmäßig einen eingeschränkten Unpickler. Während der Deserialisierung sind nur ein integrierter Satz sicherer Python-Typen (Grundtypen, datetime, uuid, Decimal, gängige Kollektionen usw.) sowie unterstützte Typen des Agent Frameworks oder OpenAI SDK zulässig. Die Allowlist mit Modulpräfix ist nur auf Typen beschränkt: Hilfsfunktionen und andere globale Elemente, die keine Typen sind, werden abgelehnt. Jeder nicht unterstützte Typ führt dazu, dass die Deserialisierung mit einem WorkflowCheckpointException fehlschlägt.
Um zusätzliche anwendungsspezifische Typen zuzulassen, übergeben Sie diese über den allowed_checkpoint_types-Parameter mithilfe des "module:qualname"-Formats.
from agent_framework import FileCheckpointStorage
storage = FileCheckpointStorage(
"/tmp/checkpoints",
allowed_checkpoint_types=[
"my_app.models:SafeState",
"my_app.models:UserProfile",
],
)
Jeder allowed_checkpoint_types Eintrag muss in einen Typ aufgelöst werden. Das Hinzufügen einer Funktion auf Modulebene oder einer anderen globalen Nichttypfunktion macht diese globale Deserialisierbarkeit nicht möglich.
CosmosCheckpointStorage akzeptiert denselben Parameter:
from azure.identity.aio import DefaultAzureCredential
from agent_framework_azure_cosmos import CosmosCheckpointStorage
storage = CosmosCheckpointStorage(
endpoint="https://my-account.documents.azure.com:443/",
credential=DefaultAzureCredential(),
database_name="agent-db",
container_name="checkpoints",
allowed_checkpoint_types=[
"my_app.models:SafeState",
"my_app.models:UserProfile",
],
)
Wenn Ihr Bedrohungsmodell überhaupt keine pickle-basierte Serialisierung zulässt, verwenden Sie InMemoryCheckpointStorage oder implementieren Sie eine benutzerdefinierte CheckpointStorage mit einer alternativen Serialisierungsstrategie.
Verantwortung für den Speicherort
FileCheckpointStorage erfordert einen expliziten storage_path Parameter – es gibt kein Standardverzeichnis. Während das Framework gegen Pfad-Traversalangriffe validiert, liegt es in der Verantwortung des Entwicklers, das Speicherverzeichnis selbst (Dateiberechtigungen, Verschlüsselung im Ruhezustand, Zugriffssteuerungen) zu sichern. Nur autorisierte Prozesse sollten Lese- oder Schreibzugriff auf das Prüfpunktverzeichnis haben.
CosmosCheckpointStorage basiert auf Azure Cosmos DB für den Speicher. Verwenden Sie, wo möglich, verwaltete Identität/RBAC, begrenzen Sie den Geltungsbereich der Datenbank und des Containers auf den Workflow-Dienst und wechseln Sie die Kontoschlüssel, wenn Sie die schlüsselbasierte Authentifizierung verwenden. Wie bei der Dateispeicherung sollten nur autorisierte Entitäten Lese- oder Schreibzugriff auf den Cosmos DB-Container haben, der Checkpoint-Dokumente enthält.
Checkpoint-Manager in Go serialisieren den Checkpoint-Zustand als JSON, aber der Checkpoint-Speicher gilt weiterhin als vertrauenswürdiger Anwendungszustand. Wenn Sie checkpoint.NewFileSystemJSONStore verwenden, speichern Sie Prüfpunktdateien in einem geschützten Verzeichnis und beschränken Sie den Lese-/Schreibzugriff nur auf autorisierte Prozesse. Kundenspezifische Stores sind selbst für die Garantien hinsichtlich Zugriffssteuerung, Integrität und Dauerhaftigkeit verantwortlich.