تنسيقات مهام سير عمل Microsoft Agent Framework - متزامنة

يمكن التزامن المتزامن عوامل متعددة من العمل على نفس المهمة بالتوازي. يعالج كل عامل الإدخال بشكل مستقل، ويتم جمع نتائجه وتجميعها. هذا النهج مناسب تماما للسيناريوهات التي تكون فيها وجهات النظر أو الحلول المتنوعة ذات قيمة، مثل وضع الأفكار أو التفكير أو أنظمة التصويت.

التزامن المتزامن

ما ستتعلمه

  • كيفية تحديد عوامل متعددة ذات خبرة مختلفة
  • كيفية تنسيق هؤلاء العوامل للعمل بشكل متزامن على مهمة واحدة
  • كيفية جمع النتائج ومعالجتها

في التزامن المتزامن، يعمل العديد من العوامل على نفس المهمة بشكل متزامن ومستقل، ما يوفر وجهات نظر متنوعة حول نفس الإدخال.

إعداد عميل Azure OpenAI

using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading.Tasks;
using Azure.AI.Projects;
using Azure.Identity;
using Microsoft.Agents.AI.Workflows;
using Microsoft.Extensions.AI;
using Microsoft.Agents.AI;

// 1) Set up the Azure OpenAI client
var endpoint = Environment.GetEnvironmentVariable("AZURE_OPENAI_ENDPOINT") ??
    throw new InvalidOperationException("AZURE_OPENAI_ENDPOINT is not set.");
var deploymentName = Environment.GetEnvironmentVariable("AZURE_OPENAI_DEPLOYMENT_NAME") ?? "gpt-4o-mini";
var client = new AIProjectClient(new Uri(endpoint), new DefaultAzureCredential())
    .GetProjectOpenAIClient()
    .GetProjectResponsesClient()
    .AsIChatClient(deploymentName);

تحذير

DefaultAzureCredential مناسب للتنمية ولكنه يتطلب دراسة متأنية في الإنتاج. في الإنتاج، ضع في اعتبارك استخدام بيانات اعتماد محددة (على سبيل المثال، ManagedIdentityCredential) لتجنب مشكلات زمن الانتقال، وبحث بيانات الاعتماد غير المقصودة، والمخاطر الأمنية المحتملة من الآليات الاحتياطية.

تعريف وكلاءك

إنشاء عوامل متخصصة متعددة ستعمل على نفس المهمة بشكل متزامن:

// 2) Helper method to create translation agents
static ChatClientAgent GetTranslationAgent(string targetLanguage, IChatClient chatClient) =>
    new(chatClient,
        $"You are a translation assistant who only responds in {targetLanguage}. Respond to any " +
        $"input by outputting the name of the input language and then translating the input to {targetLanguage}.");

// Create translation agents for concurrent processing
var translationAgents = (from lang in (string[])["French", "Spanish", "English"]
                         select GetTranslationAgent(lang, client));

إعداد التنسيق المتزامن

إنشاء سير العمل باستخدام AgentWorkflowBuilder لتشغيل العوامل بالتوازي:

// 3) Build concurrent workflow
var workflow = AgentWorkflowBuilder.BuildConcurrent(translationAgents);

تشغيل سير العمل المتزامن وجمع النتائج

تنفيذ سير العمل ومعالجة الأحداث من جميع العوامل التي تعمل في وقت واحد:

// 4) Run the workflow
var messages = new List<ChatMessage> { new(ChatRole.User, "Hello, world!") };

await using StreamingRun run = await InProcessExecution.RunStreamingAsync(workflow, messages);
await run.TrySendMessageAsync(new TurnToken(emitEvents: true));

List<ChatMessage> result = new();
await foreach (WorkflowEvent evt in run.WatchStreamAsync())
{
    if (evt is AgentResponseUpdateEvent e)
    {
        Console.WriteLine($"{e.ExecutorId}: {e.Update.Text}");
    }
    else if (evt is WorkflowOutputEvent outputEvt)
    {
        result = outputEvt.As<List<ChatMessage>>()!;
        break;
    }
}

// Display aggregated results from all agents
Console.WriteLine("===== Final Aggregated Results =====");
foreach (var message in result)
{
    Console.WriteLine($"{message.Role}: {message.Text}");
}

إخراج العينة

French_Agent: English detected. Bonjour, le monde !
Spanish_Agent: English detected. ¡Hola, mundo!
English_Agent: English detected. Hello, world!

===== Final Aggregated Results =====
User: Hello, world!
Assistant: English detected. Bonjour, le monde !
Assistant: English detected. ¡Hola, mundo!
Assistant: English detected. Hello, world!

المفاهيم الأساسية

  • التنفيذ المتوازي: يقوم جميع الوكلاء بمعالجة الإدخال في وقت واحد وبشكل مستقل
  • AgentWorkflowBuilder.BuildConcurrent(): إنشاء سير عمل متزامن من مجموعة من العوامل
  • التجميع التلقائي: يتم جمع النتائج من جميع العوامل تلقائيا في النتيجة النهائية
  • دفق الأحداث: المراقبة في الوقت الحقيقي لتقدم العامل من خلال AgentResponseUpdateEvent
  • وجهات نظر متنوعة: يجلب كل وكيل خبرته الفريدة لنفس المشكلة

الوكلاء هم كيانات متخصصة يمكنها معالجة المهام. تحدد التعليمات البرمجية التالية ثلاثة وكلاء: خبير أبحاث وخبير تسويق وخبير قانوني.

import os

from agent_framework.foundry import FoundryChatClient
from azure.identity import AzureCliCredential

# 1) Create three domain agents using FoundryChatClient
chat_client = FoundryChatClient(
    project_endpoint=os.environ["FOUNDRY_PROJECT_ENDPOINT"],
    model=os.environ["FOUNDRY_MODEL"],
    credential=AzureCliCredential(),
)

researcher = chat_client.as_agent(
    instructions=(
        "You're an expert market and product researcher. Given a prompt, provide concise, factual insights,"
        " opportunities, and risks."
    ),
    name="researcher",
)

marketer = chat_client.as_agent(
    instructions=(
        "You're a creative marketing strategist. Craft compelling value propositions and target messaging"
        " aligned to the prompt."
    ),
    name="marketer",
)

legal = chat_client.as_agent(
    instructions=(
        "You're a cautious legal/compliance reviewer. Highlight constraints, disclaimers, and policy concerns"
        " based on the prompt."
    ),
    name="legal",
)

إعداد التنسيق المتزامن

ConcurrentBuilder تسمح لك الفئة بإنشاء سير عمل لتشغيل عوامل متعددة بالتوازي. يمكنك تمرير قائمة الوكلاء كمشاركين.

from agent_framework.orchestrations import ConcurrentBuilder

# 2) Build a concurrent workflow
# Participants are either Agents (type of SupportsAgentRun) or Executors
workflow = ConcurrentBuilder(participants=[researcher, marketer, legal]).build()

تشغيل سير العمل المتزامن وجمع النتائج

ينتج عن المجمع الافتراضي رسالة مساعدة واحدة AgentResponse تحتوي على رسالة مساعدة واحدة لكل مشارك:

from agent_framework import AgentResponse

# 3) Run with a single prompt and print the aggregated agent responses
events = await workflow.run("We are launching a new budget-friendly electric bike for urban commuters.")
outputs = events.get_outputs()

if outputs:
    print("===== Final Aggregated Results =====")
    final: AgentResponse = outputs[0]
    for msg in final.messages:
        name = msg.author_name or "assistant"
        print(f"{'-' * 60}\n\n[{name}]:\n{msg.text}")

إخراج العينة

===== Final Aggregated Results =====
------------------------------------------------------------

[researcher]:
**Insights:**

- **Target Demographic:** Urban commuters seeking affordable, eco-friendly transport;
    likely to include students, young professionals, and price-sensitive urban residents.
- **Market Trends:** E-bike sales are growing globally, with increasing urbanization,
    higher fuel costs, and sustainability concerns driving adoption.
...
------------------------------------------------------------

[marketer]:
**Value Proposition:**
"Empowering your city commute: Our new electric bike combines affordability, reliability, and
    sustainable design—helping you conquer urban journeys without breaking the bank."
...
------------------------------------------------------------

[legal]:
**Constraints, Disclaimers, & Policy Concerns for Launching a Budget-Friendly Electric Bike for Urban Commuters:**

**1. Regulatory Compliance**
- Verify that the electric bike meets all applicable federal, state, and local regulations
    regarding e-bike classification, speed limits, power output, and safety features.

Advanced: Custom Agent Executors

يدعم التزامن المتزامن المنفذين المخصصين الذين يغلفون العوامل بمنطق إضافي. يكون هذا مفيدا عندما تحتاج إلى مزيد من التحكم في كيفية تهيئة الوكلاء وكيفية معالجة الطلبات:

تعريف منفذي العامل المخصص

from agent_framework import (
    AgentExecutorRequest,
    AgentExecutorResponse,
    Agent,
    Executor,
    WorkflowContext,
    handler,
)

class ResearcherExec(Executor):
    def __init__(self, chat_client: FoundryChatClient, id: str = "researcher"):
        self.agent = chat_client.as_agent(
            instructions=(
                "You're an expert market and product researcher. Given a prompt, provide concise, factual insights,"
                " opportunities, and risks."
            ),
            name=id,
        )
        super().__init__(id=id)

    @handler
    async def run(self, request: AgentExecutorRequest, ctx: WorkflowContext[AgentExecutorResponse]) -> None:
        response = await self.agent.run(request.messages)
        full_conversation = list(request.messages) + list(response.messages)
        await ctx.send_message(AgentExecutorResponse(self.id, response, full_conversation=full_conversation))

class MarketerExec(Executor):
    def __init__(self, chat_client: FoundryChatClient, id: str = "marketer"):
        self.agent = chat_client.as_agent(
            instructions=(
                "You're a creative marketing strategist. Craft compelling value propositions and target messaging"
                " aligned to the prompt."
            ),
            name=id,
        )
        super().__init__(id=id)

    @handler
    async def run(self, request: AgentExecutorRequest, ctx: WorkflowContext[AgentExecutorResponse]) -> None:
        response = await self.agent.run(request.messages)
        full_conversation = list(request.messages) + list(response.messages)
        await ctx.send_message(AgentExecutorResponse(self.id, response, full_conversation=full_conversation))

إنشاء سير عمل باستخدام المنفذين المخصصين

chat_client = FoundryChatClient(
    project_endpoint=os.environ["FOUNDRY_PROJECT_ENDPOINT"],
    model=os.environ["FOUNDRY_MODEL"],
    credential=AzureCliCredential(),
)

researcher = ResearcherExec(chat_client)
marketer = MarketerExec(chat_client)
legal = LegalExec(chat_client)

workflow = ConcurrentBuilder(participants=[researcher, marketer, legal]).build()

متقدم: مجمع مخصص

بشكل افتراضي، يجمع التزامن المتزامن جميع استجابات العامل في استجابة واحدة AgentResponse مع رسالة مساعد واحدة لكل مشارك. يمكنك تجاوز هذا السلوك باستخدام مجمع مخصص يعالج النتائج بطريقة معينة:

تعريف مجمع مخصص

from agent_framework import AgentExecutorResponse

# Create a summarizer agent for the aggregator
summarizer_agent = chat_client.as_agent(
    instructions=(
        "You are a helpful assistant that consolidates multiple domain expert outputs "
        "into one cohesive, concise summary with clear takeaways. Keep it under 200 words."
    ),
    name="summarizer",
)

# Define a custom aggregator callback
async def summarize_results(results: list[AgentExecutorResponse]) -> str:
    # Extract one final assistant message per agent
    expert_sections: list[str] = []
    for r in results:
        try:
            messages = getattr(r.agent_response, "messages", [])
            final_text = messages[-1].text if messages and hasattr(messages[-1], "text") else "(no content)"
            expert_sections.append(f"{r.executor_id}:\n{final_text}")
        except Exception as e:
            expert_sections.append(f"{r.executor_id}: (error: {type(e).__name__}: {e})")

    # Ask the model to synthesize a concise summary of the experts' outputs
    prompt = "\n\n".join(expert_sections)
    response = await summarizer_agent.run(prompt)
    # Return the model's final assistant text as the completion result
    return response.messages[-1].text if response.messages else ""

إنشاء سير عمل باستخدام مجمع مخصص

workflow = (
    ConcurrentBuilder(participants=[researcher, marketer, legal])
    .with_aggregator(summarize_results)
    .build()
)

output = None
async for event in workflow.run("We are launching a new budget-friendly electric bike for urban commuters.", stream=True):
    if event.type == "output":
        output = event.data

if output:
    print("===== Final Consolidated Output =====")
    print(output)

عينة الإخراج مع مجمع مخصص

===== Final Consolidated Output =====
Urban e-bike demand is rising rapidly due to eco-awareness, urban congestion, and high fuel costs,
with market growth projected at a ~10% CAGR through 2030. Key customer concerns are affordability,
easy maintenance, convenient charging, compact design, and theft protection. Differentiation opportunities
include integrating smart features (GPS, app connectivity), offering subscription or leasing options, and
developing portable, space-saving designs. Partnering with local governments and bike shops can boost visibility.

Risks include price wars eroding margins, regulatory hurdles, battery quality concerns, and heightened expectations
for after-sales support. Accurate, substantiated product claims and transparent marketing (with range disclaimers)
are essential. All e-bikes must comply with local and federal regulations on speed, wattage, safety certification,
and labeling. Clear warranty, safety instructions (especially regarding batteries), and inclusive, accessible
marketing are required. For connected features, data privacy policies and user consents are mandatory.

Effective messaging should target young professionals, students, eco-conscious commuters, and first-time buyers,
emphasizing affordability, convenience, and sustainability. Slogan suggestion: "Charge Ahead—City Commutes Made
Affordable." Legal review in each target market, compliance vetting, and robust customer support policies are
critical before launch.

مخرجات متوسطة

بشكل افتراضي، يظهر إخراج المجمع فقط كحدث سير عمل "output" (المحطة الطرفية). مرر intermediate_output_from مع المشاركين الذين تريد تعيينهم كمصادر وسيطة لعرض مخرجاتهم الفردية كأحداث "intermediate" :

workflow = ConcurrentBuilder(
    participants=[researcher, marketer, legal],
    intermediate_output_from=[researcher, marketer, legal],
).build()

يمكنك التعامل مع هذه الأحداث في الوقت الحقيقي في وضع الدفق:

from agent_framework import AgentResponseUpdate

# Track the last author to format streaming output.
last_author: str | None = None

async for event in workflow.run("Analyze our new product launch strategy.", stream=True):
    if event.type == "intermediate" and isinstance(event.data, AgentResponseUpdate):
        update = event.data
        author = update.author_name
        if author != last_author:
            if last_author is not None:
                print()  # Newline between different authors
            print(f"{author}: {update.text}", end="", flush=True)
            last_author = author
        else:
            print(update.text, end="", flush=True)

المفاهيم الأساسية

  • التنفيذ المتوازي: يعمل جميع الوكلاء على المهمة في وقت واحد وبشكل مستقل
  • AgentResponse Output: ينتج عن التجميع الافتراضي رسالة مساعدة واحدة AgentResponse لكل مشارك (لم يتم تضمين مطالبة المستخدم)
  • وجهات نظر متنوعة: يجلب كل وكيل خبرته الفريدة لنفس المشكلة
  • المشاركون المرنون: يمكنك استخدام الوكلاء مباشرة أو تضمينهم في منفذين مخصصين
  • معالجة مخصصة: تجاوز المجمع الافتراضي لتجميع النتائج بطرق خاصة بالمجال
  • المخرجات المتوسطة: تمرير intermediate_output_from=[participant, ...] لعرض إخراج كل مشارك مدرج كأحداث "intermediate" ، بالإضافة إلى الحدث الطرفي "output" للمجمع

يدعم Go مهام سير عمل العامل المتزامنة مع agentworkflow.NewConcurrentWorkflowBuilder. يمكنك أيضا إنشاء نفس النمط يدويا باستخدام حواف المروحة والمروحة عندما تحتاج إلى سلوك منفذ مخصص.

إعداد تكوين Foundry

تكوين نقطة نهاية مشروع Foundry ونشر النموذج والمصادقة:

endpoint := os.Getenv("FOUNDRY_PROJECT_ENDPOINT")
model := cmp.Or(os.Getenv("FOUNDRY_MODEL"), "gpt-4o-mini")

token, err := azidentity.NewDefaultAzureCredential(nil)
if err != nil {
    return err
}

تحذير

azidentity.NewDefaultAzureCredential مناسب للتنمية ولكنه يتطلب دراسة متأنية في الإنتاج. في الإنتاج، ضع في اعتبارك استخدام بيانات اعتماد معينة، مثل azidentity.NewManagedIdentityCredential، لتجنب مشكلات زمن الانتقال، وبحث بيانات الاعتماد غير المقصودة، والمخاطر الأمنية المحتملة من الآليات الاحتياطية.

تعريف وكلاءك

إنشاء عوامل متخصصة متعددة ستعمل على نفس المهمة بشكل متزامن:

newTranslationAgent := func(language string) *agent.Agent {
    return foundryprovider.NewAgent(
        endpoint,
        token,
        foundryprovider.ModelDeployment(model),
        foundryprovider.AgentConfig{
            Instructions: fmt.Sprintf(
                "You are a translation assistant who only responds in %s. Respond to any input by outputting the name of the input language and then translating the input to %s.",
                language,
                language,
            ),
            Config: agent.Config{Name: language},
        },
    )
}

agents := []*agent.Agent{
    newTranslationAgent("French"),
    newTranslationAgent("Spanish"),
    newTranslationAgent("English"),
}

إعداد التنسيق المتزامن

إنشاء سير العمل باستخدام agentworkflow.NewConcurrentWorkflowBuilder:

wf, err := agentworkflow.NewConcurrentWorkflowBuilder(agents...).
    WithName("translation-concurrent").
    Build()
if err != nil {
    return err
}

تشغيل سير العمل المتزامن وجمع النتائج

تشغيل سير العمل مع رسالة مستخدم ورمز تحويل. عند تمكين انبعاث الحدث، تظهر تحديثات العامل كأحداث إخراج سير العمل قبل الإخراج المجمع النهائي.

run, err := inproc.Default.RunStreaming(ctx, wf, []*message.Message{message.NewText("Hello, world!")})
if err != nil {
    return err
}
defer run.Close(ctx)

emitEvents := true
if err := run.SendMessage(ctx, workflow.TurnToken{EmitEvents: &emitEvents}); err != nil {
    return err
}

for evt, err := range run.WatchStream(ctx) {
    if err != nil {
        return err
    }
    if output, ok := evt.(workflow.OutputEvent); ok {
        switch value := output.Output.(type) {
        case *agent.ResponseUpdate:
            fmt.Printf("%s: %s\n", output.ExecutorID, value.String())
        case []*message.Message:
            fmt.Println("===== Final Aggregated Results =====")
            for _, msg := range value {
                fmt.Printf("%s: %s\n", msg.Role, msg.String())
            }
        }
    }
}

إخراج العينة

French: English detected. Bonjour, le monde !
Spanish: English detected. ¡Hola, mundo!
English: English detected. Hello, world!

===== Final Aggregated Results =====
assistant: English detected. Bonjour, le monde !
assistant: English detected. ¡Hola, mundo!
assistant: English detected. Hello, world!

Advanced: Custom Agent Executors

إنشاء مهام سير عمل متزامنة يدويا عندما تحتاج إلى سلوك منفذ مخصص. يمكن للمنفذ المخصص استدعاء عامل ثم المشاركة في سير عمل fan-out/fan-in.

agentExecutor := func(id string, ag *agent.Agent) workflow.ExecutorBinding {
    return workflow.BindNewExecutorFunc(id, func(_ string, executorID string) (*workflow.Executor, error) {
        return workflow.NewExecutor(executorID, func(ctx *workflow.Context, prompt string) (string, error) {
            response, err := ag.RunText(ctx, prompt).Collect()
            if err != nil {
                return "", err
            }
            return response.String(), nil
        }), nil
    })
}

researcher := agentExecutor("researcher", researcherAgent)
marketer := agentExecutor("marketer", marketerAgent)
aggregate := aggregateStrings("ConcurrentAggregationExecutor")

wf, err := workflow.NewBuilder(start).
    AddFanOutEdge(start, []workflow.ExecutorBinding{researcher, marketer}).
    AddFanInBarrierEdge([]workflow.ExecutorBinding{researcher, marketer}, aggregate).
    WithOutputFrom(aggregate).
    Build()

متقدم: مجمع مخصص

استخدم WithAggregator لاستبدال سلوك تجميع الرسائل الافتراضي:

wf, err := agentworkflow.NewConcurrentWorkflowBuilder(agents...).
    WithName("translation-concurrent").
    WithAggregator(func(_ context.Context, batches [][]*message.Message) []*message.Message {
        results := make([]*message.Message, 0, len(batches))
        for _, batch := range batches {
            if len(batch) > 0 {
                results = append(results, batch[len(batch)-1])
            }
        }
        return results
    }).
    Build()
if err != nil {
    return err
}

مخرجات متوسطة

بشكل افتراضي، NewConcurrentWorkflowBuilder يصدر مخرجات المشارك والتجميع كمخرجات سير عمل وسيطة ويصدر النتيجة المجمعة كإخراج المحطة الطرفية. بالنسبة إلى مهام سير عمل المنفذ المخصصة، ضع علامة على منفذي الفرع كمتوسط والمجمع كمحطة طرفية:

wf, err := workflow.NewBuilder(start).
    AddFanOutEdge(start, []workflow.ExecutorBinding{physics, chemistry}).
    AddFanInBarrierEdge([]workflow.ExecutorBinding{physics, chemistry}, aggregate).
    WithIntermediateOutputFrom(physics, chemistry).
    WithOutputFrom(aggregate).
    Build()

يتضمن كل workflow.OutputEvent منها ExecutorID الذي نتج عنه الإخراج. يستخدم OutputEvent.IsIntermediate() لتمييز مخرجات الفرع الوسيط عن التجميع النهائي.

المفاهيم الأساسية

  • التنفيذ المتوازي: يقوم جميع الوكلاء أو المنفذين بمعالجة الإدخال بشكل مستقل.
  • تدفق عمل العامل. NewConcurrentWorkflowBuilder(): إنشاء سير عمل متزامن من مجموعة من العوامل.
  • Fan-out/Fan-in Edges: تستخدم AddFanOutEdge مهام سير العمل المتزامنة المخصصة و AddFanInBarrierEdge.
  • تجميع الرسائل: يقوم المجمع الافتراضي بإرجاع الرسالة الأخيرة من كل مشارك؛ يمكن للمجمعات المخصصة استبدال هذا السلوك.
  • دفق الأحداث: يمكن أن تعرض أحداث الإخراج تحديثات الوكيل الفردية والنتائج المجمعة النهائية.
  • المخرجات المتوسطة: WithIntermediateOutputFrom يضع علامة على المخرجات المحددة باستخدام workflow.OutputTagIntermediate.

Tip

راجع نموذج سير العمل المتزامنونموذج أنماط سير عمل العامل للحصول على أمثلة كاملة قابلة للتشغيل.

الخطوات التالية