Workflow 实战:微博情绪量化股票

前两篇分别介绍了 Workflow 的基础概念和进阶能力,这一篇把它们放进一个完整场景:根据微博舆情生成股票买入、卖出或持有的参考信号。

整个过程分为抓取数据、逐条分析情绪、汇总量化、人工确认和生成报告。流程顺序是固定的,模型只负责情绪判断和报告生成,因此很适合用 Workflow 编排。

抓取微博涉及鉴权和反爬,文章不会展开这部分实现,而是使用 mock 数据。重点是业务如何拆成节点、节点之间如何连接,以及代码节点、agent 节点和 HITL 分别放在哪里。文中的部分代码用于说明拓扑和职责,并不是一个复制后即可运行的完整项目。


流程设计

先看为什么这里使用 Workflow,而不是把股票名称交给一个 agent,让它自行抓取、分析并给出建议。

这条流程需要固定且可审计,但不同步骤适合的实现方式并不相同:

  • 抓取数据需要调用 API 或爬虫,属于确定性逻辑,使用代码节点。
  • 判断每条微博的情绪倾向需要理解自然语言,使用 agent 节点。
  • 汇总分数和计算置信度必须可解释、可审计,使用代码节点。
  • 低置信度结果需要人工判断,使用 HITL。
  • 将结构化结果组织成自然语言报告,再交给 agent 节点。


如果整条流程都由模型自由决定,既无法保证数据真实,也难以解释最终结论。将模型限制在情绪分析和报告生成两个节点中,其余步骤由代码和连边控制,结果会更稳定。

整体结构是一条带条件分支和人工确认的管道:

正在渲染 Mermaid 图表...


对应到上一篇介绍的能力:

阶段零件为什么用
批量情绪子工作流 + FanOut/FanIn多条微博需要并发分析后再汇总
情绪分析Agent 节点 + 结构化输出将语义判断交给模型,并约束返回字段
分流条件边按置信度决定要不要打断用户
用户确认HITL / RequestPort低置信度信号先交给用户确认
中间数据共享状态在多个节点之间共享微博原文和分析结果
监控事件 + 自定义事件实时观察进度、报告完成通知
恢复检查点保存执行状态,支持中断后继续运行

数据与抓取

工作流中的消息使用强类型对象。这个示例需要微博原文、单条情绪分析结果、聚合后的交易信号和工作流输入:

using System.Text.Json.Serialization;

/// <summary>一条微博。</summary>
public sealed class WeiboPost
{
    public string Id { get; set; } = "";
    public string Content { get; set; } = "";
    public int Likes { get; set; }      // 点赞数,用作情绪权重
    public DateTime CreatedAt { get; set; }
}

/// <summary>情绪标签。</summary>
public enum Sentiment { Bullish, Bearish, Neutral }   // 利好 / 利空 / 中性

/// <summary>单条微博的情绪分析结果(模型结构化输出)。</summary>
public sealed class SentimentResult
{
    [JsonPropertyName("sentiment")]
    public Sentiment Sentiment { get; set; }

    [JsonPropertyName("score")]       // -1.0(极空) ~ 1.0(极多)
    public double Score { get; set; }

    [JsonPropertyName("reason")]
    public string Reason { get; set; } = "";
}

/// <summary>聚合后的交易信号。</summary>
public sealed record StockSignal(
    string Topic,
    string Action,          // "buy" / "sell" / "hold"
    double Confidence,      // 0.0 ~ 1.0
    double AggregateScore,
    int SampleCount,
    string Rationale);      // 人类可读的依据

/// <summary>工作流输入:要分析的股票/公司话题。</summary>
public sealed record TopicRequest(string Topic);


SentimentResult 会作为模型的结构化输出,通过 [JsonPropertyName] 固定 JSON 字段名。WeiboPost.Likes 可以作为后续量化的权重,StockSignal.Confidence 则用于决定是否进入人工确认。


抓取节点属于确定性逻辑。真实项目可以在这里调用微博 API 或爬虫,示例只返回一组 mock 数据,并将原文写入共享状态:

using Microsoft.Agents.AI.Workflows;

internal sealed partial class FetchNode : Executor
{
    public FetchNode() : base("Fetch") { }

    [MessageHandler]
    private async ValueTask<List<WeiboPost>> HandleAsync(
        TopicRequest req, IWorkflowContext context, CancellationToken ct = default)
    {
        // 真实场景:调微博 API / 爬虫。这里 mock 一批数据。
        List<WeiboPost> posts =
        [
            new() { Id = "p1", Content = "宁德时代新电池能量密度又创新高,稳了!", Likes = 1200, CreatedAt = DateTime.UtcNow },
            new() { Id = "p2", Content = "这股价跌得我心慌,是不是该跑了", Likes = 80, CreatedAt = DateTime.UtcNow },
            new() { Id = "p3", Content = "Q3 财报超预期,看好后续走势", Likes = 560, CreatedAt = DateTime.UtcNow },
            new() { Id = "p4", Content = "新能源内卷太严重,龙头也不好做", Likes = 320, CreatedAt = DateTime.UtcNow },
            new() { Id = "p5", Content = "刚提车,电池确实耐用,支持一波", Likes = 45, CreatedAt = DateTime.UtcNow },
        ];

        // 把原文存进共享状态,下游量化节点要用(避免在消息里反复传大对象)
        await context.QueueStateUpdateAsync(
            $"posts:{req.Topic}", posts, scopeName: "WeiboStock");

        await context.AddEventAsync(new ProgressEvent($"抓取到 {posts.Count} 条微博"));
        return posts;   // 沿边送给情绪分析节点
    }
}

/// <summary>自定义进度事件,用于实时观察。</summary>
internal sealed class ProgressEvent(string message) : WorkflowEvent(message) { }


FetchNode 接收 TopicRequest 并返回 List<WeiboPost>,返回值会自动沿边发送。微博原文保存在 "WeiboStock" scope 中,下游需要按点赞数加权时可以再读取,不必把完整对象反复放进消息。ProgressEvent 只用于通知调用方抓取进度,不参与节点之间的消息传递。


批量情绪分析

多条微博需要逐条分析,并在全部完成后汇总。这部分可以封装成子工作流,使主流程只关心输入的 List<WeiboPost> 和输出的 List<SentimentResult>

单条微博的情绪判断交给 ChatClientAgent,并使用结构化输出约束返回的 SentimentResult

using Azure.AI.OpenAI;
using Azure.Identity;
using Microsoft.Agents.AI;
using Microsoft.Extensions.AI;

// 创建一个情绪分析 agent
AIAgent sentimentAgent = chatClient.AsAIAgent(new ChatClientAgentOptions
{
    Name = "SentimentAgent",
    ChatOptions = new()
    {
        Instructions = """
            你是一个 A 股舆情分析师。判断给定微博对相关公司的情绪倾向。
            - bullish(利好):对公司/股价看好
            - bearish(利空):看空、担忧、负面
            - neutral(中性):陈述事实、无明显倾向
            score 范围 -1.0(极度利空)到 1.0(极度利好)。
            reason 用一句话说明判断依据。
            """,
        ResponseFormat = ChatResponseFormat.ForJsonSchema<SentimentResult>()
    }
});


ForJsonSchema<SentimentResult>() 会约束模型生成 sentimentscorereason 字段,调用方再将 JSON 反序列化成强类型对象。

批量处理的拓扑可以拆成 Split、Adapter 和 Aggregate 三类节点。Split 拆分列表,Adapter 调用 agent,Aggregate 收集结果:

// ── 子工作流:BatchSentimentWorkflow ──
// 入口:接收 List<WeiboPost>
// 内部:拆分 → 分发给分析节点 → 汇总
// 出口:产出 List<SentimentResult>

public sealed record PostBatchItem(int Index, WeiboPost Post);
public sealed record PostResultItem(int Index, SentimentResult Result);

// 1. 拆分节点:List<WeiboPost> → 逐条 PostBatchItem 发出去(fan-out 触发点)
internal sealed partial class SplitNode : Executor
{
    public SplitNode() : base("Split") { }

    [MessageHandler]
    private async ValueTask HandleAsync(
        List<WeiboPost> posts, IWorkflowContext context, CancellationToken ct = default)
    {
        for (int i = 0; i < posts.Count; i++)
            await context.SendMessageAsync(new PostBatchItem(i, posts[i]));  // 逐条发,触发 fan-out
    }
}

// 2. 汇总节点:收集所有 PostResultItem → List<SentimentResult>(fan-in 接收点)
internal sealed partial class AggregateNode : Executor
{
    private readonly List<(int Index, SentimentResult Result)> _collected = [];

    public AggregateNode() : base("Aggregate") { }

    [MessageHandler]
    private ValueTask<List<SentimentResult>> HandleAsync(
        PostResultItem item, IWorkflowContext context, CancellationToken ct = default)
    {
        _collected.Add((item.Index, item.Result));
        // 注意:这里有个"何时算收齐"的问题。
        // FanIn barrier 边保证:所有源都产出后,这些消息才送到本节点。
        // 但本节点要等到"全部"到齐才能输出,实际实现需要配合 FanIn 的语义。
        return ValueTask.FromResult(
            _collected.OrderBy(x => x.Index).Select(x => x.Result).ToList());
    }
}

// 3. 把 agent 包成"吃 PostBatchItem,吐 PostResultItem"的适配节点
//    (因为 agent 直接吃的是文本,需要一个适配层做 WeiboPost → 文本 → SentimentResult 的转换)
internal sealed partial class SentimentAdapter(AIAgent agent) : Executor("SentimentAdapter")
{
    [MessageHandler]
    private async ValueTask<PostResultItem> HandleAsync(
        PostBatchItem item, IWorkflowContext context, CancellationToken ct = default)
    {
        var resp = await agent.RunAsync(item.Post.Content, cancellationToken: ct);
        var result = System.Text.Json.JsonSerializer.Deserialize<SentimentResult>(resp.Text)!;
        return new PostResultItem(item.Index, result);
    }
}

// 4. 组装子工作流
static Workflow BuildBatchSentimentWorkflow(AIAgent sentimentAgent)
{
    var split = new SplitNode();
    var adapter = new SentimentAdapter(sentimentAgent);
    var aggregate = new AggregateNode();

    return new WorkflowBuilder(split)
        .AddFanOutEdge(split, [adapter])           // 拆分节点逐条发 → adapter 并行处理
        .AddFanInBarrierEdge([adapter], aggregate) // adapter 全部完成 → 汇总
        .WithOutputFrom(aggregate)
        .Build();
}


这段代码只用于说明拓扑,其中只有一个 Adapter,因此不能代表按微博数量动态扩展的 worker 池。AggregateNode 还需要知道本批次的总数,收齐后再通过 YieldOutputAsync 产生一次完整结果。实际项目可以创建固定数量的 Adapter 节点分担消息,也可以在一个批处理 Executor 中使用受控并发。无论采用哪种方式,都要限制并发数,避免同时向模型服务发出过多请求。

子工作流建好后,使用 BindAsExecutor 包装成主流程中的一个节点:

Workflow batchWorkflow = BuildBatchSentimentWorkflow(sentimentAgent);
ExecutorBinding batchNode = batchWorkflow.BindAsExecutor("BatchSentiment");
// 后面主流程里,FetchNode → batchNode → QuantifyNode


封装后,主流程只看到批量分析的输入和输出。相同逻辑还可以用于新闻或评论分析,也更容易单独测试。


量化、确认与报告

拿到所有 SentimentResult 后,QuantifyNode 通过确定性公式生成交易信号。这部分需要可解释、可审计,因此使用代码节点:

internal sealed partial class QuantifyNode : Executor
{
    private const double BuyThreshold = 0.3;    // 综合分 > 0.3 → buy
    private const double SellThreshold = -0.3;  // 综合分 < -0.3 → sell

    public QuantifyNode() : base("Quantify") { }

    [MessageHandler]
    private async ValueTask<StockSignal> HandleAsync(
        List<SentimentResult> results, IWorkflowContext context, CancellationToken ct = default)
    {
        if (results.Count == 0)
            return new StockSignal("", "hold", 0, 0, 0, "无样本");

        // 简单聚合:平均分(真实场景会用_likes 加权,从共享状态取原文)
        double avg = results.Average(r => r.Score);
        double confidence = Math.Min(1.0, results.Count / 20.0);  // 样本越多越自信,20 条封顶

        string action = avg switch
        {
            > BuyThreshold => "buy",
            < SellThreshold => "sell",
            _ => "hold"
        };

        var rationale = $"共 {results.Count} 条样本,平均情绪分 {avg:F2}。" +
                        $"看多 {results.Count(r => r.Sentiment == Sentiment.Bullish)} 条," +
                        $"看空 {results.Count(r => r.Sentiment == Sentiment.Bearish)} 条," +
                        $"中性 {results.Count(r => r.Sentiment == Sentiment.Neutral)} 条。";

        await context.AddEventAsync(new ProgressEvent($"量化完成: {action} (置信度 {confidence:F2})"));
        return new StockSignal("", action, confidence, avg, results.Count, rationale);
    }
}


confidence 决定下一步是否需要人工确认。示例只根据样本数量计算置信度,真实项目还可以加入点赞权重、时间衰减和内容类型等因素。这些仍然是确定性公式,不需要交给模型。

置信度较低时,工作流通过 RequestPort 暂停并询问用户是否继续。先定义请求和响应类型:

// 请求中保留原始信号,用户响应后还要继续交给报告节点
public sealed record ConfirmRequest(string Message, StockSignal Signal);
public sealed record ConfirmResponse(bool Proceed, StockSignal Signal);

// 在构建主工作流时定义端口
RequestPort confirmPort = RequestPort.Create<ConfirmRequest, ConfirmResponse>("ConfirmProceed");

<br />QuantifyNode 输出 StockSignal,而端口接收 ConfirmRequest,因此需要一个适配节点转换消息类型:

internal sealed class CreateConfirmRequestNode()
    : Executor<StockSignal, ConfirmRequest>("CreateConfirmRequest")
{
    public override ValueTask<ConfirmRequest> HandleAsync(
        StockSignal signal,
        IWorkflowContext context,
        CancellationToken cancellationToken = default)
    {
        return ValueTask.FromResult(
            new ConfirmRequest(
                $"当前信号置信度为 {signal.Confidence:F2},是否继续生成报告?",
                signal));
    }
}


然后按照 Confidence 分流。低置信度结果先进入适配节点,再交给确认端口;高置信度结果直接进入报告节点:

// 两个出口:低置信度 → 确认请求;高置信度 → 报告节点
const double ConfidenceThreshold = 0.5;

builder.AddEdge<StockSignal>(quantify, createConfirmRequest,
    condition: signal => signal.Confidence < ConfidenceThreshold);   // 弱信号 → 确认
builder.AddEdge<StockSignal>(quantify, reportNode,
    condition: signal => signal.Confidence >= ConfidenceThreshold);  // 强信号 → 直接报告
builder.AddEdge(createConfirmRequest, confirmPort);


用户通过 RequestPort 返回 ConfirmResponse 后,使用一个 gate 节点判断是否继续。用户同意时,节点将原始 StockSignal 发送给报告节点;不同意时不再发送消息,这条路径随即结束:

internal sealed partial class ConfirmResponseNode() : Executor("ConfirmResponse")
{
    [MessageHandler]
    private async ValueTask HandleAsync(
        ConfirmResponse response,
        IWorkflowContext context,
        CancellationToken cancellationToken = default)
    {
        if (response.Proceed)
        {
            await context.SendMessageAsync(response.Signal, cancellationToken);
        }
    }
}


运行到 confirmPort 时,工作流发布 RequestInfoEvent 并暂停。调用方通过事件中的 ExternalRequest 创建关联响应,再调用 SendResponseAsync 发送,工作流便可以继续执行。

最后,报告节点将结构化的 StockSignal 组织成自然语言:

AIAgent reportAgent = chatClient.AsAIAgent(new ChatClientAgentOptions
{
    Name = "ReportAgent",
    ChatOptions = new()
    {
        Instructions = """
            你是一个 A 股舆情分析师。根据给定的交易信号,写一段简短(150 字内)的投资参考。
            要客观、提示风险,不要夸大其词。包含:结论、主要依据、风险提示。
            """,
        // 报告不需要结构化输出,自然语言即可
    }
});

// 用一个适配节点把 StockSignal 喂给 reportAgent,把输出作为工作流输出 yield 出去
internal sealed partial class ReportNode(AIAgent reportAgent) : Executor("Report")
{
    [MessageHandler]
    private async ValueTask HandleAsync(
        StockSignal signal, IWorkflowContext context, CancellationToken ct = default)
    {
        string prompt = $"""
            话题: {signal.Topic}
            建议: {signal.Action}
            置信度: {signal.Confidence:F2}
            综合情绪分: {signal.AggregateScore:F2}
            依据: {signal.Rationale}
            """;
        var resp = await reportAgent.RunAsync(prompt, cancellationToken: ct);
        await context.YieldOutputAsync(resp.Text);   // 作为工作流最终输出
    }
}

// 注册为输出源
builder.WithOutputFrom(reportNode);


ReportNode 通过 YieldOutputAsync 产生工作流输出。配合 WithOutputFrom(reportNode),报告会以 WorkflowOutputEvent 暴露给调用方。

组装与运行

各个节点准备好后,将它们组装成主工作流:

using Microsoft.Agents.AI.Workflows;

Workflow BuildMainWorkflow(IChatClient chatClient)
{
    // 创建 agents
    AIAgent sentimentAgent = chatClient.AsAIAgent(new ChatClientAgentOptions { /* 同前 */ });
    AIAgent reportAgent    = chatClient.AsAIAgent(new ChatClientAgentOptions { /* 同前 */ });

    // 构建子工作流并包成节点
    Workflow batchWf = BuildBatchSentimentWorkflow(sentimentAgent);
    ExecutorBinding batchNode = batchWf.BindAsExecutor("BatchSentiment");

    // 主流程节点
    var fetch     = new FetchNode();
    var quantify  = new QuantifyNode();
    var report    = new ReportNode(reportAgent);
    var createConfirmRequest = new CreateConfirmRequestNode();
    var confirmResponse = new ConfirmResponseNode();
    RequestPort confirmPort = RequestPort.Create<ConfirmRequest, ConfirmResponse>("ConfirmProceed");

    const double ConfidenceThreshold = 0.5;

    return new WorkflowBuilder(fetch)
        // 抓取 → 批量情绪分析(子工作流)
        .AddEdge(fetch, batchNode)
        // 批量结果 → 量化
        .AddEdge(batchNode, quantify)
        // 量化结果分流:弱信号 → 创建确认请求
        .AddEdge<StockSignal>(quantify, createConfirmRequest,
            condition: s => s.Confidence < ConfidenceThreshold)
        // 强信号 → 直接报告
        .AddEdge<StockSignal>(quantify, report,
            condition: s => s.Confidence >= ConfidenceThreshold)
        // 请求外部确认,确认通过后恢复 StockSignal
        .AddEdge(createConfirmRequest, confirmPort)
        .AddEdge(confirmPort, confirmResponse)
        .AddEdge(confirmResponse, report)
        // 报告节点是输出源
        .WithOutputFrom(report)
        .WithName("微博情绪量化股票")
        .WithDescription("抓取→情绪分析→量化→确认→报告")
        .Build();
}


运行前可以先生成 Mermaid 流程图,检查连边是否符合设计:

Workflow wf = BuildMainWorkflow(chatClient);
Console.WriteLine(wf.ToMermaidString());


将输出放进支持 Mermaid 的 Markdown 渲染器,就能看到工作流结构。确认拓扑后开始运行。带 HITL 的工作流需要监听请求事件,并通过 ExternalResponse 返回用户的选择:

Workflow workflow = BuildMainWorkflow(chatClient);

await using StreamingRun run = await InProcessExecution.RunStreamingAsync(
    workflow, input: new TopicRequest("宁德时代"));

await foreach (WorkflowEvent evt in run.WatchStreamAsync())
{
    switch (evt)
    {
        case ProgressEvent progress:
            Console.WriteLine($"[进度] {progress.Data}");
            break;

        case RequestInfoEvent requestEvent:
            if (!requestEvent.Request.TryGetDataAs<ConfirmRequest>(out var request))
                throw new InvalidOperationException("无法读取确认请求");

            Console.WriteLine(request.Message);
            bool proceed = Console.ReadLine()?.Trim().Equals("y", StringComparison.OrdinalIgnoreCase) == true;
            ExternalResponse response = requestEvent.Request.CreateResponse(
                new ConfirmResponse(proceed, request.Signal));
            await run.SendResponseAsync(response);
            break;

        case WorkflowOutputEvent output:
            Console.WriteLine($"\n===== 舆情报告 =====\n{output.Data}");
            break;

        case ExecutorFailedEvent failed:
            Console.Error.WriteLine($"节点失败: {failed.Data}");
            break;

        case WorkflowErrorEvent error:
            Console.Error.WriteLine($"工作流错误: {error.Data}");
            break;
    }
}


同一个 WatchStreamAsync 循环处理进度、外部请求、最终输出和错误。用户拒绝继续时,ConfirmResponseNode 不会向报告节点发送消息,这条执行路径会自然结束。


恢复与扩展

如果流程可能长时间等待确认,或者需要在进程重启后继续运行,可以加入检查点:

using Microsoft.Agents.AI.Workflows.Checkpointing;

// 用 JSON 文件存储检查点(生产可换成云存储)
// 注意:构造函数接收 DirectoryInfo,不是 string
DirectoryInfo checkpointDir = Directory.CreateDirectory("./checkpoints");
using FileSystemJsonCheckpointStore store = new(checkpointDir);
CheckpointManager mgr = CheckpointManager.CreateJson(store);

// 把 mgr 传给执行入口
StreamingRun run = await InProcessExecution.RunStreamingAsync(
    workflow, new TopicRequest("宁德时代"), mgr);

// 监听 SuperStepCompletedEvent 取检查点(每个超级步边界都会存一次)
List<CheckpointInfo> checkpoints = [];
await foreach (WorkflowEvent evt in run.WatchStreamAsync())
{
    if (evt is SuperStepCompletedEvent step && step.CompletionInfo?.Checkpoint is not null)
        checkpoints.Add(step.CompletionInfo.Checkpoint);
    // ... 其他事件处理
}

// 崩了之后,从最近检查点恢复:
// StreamingRun resumed = await InProcessExecution.ResumeStreamingAsync(
//     workflow, checkpoints[^1], mgr);


检查点在超级步边界创建。HITL 暂停后,待处理请求和共享状态都会保存下来;恢复运行时,工作流仍会等待对应的外部响应。微博原文放在共享状态中,也正是为了让恢复后的节点能够继续读取。

这套示例还可以沿几个方向扩展:

  • QuantifyNode 读取 WeiboPost.Likes,加入点赞权重和时间衰减。
  • 同时抓取雪球、股吧等来源,再使用 FanIn 汇总。
  • 让看多和看空 agent 分别分析争议内容,再交给汇总节点。
  • 使用 Microsoft.Agents.AI.Workflows.Declarative 将流程改为声明式配置。
  • 接入 Microsoft.Agents.AI.DevUI 查看和调试运行过程。


无论场景如何扩展,拆分原则没有变化:确定且需要审计的步骤使用代码节点,需要理解和生成自然语言的步骤使用 agent,需要外部决策时使用 HITL,需要中断恢复时保存检查点。