用 Microsoft Agent Framework 实现会话记录三方存储,让对话持久化不丢失
书接上篇《用Microsoft Agent Framework 实现函数调用人工批准:让 AI 操作更可控》,让我来继续聊聊 Microsoft Agent Framework 。
在使用 AI Agent 开发对话类应用时,你是否遇到过这样的困扰:默认的会话记录存在内存中,服务重启就丢失,多实例部署无法共享对话上下文,甚至无法满足数据持久化的合规要求?
Mircrosoft Agent Framework 的ChatMessageStore机制完美解决了这个问题 —— 通过自定义三方存储,让 Agent 的会话记录可以持久化到向量库、Redis、数据库等外部存储中。今天就结合实战 Demo,聊聊如何实现 Agent 会话记录的三方存储。
为什么需要 “三方存储” 会话记录?
默认情况下,Agent 的会话历史会存储在AgentThread的内存中,或依赖底层推理服务的存储能力,这在实际应用中存在明显局限:
-
内存存储易丢失:服务重启、进程销毁后,对话历史直接消失;
-
不支持多实例共享:分布式部署时,不同实例无法获取同一用户的历史对话;
-
数据管控不足:无法满足数据备份、合规审计、长期存储的业务需求;
-
容量受限:内存存储难以支撑大量长会话的历史数据管理。
而自定义ChatMessageStore,可以将会话记录迁移到外部存储,彻底解决这些问题。
而自定义ChatMessageStore,可以将会话记录迁移到外部存储,彻底解决这些问题。
Demo 核心:基于向量库的会话持久化实现
本次 Demo 以 “向量库存储会话记录” 为例,实现了会话的持久化存储、历史查询和线程序列化功能。核心目标是:让 Agent 的对话历史存储在外部向量库中,即使重启服务,也能通过线程序列化恢复对话上下文。
下面一步步拆解实现逻辑。
1. 核心依赖:先安装必要的 NuGet 包
要实现三方存储,需要依赖向量库连接器和 Agent Framework 相关包,通过.NET CLI 安装:
<PackageReference Include="Microsoft.Agents.AI.OpenAI" Version="1.0.0-preview.251105.1" /><PackageReference Include="Microsoft.SemanticKernel.Connectors.InMemory" Version="1.67.0-preview" />
2. 自定义存储核心:实现 VectorChatMessageStore
要对接外部存储,必须继承抽象类ChatMessageStore,并实现三个关键方法:AddMessagesAsync(存储消息)、GetMessagesAsync(查询消息)、Serialize(序列化状态)。
Demo 中实现了基于向量库的存储类VectorChatMessageStore,核心逻辑如下:
internal sealed class VectorChatMessageStore : ChatMessageStore{ private readonly VectorStore _vectorStore; // 外部存储载体(向量库) public string? ThreadDbKey { get; private set; } // 会话唯一标识 // 构造函数:初始化向量库,反序列化会话标识 public VectorChatMessageStore(VectorStore vectorStore, JsonElement serializedStoreState, JsonSerializerOptions? jsonSerializerOptions = null) { this._vectorStore = vectorStore ?? throw new ArgumentNullException(nameof(vectorStore)); // 从序列化状态中恢复会话标识(支持线程重启后复用历史) if (serializedStoreState.ValueKind is JsonValueKind.String) { this.ThreadDbKey = serializedStoreState.Deserialize<string>(); } } // 关键方法1:将新消息存储到向量库 public override async Task AddMessagesAsync(IEnumerable<ChatMessage> messages, CancellationToken cancellationToken = default) { // 为会话生成唯一标识(首次存储时创建) this.ThreadDbKey ??= Guid.NewGuid().ToString("N"); // 获取向量库中的会话集合(不存在则自动创建) var collection = _vectorStore.GetCollection<string, ChatHistoryItem>("ChatHistory"); await collection.EnsureCollectionExistsAsync(cancellationToken); // 将消息序列化后存入向量库,关联会话标识和时间戳 await collection.UpsertAsync(messages.Select(x => new ChatHistoryItem() { Key = ThreadDbKey + x.MessageId, // 消息唯一键(会话标识+消息ID) Timestamp = DateTimeOffset.UtcNow, // 存储时间戳 ThreadId = ThreadDbKey, // 关联当前会话 SerializedMessage = JsonSerializer.Serialize(x), // 序列化消息内容 MessageText = x.Text // 消息文本(用于后续检索) }), cancellationToken); } // 关键方法2:从向量库查询会话历史 public override async Task<IEnumerable<ChatMessage>> GetMessagesAsync(CancellationToken cancellationToken = default) { var collection = _vectorStore.GetCollection<string, ChatHistoryItem>("ChatHistory"); await collection.EnsureCollectionExistsAsync(cancellationToken); // 按会话标识查询,按时间戳倒序取最新10条(可根据模型上下文窗口调整) var records = collection.GetAsync( x => x.ThreadId == ThreadDbKey, 10, new() { OrderBy = x => x.Descending(y => y.Timestamp) }, cancellationToken); var messages = new List<ChatMessage>(); await foreach (var record in records) { // 反序列化消息内容,恢复ChatMessage对象 messages.Add(JsonSerializer.Deserialize<ChatMessage>(record.SerializedMessage!)!); } // 反转顺序,返回按时间升序的历史(符合Agent上下文处理逻辑) messages.Reverse(); return messages; } // 关键方法3:序列化会话状态(支持线程持久化) public override JsonElement Serialize(JsonSerializerOptions? jsonSerializerOptions = null) { // 序列化会话标识,让线程重启后能找到对应历史 return JsonSerializer.SerializeToElement(this.ThreadDbKey); } // 向量库存储模型:定义消息在向量库中的存储结构 private sealed class ChatHistoryItem { [VectorStoreKey] public string? Key { get; set; } // 唯一键 [VectorStoreData] public string? ThreadId { get; set; } // 会话标识 [VectorStoreData] public DateTimeOffset? Timestamp { get; set; } // 时间戳 [VectorStoreData] public string? SerializedMessage { get; set; } // 序列化消息 [VectorStoreData] public string? MessageText { get; set; } // 消息文本 }}
核心设计亮点:
-
每个会话生成唯一
ThreadDbKey,确保多会话数据隔离; -
消息存储时关联时间戳,查询时支持按时间排序;
-
支持序列化
ThreadDbKey,让AgentThread可持久化和恢复。
3. 集成到 Agent:通过工厂函数绑定自定义存储
创建 Agent 时,需要通过ChatMessageStoreFactory将自定义存储与AgentThread绑定,让每个线程自动使用三方存储。
public static async Task DemoAsync(string apiKey, string modelName, string endpoint){ // 初始化OpenAI客户端 var clientOptions = new OpenAIClientOptions { Endpoint = new Uri(endpoint) }; var agent = new OpenAIClient(new ApiKeyCredential(apiKey), clientOptions) .GetChatClient(modelName) .CreateAIAgent(new ChatClientAgentOptions { Instructions = "你是一个擅长讲笑话的Agent", Name = "ZerekZhang", // 关键:绑定自定义存储工厂,每个线程创建独立存储实例 ChatMessageStoreFactory = ctx => { // 这里用内存向量库示例,可替换为Redis、PostgreSQL等实际存储 return new VectorChatMessageStore( new InMemoryVectorStore(), ctx.SerializedState, ctx.JsonSerializerOptions); } }); // 1. 创建新线程(自动关联自定义存储) AgentThread thread = agent.GetNewThread(); // 2. 序列化线程(可存储到数据库/文件,支持后续恢复) JsonElement serializedThread = thread.Serialize(); // 3. 恢复线程(模拟服务重启后恢复对话) AgentThread resumedThread = agent.DeserializeThread(serializedThread); // 交互循环:接收用户输入,调用Agent并持久化对话 while (true) { var userInput = Console.ReadLine(); if (userInput == "Exit") { // 退出时,从三方存储中查询并打印所有历史消息 var messageStore = resumedThread.GetService<VectorChatMessageStore>()!; var messages = await messageStore.GetMessagesAsync(); foreach (var item in messages) { Console.WriteLine($"历史消息:{item.Role} - {item.Text}"); } break; } // 调用Agent,对话历史自动存入三方存储 var response = await agent.RunAsync(userInput, resumedThread); Console.WriteLine("Agent Output: " + response); }}
关键流程说明:
-
工厂函数
ChatMessageStoreFactory:为每个新线程创建独立的VectorChatMessageStore实例,确保会话隔离; -
线程序列化:
thread.Serialize()可将线程状态(含ThreadDbKey)序列化,支持持久化到文件 / 数据库; -
历史查询:通过
resumedThread.GetService<VectorChatMessageStore>()获取存储实例,调用GetMessagesAsync查询历史。
4. 灵活替换:将存储载体换成 Redis / 数据库
Demo 中用的是InMemoryVectorStore(内存向量库),仅作示例。实际开发中,可轻松替换为其他存储载体:
-
分布式存储:Redis(适合高并发场景)、MongoDB(文档型存储);
-
关系型数据库:PostgreSQL(支持向量存储)、MySQL;
-
云存储:Azure Cosmos DB、AWS DynamoDB。
实际应用场景:哪里需要会话三方存储?
这个方案的实用性极强,适用于大多数 Agent 对话场景:
-
客服机器人:持久化用户咨询历史,支持跨会话上下文关联;
-
企业内部助手:满足数据合规要求,对话历史可审计、可备份;
-
多端同步应用:手机、PC、网页端共享同一对话上下文;
-
长会话场景:比如代码助手、文档问答,需要长期保留对话逻辑。
核心要点总结
用微软 Agent Framework 实现会话三方存储,核心只需三步:
-
继承
ChatMessageStore,实现AddMessagesAsync(存)、GetMessagesAsync(查)、Serialize(序列化); -
通过
ChatMessageStoreFactory将自定义存储绑定到 Agent; -
利用
AgentThread的序列化能力,支持会话上下文的持久化与恢复。
这个方案既保留了 Agent 的原生对话能力,又解决了内存存储的诸多痛点,让 Agent 的对话数据 “可控、可存、可复用”。
本Demo全部代码
using Microsoft.Agents.AI;using Microsoft.Extensions.AI;using Microsoft.Extensions.VectorData;using Microsoft.SemanticKernel.Connectors.InMemory;using OpenAI;using System.ClientModel;using System.Text.Json;using VectorStore = Microsoft.Extensions.VectorData.VectorStore;
namespace AgentDemo{#pragma warning disable OPENAI001 /// <summary> /// Agent 会话记录三方存储 /// </summary> internal static partial class AgentConversationSaveBase { public static async Task DemoAsync(string apiKey, string modelName, string endpoint) { var clientOptions = new OpenAIClientOptions { Endpoint = new Uri(endpoint) }; var agent = new OpenAIClient(new ApiKeyCredential(apiKey), clientOptions) .GetChatClient(modelName) .CreateAIAgent(new ChatClientAgentOptions { Instructions = "你是一个擅长讲笑话的Agent", Name = "ZerekZhang", ChatMessageStoreFactory = ctx => { return new VectorChatMessageStore(new InMemoryVectorStore(), ctx.SerializedState, ctx.JsonSerializerOptions); } });
AgentThread thread = agent.GetNewThread(); JsonElement serializedThread = thread.Serialize(); AgentThread resumedThread = agent.DeserializeThread(serializedThread); while (true) { var userInput = Console.ReadLine(); if (userInput == "Exit") { var messageStore = resumedThread.GetService<VectorChatMessageStore>()!; var messages = await messageStore.GetMessagesAsync(); foreach (var item in messages) { Console.WriteLine(item); } break; } Console.WriteLine("Agent Output" + await agent.RunAsync(userInput, resumedThread)); } } }
internal sealed class VectorChatMessageStore : ChatMessageStore { private readonly VectorStore _vectorStore; public string? ThreadDbKey { get; private set; } public VectorChatMessageStore(VectorStore vectorStore, JsonElement serializedStoreState, JsonSerializerOptions? jsonSerializerOptions = null) { this._vectorStore = vectorStore ?? throw new ArgumentNullException(nameof(vectorStore));
if (serializedStoreState.ValueKind is JsonValueKind.String) { this.ThreadDbKey = serializedStoreState.Deserialize<string>(); } } public override async Task AddMessagesAsync(IEnumerable<ChatMessage> messages, CancellationToken cancellationToken = default) { this.ThreadDbKey ??= Guid.NewGuid().ToString("N"); var collection = this._vectorStore.GetCollection<string, ChatHistoryItem>("ChatHistory"); await collection.EnsureCollectionExistsAsync(cancellationToken); await collection.UpsertAsync(messages.Select(x => new ChatHistoryItem() { Key = this.ThreadDbKey + x.MessageId, Timestamp = DateTimeOffset.UtcNow, ThreadId = this.ThreadDbKey, SerializedMessage = JsonSerializer.Serialize(x), MessageText = x.Text }), cancellationToken); }
public override async Task<IEnumerable<ChatMessage>> GetMessagesAsync(CancellationToken cancellationToken = default) { var collection = this._vectorStore.GetCollection<string, ChatHistoryItem>("ChatHistory"); await collection.EnsureCollectionExistsAsync(cancellationToken); var records = collection .GetAsync(x => x.ThreadId == this.ThreadDbKey, 10, new() { OrderBy = x => x.Descending(y => y.Timestamp) }, cancellationToken);
List<ChatMessage> messages = []; await foreach (var record in records) { messages.Add(JsonSerializer.Deserialize<ChatMessage>(record.SerializedMessage!)!); }
messages.Reverse(); return messages; }
public override JsonElement Serialize(JsonSerializerOptions? jsonSerializerOptions = null) { return JsonSerializer.SerializeToElement(this.ThreadDbKey); }
private sealed class ChatHistoryItem { [VectorStoreKey] public string? Key { get; set; } [VectorStoreData] public string? ThreadId { get; set; } [VectorStoreData] public DateTimeOffset? Timestamp { get; set; } [VectorStoreData] public string? SerializedMessage { get; set; } [VectorStoreData] public string? MessageText { get; set; } } }#pragma warning restore OPENAI001}
Demo演示结果

更多推荐


所有评论(0)