书接上篇《用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 实现会话三方存储,核心只需三步:

  1. 继承ChatMessageStore,实现AddMessagesAsync(存)、GetMessagesAsync(查)、Serialize(序列化);

  2. 通过ChatMessageStoreFactory将自定义存储绑定到 Agent;

  3. 利用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演示结果


Logo

中国智能体开发者社区,聚焦智能体与大模型开发,提供前沿资讯、实用工具链、开源项目及行业案例。通过技术沙龙、开发者大赛等活动,促进经验交流与协作,助力开发者快速构建创新智能应用。

更多推荐