Workflow 实战:微博情绪量化股票
前两篇分别介绍了 Workflow 的基础概念和进阶能力,这一篇把它们放进一个完整场景:根据微博舆情生成股票买入、卖出或持有的参考信号。
整个过程分为抓取数据、逐条分析情绪、汇总量化、人工确认和生成报告。流程顺序是固定的,模型只负责情绪判断和报告生成,因此很适合用 Workflow 编排。
抓取微博涉及鉴权和反爬,文章不会展开这部分实现,而是使用 mock 数据。重点是业务如何拆成节点、节点之间如何连接,以及代码节点、agent 节点和 HITL 分别放在哪里。文中的部分代码用于说明拓扑和职责,并不是一个复制后即可运行的完整项目。
流程设计
先看为什么这里使用 Workflow,而不是把股票名称交给一个 agent,让它自行抓取、分析并给出建议。
这条流程需要固定且可审计,但不同步骤适合的实现方式并不相同:
- 抓取数据需要调用 API 或爬虫,属于确定性逻辑,使用代码节点。
- 判断每条微博的情绪倾向需要理解自然语言,使用 agent 节点。
- 汇总分数和计算置信度必须可解释、可审计,使用代码节点。
- 低置信度结果需要人工判断,使用 HITL。
- 将结构化结果组织成自然语言报告,再交给 agent 节点。
如果整条流程都由模型自由决定,既无法保证数据真实,也难以解释最终结论。将模型限制在情绪分析和报告生成两个节点中,其余步骤由代码和连边控制,结果会更稳定。
整体结构是一条带条件分支和人工确认的管道:
对应到上一篇介绍的能力:
| 阶段 | 零件 | 为什么用 |
|---|---|---|
| 批量情绪 | 子工作流 + 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>() 会约束模型生成 sentiment、score 和 reason 字段,调用方再将 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,需要中断恢复时保存检查点。